3.2 duplicatedetector 模块

3.2.1 模块入口 — CheckDuplicateDevices()

文件: manager.go L34-L46

失败

成功

nil

非nil

CheckDuplicateDevices(ctx, config)

once.Do: 首次调用

NewManager创建管理器

日志错误
manager = nil

manager保存到包级变量

manager != nil?

日志错误, return

go manager.Start(ctx)
异步启动

func CheckDuplicateDevices(ctx context.Context, config *types.DetectorConfig) {
    // 使用sync.Once确保Manager只创建一次(单例模式)
    once.Do(func() {
        var err error
        manager, err = NewManager(config)
        if err != nil {
            hwlog.RunLog.Errorf("Failed to create manager: %v", err)
        }
    })
    
    // 检查Manager是否创建成功
    if manager == nil {
        hwlog.RunLog.Error("manager is nil")
        return
    }
    
    // 异步启动检测,不阻塞主流程
    go manager.Start(ctx)
}

设计意图

  1. sync.Once 保证单例,即使 CheckDuplicateDevices 被多次调用,Manager 也只创建一次
  2. go 异步启动,不阻塞 Device Plugin 主初始化流程
  3. manageronce 是包级变量,全局唯一
3.2.2 Manager 创建 — NewManager()

文件: manager.go L57-L72

func NewManager(config *types.DetectorConfig) (*Manager, error) {
    if config == nil {
        return nil, errors.New("config is nil")
    }
    
    // 根据配置创建容器运行时客户端
    // config.RuntimeType 决定创建 Docker 还是 containerd 客户端
    client, err := containerruntime.NewClient(config)
    if err != nil {
        return nil, fmt.Errorf("failed to create container runtime client: %w", err)
    }
    
    return &Manager{
        client: client,
        cache:  cache.NewContainerCache(),  // 初始化空缓存
    }, nil
}
3.2.3 启动检测 — Start()

文件: manager.go L74-L91

true

false

失败

成功

Start 入口

isRunning?

日志: 已在运行
return

日志: starting...

scanAllContainers
全量扫描所有容器

日志错误, return

启动 watchContainerEvents
goroutine 监听容器事件

isRunning = true

日志: started successfully

func (m *Manager) Start(ctx context.Context) {
    if m.isRunning {
        return  // 防止重复启动
    }
    
    // 步骤1: 全量扫描 — 获取当前所有容器的NPU设备挂载信息
    if err := m.scanAllContainers(ctx); err != nil {
        return  // 初始化失败则不启动
    }
    
    // 步骤2: 启动事件监听 — 持续监控容器创建/销毁事件
    go m.watchContainerEvents(ctx)
    
    m.isRunning = true
}
3.2.4 全量扫描 — scanAllContainers()

文件: manager.go L101-L112

失败

成功

scanAllContainers

client.ParseAllContainers
获取所有容器NPU信息

返回错误

cache.StoreAllAndFindDuplicates
全量存储并检测重复

获取重复列表 duplicates

duplicates非空?

遍历重复项
logDuplicate逐条告警

日志: 0 duplicates

返回nil

func (m *Manager) scanAllContainers(ctx context.Context) error {
    // 调用容器运行时客户端,解析所有容器的NPU设备信息
    // 返回: map[containerID]*ContainerNPUInfo
    result, err := m.client.ParseAllContainers(ctx)
    if err != nil {
        return fmt.Errorf("failed to parse containers: %w", err)
    }
    
    // 将所有容器信息存入缓存,并查找重复挂载
    duplicates := m.cache.StoreAllAndFindDuplicates(result)
    
    // 记录所有已存在的重复挂载
    for _, dup := range duplicates {
        m.logDuplicate(dup)
    }
    
    return nil
}
3.2.5 容器事件监听 — watchContainerEvents()

文件: manager.go L93-L99

ContainerEventCreate

ContainerEventDestroy

ctx.Done

未知事件类型

watchContainerEvents
goroutine

client.WatchContainerEvents
注册回调函数

等待事件

HandleNewContainer
解析新容器NPU信息

StoreSingleAndFindDuplicates
增量存储并检测重复

有重复?

logDuplicate告警

HandleContainerRemoval
从缓存移除容器

停止监听

日志警告

func (m *Manager) watchContainerEvents(ctx context.Context) {
    m.client.WatchContainerEvents(ctx, func(event types.ContainerEvent) {
        switch event.Type {
        case types.ContainerEventCreate:
            // 新容器创建: 解析设备信息并检测重复
            if err := m.HandleNewContainer(ctx, event.ContainerID, event.Namespace); err != nil {
                hwlog.RunLog.Warnf("failed to handle new container %s: %v", event.ContainerID, err)
            }
        case types.ContainerEventDestroy:
            // 容器销毁: 从缓存中移除
            m.HandleContainerRemoval(event.ContainerID)
        }
    })
}
3.2.6 新容器处理 — HandleNewContainer()

文件: manager.go L114-L138

成功

失败

HandleNewContainer
containerID, namespace

设置containerd namespace上下文

重试循环 maxRetries=5

client.ParseSingleContainer
解析容器NPU信息

break 跳出重试

i < maxRetries-1?

Sleep 100ms

返回错误
超过最大重试次数

info为nil或
Devices为空?

返回nil
无NPU设备

设置namespace

cache.StoreSingleAndFindDuplicates
增量存储并检测重复

有重复?

遍历重复项
logDuplicate告警

返回nil

func (m *Manager) HandleNewContainer(ctx context.Context, containerID string, namespace string) error {
    const maxRetries = 5
    const retryDelay = 100 * time.Millisecond
    
    // 设置containerd namespace
    // containerd使用namespace隔离容器(如k8s使用"k8s.io",docker使用"moby")
    ctx = namespaces.WithNamespace(ctx, namespace)
    
    var info *types.ContainerNPUInfo
    var err error
    
    // 重试机制: 新创建的容器可能需要短暂时间才能被查询到
    // 例如: Docker返回start事件后,容器可能还未完全就绪
    for i := 0; i < maxRetries; i++ {
        info, err = m.client.ParseSingleContainer(ctx, containerID)
        if err == nil {
            break
        }
        if i < maxRetries-1 {
            time.Sleep(retryDelay)  // 100ms间隔重试
        }
    }
    
    if err != nil {
        return fmt.Errorf("failed to parse container %s after %d retries: %w", containerID, maxRetries, err)
    }
    
    // 无NPU设备的容器不需要处理
    if info == nil || len(info.Devices) == 0 {
        return nil
    }
    
    info.Namespace = namespace
    
    // 增量存储并检测重复
    for _, dup := range m.cache.StoreSingleAndFindDuplicates(info) {
        m.logDuplicate(dup)
    }
    
    return nil
}
3.2.7 重复挂载告警 — logDuplicate()

文件: manager.go L150-L162

func (m *Manager) logDuplicate(dup *types.DuplicateMountInfo) {
    const maxPrintLength = 12  // 容器ID截断长度,避免日志过长
    
    var containerInfos []string
    for _, c := range dup.Containers {
        // 构建容器信息字符串: ID(前12位) + Name + Namespace
        info := fmt.Sprintf("ID=%s, Name=%s, Namespace=%s",
            c.ID[:min(maxPrintLength, len(c.ID))],  // 类似Docker的短ID显示
            c.Name, c.Namespace)
        
        // 如果是K8s Pod,附加Pod信息
        if c.PodName != "" {
            info += fmt.Sprintf(", Pod=%s/%s", c.PodNS, c.PodName)
        }
        containerInfos = append(containerInfos, info)
    }
    
    // 输出告警日志: 哪个设备被哪些容器同时挂载
    hwlog.RunLog.Warnf("detected duplicate NPU device mount: device /dev/davinci%d is mounted by multiple containers: %s",
        dup.DeviceID, strings.Join(containerInfos, "; "))
}
3.2.8 ContainerCache 缓存机制

文件: cache/container_cache.go

ContainerCache

-containers: map[string]*ContainerNPUInfo

-deviceMap: map[int][]string

-mutex: sync.RWMutex

双向索引结构:\ncontainers: ContainerID → ContainerNPUInfo\ndeviceMap: DeviceID → []ContainerID\n\n两个索引保持同步:\n- 存储容器时更新两个索引\n- 删除容器时清理两个索引

数据结构设计

type ContainerCache struct {
    // 正向索引: 容器ID → 容器NPU信息
    // 用于: 快速查找容器信息、删除时获取设备列表
    containers map[string]*types.ContainerNPUInfo
    
    // 反向索引: 设备ID → 使用该设备的容器ID列表
    // 用于: 快速检测重复挂载(列表长度>1即重复)
    deviceMap map[int][]string
    
    // 读写锁保护并发访问
    // StoreAll/StoreSingle/Remove 使用写锁
    // findDuplicates 在 StoreAll 内部调用,已持有写锁
    mutex sync.RWMutex
}
3.2.8.1 全量存储与重复检测 — StoreAllAndFindDuplicates()

文件: cache/container_cache.go L26-L41

func (cc *ContainerCache) StoreAllAndFindDuplicates(infos map[string]*types.ContainerNPUInfo) []*types.DuplicateMountInfo {
    cc.mutex.Lock()
    defer cc.mutex.Unlock()
    
    // 全量替换: 直接用新数据覆盖旧数据
    cc.containers = infos
    
    // 重建反向索引: 遍历所有容器,为每个设备建立容器列表
    for _, info := range infos {
        for _, deviceID := range info.Devices {
            cc.deviceMap[deviceID] = append(cc.deviceMap[deviceID], info.ID)
        }
    }
    
    // 查找重复: deviceMap中容器数量>1的设备即为重复挂载
    return cc.findDuplicates()
}
// findDuplicates: 遍历deviceMap,找出被多个容器使用的设备
func (cc *ContainerCache) findDuplicates() []*types.DuplicateMountInfo {
    var duplicates []*types.DuplicateMountInfo
    for deviceID, containerIDs := range cc.deviceMap {
        // 日志记录每个设备的容器数量(便于调试)
        hwlog.RunLog.Infof("checking device %d, containers: %d", deviceID, len(containerIDs))
        
        if len(containerIDs) == 1 {
            continue  // 只有一个容器使用,正常
        }
        
        // 多个容器使用同一设备: 收集所有容器信息
        containers := make([]*types.ContainerNPUInfo, 0, len(containerIDs))
        for _, id := range containerIDs {
            if info, ok := cc.containers[id]; ok {
                containers = append(containers, info)
            }
        }
        
        duplicates = append(duplicates, &types.DuplicateMountInfo{
            DeviceID:   deviceID,
            Containers: containers,
        })
    }
    return duplicates
}
3.2.8.2 增量存储与重复检测 — StoreSingleAndFindDuplicates()

文件: cache/container_cache.go L43-L73

StoreSingleAndFindDuplicates
info

加锁

遍历info.Devices

deviceMap中
已有该设备?

收集已存在的容器信息

将当前容器加入重复列表

构建DuplicateMountInfo

追加到duplicates

无重复

遍历完所有设备?

将容器加入containers缓存

将容器ID加入deviceMap
所有设备索引

返回duplicates

func (cc *ContainerCache) StoreSingleAndFindDuplicates(info *types.ContainerNPUInfo) []*types.DuplicateMountInfo {
    cc.mutex.Lock()
    defer cc.mutex.Unlock()
    
    var duplicates []*types.DuplicateMountInfo
    
    // 步骤1: 检查新容器的每个设备是否已被其他容器使用
    for _, deviceID := range info.Devices {
        existingContainers, ok := cc.deviceMap[deviceID]
        if !ok || len(existingContainers) == 0 {
            continue  // 设备未被使用,无重复
        }
        
        // 设备已被其他容器使用: 收集所有相关容器信息
        var dupContainers []*types.ContainerNPUInfo
        for _, existingID := range existingContainers {
            if existingContainer, ok := cc.containers[existingID]; ok {
                dupContainers = append(dupContainers, existingContainer)
            }
        }
        // 将当前容器也加入重复列表
        dupContainers = append(dupContainers, info)
        
        duplicates = append(duplicates, &types.DuplicateMountInfo{
            DeviceID:   deviceID,
            Containers: dupContainers,
        })
    }
    
    // 步骤2: 将新容器存入缓存(即使检测到重复也要存储,保持缓存完整)
    cc.containers[info.ID] = info
    for _, deviceID := range info.Devices {
        cc.deviceMap[deviceID] = append(cc.deviceMap[deviceID], info.ID)
    }
    
    return duplicates
}

设计要点:即使检测到重复挂载,仍然将容器信息存入缓存。这确保后续如果重复容器被销毁,缓存能正确更新。

3.2.8.3 容器移除 — RemoveContainer()

文件: cache/container_cache.go L75-L96

func (cc *ContainerCache) RemoveContainer(containerID string) {
    cc.mutex.Lock()
    defer cc.mutex.Unlock()
    
    // 获取容器信息(需要知道该容器使用了哪些设备)
    info, ok := cc.containers[containerID]
    if !ok {
        return  // 容器不在缓存中,无需处理
    }
    
    // 从正向索引删除
    delete(cc.containers, containerID)
    
    // 从反向索引删除: 遍历该容器使用的所有设备
    for _, deviceID := range info.Devices {
        containers := cc.deviceMap[deviceID]
        // 从容器列表中移除该容器ID(保持列表顺序)
        newList := make([]string, 0, len(containers))
        for _, id := range containers {
            if id != containerID {
                newList = append(newList, id)
            }
        }
        cc.deviceMap[deviceID] = newList
    }
}
3.2.9 容器运行时客户端
3.2.9.1 客户端工厂 — NewClient()

文件: containerruntime/interface.go L58-L73

docker

containerd

其他

NewClient
config

config为nil?

返回错误

autoDetectOciEndpoint
自动探测OCI socket

检查路径存在?
/run/containerd/containerd.sock

返回默认路径

检查路径存在?
/run/docker/containerd/containerd.sock

返回Docker containerd路径

返回错误: 探测失败

config.RuntimeType?

NewDockerClient
config.CriEndpoint, ociEndpoint

NewContainerdClient
config.CriEndpoint, ociEndpoint

返回错误: 不支持的运行时

返回Client

func NewClient(config *types.DetectorConfig) (Client, error) {
    if config == nil {
        return nil, fmt.Errorf("config is nil")
    }
    
    // 自动探测OCI (Open Container Initiative) socket路径
    // OCI是容器运行时的底层标准接口
    // Docker底层也使用containerd作为OCI运行时
    ociEndpoint, err := autoDetectOciEndpoint()
    if err != nil {
        return nil, err
    }
    
    // 根据运行时类型创建对应客户端
    if config.RuntimeType == kubeclient.DockerRuntime {
        return NewDockerClient(config.CriEndpoint, ociEndpoint)
    }
    if config.RuntimeType == kubeclient.ContainerdRuntime {
        return NewContainerdClient(config.CriEndpoint, ociEndpoint)
    }
    
    return nil, fmt.Errorf("runtime type %s is not supported", config.RuntimeType)
}
func autoDetectOciEndpoint() (string, error) {
    // 优先检查K8s标准containerd路径
    if _, err := os.Stat(defaultContainerdAddr); err == nil {
        return defaultContainerdAddr, nil  // "/run/containerd/containerd.sock"
    }
    
    // 其次检查Docker内嵌containerd路径
    if _, err := os.Stat(dockerContainerdAddr); err == nil {
        return dockerContainerdAddr, nil  // "/run/docker/containerd/containerd.sock"
    }
    
    return "", errors.New("failed to auto-detect oci socket path")
}
3.2.9.2 OCI客户端 — 容器设备信息解析

文件: containerruntime/interface.go L40-L56

失败

成功

失败

成功

失败

成功

ParseSingleContainer
containerID

TaskService.Get
获取容器Task

返回错误

Process为nil?

返回错误: task未找到

LoadContainer
加载容器对象

返回错误

获取容器Spec
oci.Spec

返回错误

获取容器Labels
K8s元数据

构建ContainerNPUInfo
ID, Name, PodName, PodNS

从Spec.Process.Env
解析AscendDeviceInfo环境变量

找到
AscendDeviceInfo?

parser.ParseAscendDeviceInfo
解析设备ID列表

parser.FilterNPUDevices
从Spec的Mounts过滤NPU设备

Devices非空?

返回info

返回info
Devices为空

func (c *ociClient) ParseSingleContainer(ctx context.Context, containerID string) (*types.ContainerNPUInfo, error) {
    // 步骤1: 通过TaskService获取容器Task信息
    // Task是containerd中对容器运行实例的抽象
    task, err := c.client.TaskService().Get(ctx, &tasks.GetRequest{ContainerID: containerID})
    if err != nil {
        return nil, err
    }
    
    // 验证Task的Process存在
    if task.GetProcess() == nil {
        return nil, fmt.Errorf("task not found for container %s", containerID)
    }
    
    // 步骤2: 加载容器对象(获取Spec和Labels)
    ctr, err := c.client.LoadContainer(ctx, containerID)
    if err != nil {
        return nil, err
    }
    
    // 步骤3: 获取OCI Spec(包含容器配置:环境变量、挂载点等)
    spec, err := ctr.Spec(ctx)
    if err != nil || spec == nil {
        return nil, fmt.Errorf("failed to get container spec: %w", err)
    }
    
    // 步骤4: 获取容器Labels(K8s注入的元数据)
    labels, err := ctr.Labels(ctx)
    if err != nil {
        return nil, fmt.Errorf("failed to get container labels: %w", err)
    }
    
    // 步骤5: 构建NPU信息对象
    info := &types.ContainerNPUInfo{
        ID:      ctr.ID(),
        Name:    labels["io.kubernetes.container.name"],
        PodName: labels["io.kubernetes.pod.name"],
        PodNS:   labels["io.kubernetes.pod.namespace"],
    }
    
    // 步骤6: 解析NPU设备 — 优先从环境变量中提取
    if spec.Process != nil {
        // 逆序遍历环境变量(后定义的优先级更高)
        for i := len(spec.Process.Env) - 1; i >= 0; i-- {
            env := strings.TrimSpace(spec.Process.Env[i])
            // 查找包含 AscendDeviceInfo 的环境变量
            // 这是华为设备插件注入的专用环境变量
            if strings.Contains(env, api.AscendDeviceInfo) {
                info.Devices = parser.ParseAscendDeviceInfo(env, ctr.ID())
                break
            }
        }
    }
    
    // 步骤7: 如果环境变量中没有设备信息,从Mounts中过滤
    if len(info.Devices) != 0 {
        return info, nil
    }
    
    // FilterNPUDevices: 从Spec的Mounts中过滤 /dev/davinci* 设备
    info.Devices = parser.FilterNPUDevices(spec)
    return info, nil
}

设备解析的双重策略

  1. 优先策略:从环境变量 AscendDeviceInfo 中解析 — 这是设备插件在创建容器时注入的结构化设备信息,解析效率高
  2. 降级策略:从 OCI Spec 的 Mounts 列表中过滤 /dev/davinci* 设备 — 通用方法,适用于环境变量缺失的场景
3.2.9.3 Docker客户端

文件: containerruntime/docker_client.go

dockerClient

-client: *client.Client

-ociClient: *ociClient

+ParseAllContainers(ctx)(map, error)

+ParseSingleContainer(ctx, containerID)(*ContainerNPUInfo, error)

+WatchContainerEvents(ctx, handler) : void

Docker客户端采用组合模式:\n- 使用Docker API获取容器列表和事件\n- 使用底层containerd(ociClient)解析容器Spec\n- 这是因为Docker的ContainerInspect不暴露\n 完整的OCI Spec环境变量

Docker客户端初始化:

func NewDockerClient(criEndpoint string, ociEndpoint string) (*dockerClient, error) {
    // CRI端点默认: unix:///run/docker.sock
    if criEndpoint == "" {
        criEndpoint = defaultDockerAddress
    }
    
    // 安全检查: 验证socket文件路径和权限
    // excludePermissions = 0002: 拒绝world-writable权限
    if err := checkSockFile(criEndpoint); err != nil {
        return nil, fmt.Errorf("invalid cri endpoint(%s): %v", criEndpoint, err)
    }
    if err := checkSockFile(ociEndpoint); err != nil {
        return nil, fmt.Errorf("invalid oci endpoint(%s): %v", ociEndpoint, err)
    }
    
    // 创建Docker CLI客户端
    // WithAPIVersionNegotiation: 自动协商API版本
    cli, err := client.NewClientWithOpts(
        client.WithHost(criEndpoint),
        client.WithAPIVersionNegotiation(),
    )
    if err != nil {
        return nil, err
    }
    
    // 同时创建containerd客户端(用于解析OCI Spec)
    ctrClient, err := containerd.New(ociEndpoint)
    if err != nil {
        return nil, err
    }
    
    return &dockerClient{
        client: cli,
        ociClient: &ociClient{client: ctrClient},
    }, nil
}

Docker事件监听:

func (d *dockerClient) WatchContainerEvents(ctx context.Context, handler dtypes.EventHandler) {
    // 构建事件过滤器: 只关注容器的start和die事件
    filterArgs := filters.NewArgs()
    filterArgs.Add("type", "container")
    filterArgs.Add("event", "start")
    filterArgs.Add("event", "die")
    
    // 订阅Docker事件流
    eventChan, errChan := d.client.Events(ctx, types.EventsOptions{
        Filters: filterArgs,
    })
    
    for {
        select {
        case <-ctx.Done():
            d.client.Close()
            return
            
        case event := <-eventChan:
            switch event.Action {
            case "start":
                handler(dtypes.ContainerEvent{
                    Type:        dtypes.ContainerEventCreate,
                    ContainerID: event.Actor.ID,
                    Namespace:   dockerNamespace,  // "moby"
                    Timestamp:   time.Now(),
                })
            case "die":
                handler(dtypes.ContainerEvent{
                    Type:        dtypes.ContainerEventDestroy,
                    ContainerID: event.Actor.ID,
                    Namespace:   dockerNamespace,
                    Timestamp:   time.Now(),
                })
            }
            
        case err := <-errChan:
            hwlog.RunLog.Errorf("error receiving event: %v", err)
        }
    }
}

Docker全量容器解析:

func (d *dockerClient) ParseAllContainers(ctx context.Context) (map[string]*dtypes.ContainerNPUInfo, error) {
    // 获取所有运行中的容器
    ctrs, err := d.client.ContainerList(ctx, types.ContainerListOptions{})
    if err != nil {
        return nil, fmt.Errorf("failed to list containers: %w", err)
    }
    
    containerInfos := make(map[string]*dtypes.ContainerNPUInfo)
    
    // Docker使用"moby"作为containerd namespace
    nsCtx := namespaces.WithNamespace(ctx, dockerNamespace)
    
    for _, ctr := range ctrs {
        // 使用ociClient解析容器Spec获取NPU设备信息
        info, err := d.ParseSingleContainer(nsCtx, ctr.ID)
        if err != nil {
            // 单个容器解析失败不影响整体
            continue
        }
        
        // 从Docker Labels补充K8s元数据
        info.PodName = ctr.Labels["io.kubernetes.pod.name"]
        info.PodNS = ctr.Labels["io.kubernetes.pod.namespace"]
        info.Namespace = dockerNamespace
        info.Name = ctr.Labels["io.kubernetes.container.name"]
        
        containerInfos[ctr.ID] = info
    }
    
    return containerInfos, nil
}
3.2.9.4 containerd客户端

文件: containerruntime/containerd_client.go

containerdClient

-ociClient: *ociClient

+ParseAllContainers(ctx)(map, error)

+WatchContainerEvents(ctx, handler) : void

-handleEvent(envelope, handler) : void

containerd客户端采用嵌入模式:\n- 直接嵌入ociClient\n- containerd原生支持OCI Spec\n- 无需额外的Docker API调用

containerd全量容器解析:

func (c *containerdClient) ParseAllContainers(ctx context.Context) (map[string]*types.ContainerNPUInfo, error) {
    // containerd支持多namespace,需要遍历所有namespace
    nss, err := c.client.NamespaceService().List(ctx)
    if err != nil {
        return nil, fmt.Errorf("failed to list containers: %w", err)
    }
    
    containerInfos := make(map[string]*types.ContainerNPUInfo)
    
    // 遍历所有namespace(如k8s.io、moby等)
    for _, ns := range nss {
        nsCtx := namespaces.WithNamespace(ctx, ns)
        
        // 获取该namespace下的所有容器
        containers, err := c.client.Containers(nsCtx)
        if err != nil {
            continue  // 单个namespace失败不影响整体
        }
        
        for _, ctr := range containers {
            info, err := c.ociClient.ParseSingleContainer(nsCtx, ctr.ID())
            if err != nil {
                continue
            }
            info.Namespace = ns
            containerInfos[ctr.ID] = info
        }
    }
    
    return containerInfos, nil
}

containerd事件监听:

func (c *containerdClient) WatchContainerEvents(ctx context.Context, handler types.EventHandler) {
    // 订阅containerd事件,使用topic过滤器
    // ~=/tasks/start: 容器Task启动
    // ~=/tasks/exit: 容器Task退出
    eventChan, errChan := c.client.EventService().Subscribe(ctx,
        `topic~="/tasks/start"`,
        `topic~="/tasks/exit"`,
    )
    
    for {
        select {
        case <-ctx.Done():
            c.client.Close()
            return
            
        case envelope := <-eventChan:
            c.handleEvent(envelope, handler)
            
        case err := <-errChan:
            hwlog.RunLog.Warnf("error receiving event: %v", err)
        }
    }
}
func (c *containerdClient) handleEvent(envelope *events.Envelope, handler types.EventHandler) {
    if envelope.Event == nil {
        return
    }
    
    // 使用typeurl反序列化事件
    v, err := typeurl.UnmarshalAny(envelope.Event)
    if err != nil {
        return
    }
    
    switch event := v.(type) {
    case *apievents.TaskStart:
        // 容器启动事件
        handler(types.ContainerEvent{
            Type:        types.ContainerEventCreate,
            ContainerID: event.ContainerID,
            Namespace:   envelope.Namespace,
            Timestamp:   time.Now(),
        })
        
    case *apievents.TaskExit:
        // 容器退出事件
        // 关键判断: event.ContainerID != event.ID
        // TaskExit事件中:
        //   - ContainerID 是容器ID
        //   - ID 是进程ID(exec进程也会有TaskExit事件)
        // 只有当两者相等时,才是容器主进程退出(即容器销毁)
        if event.ContainerID != event.ID {
            return  // exec进程退出,忽略
        }
        handler(types.ContainerEvent{
            Type:        types.ContainerEventDestroy,
            ContainerID: event.ContainerID,
            Namespace:   envelope.Namespace,
            Timestamp:   time.Now(),
        })
        
    default:
        hwlog.RunLog.Warnf("unknown event type: %T", event)
    }
}

设计要点event.ContainerID != event.ID 的判断非常关键。containerd 中,容器内通过 exec 创建的额外进程也会有 TaskExit 事件,但此时 event.ID(进程ID)不等于 event.ContainerID(容器ID)。只有主进程退出才表示容器终止。

3.2.10 Socket文件安全检查 — checkSockFile()

文件: containerruntime/docker_client.go L48-L56

const (
    excludePermissions = 0002  // 世界可写权限位
    unixPre            = "unix://"
)

func checkSockFile(path string) error {
    // 去除 unix:// 前缀,获取文件系统路径
    absPath, err := utils.CheckPath(strings.TrimPrefix(path, unixPre))
    if err != nil {
        return err
    }
    
    // 检查socket文件的所有者和权限
    // excludePermissions = 0002: 拒绝world-writable
    // 第二个参数 0 表示要求root所有(uid=0)
    return utils.DoCheckOwnerAndPermission(absPath, excludePermissions, 0)
}

安全意图

  • 0002 权限位表示"其他用户可写",这是一个安全风险
  • 如果socket文件是world-writable,任何用户都可以连接该socket并操作容器
  • 要求root所有(uid=0)确保socket文件由特权用户创建
3.2.11 类型定义

文件: duplicatedetector/types/types.go

// ContainerNPUInfo: 容器NPU设备信息
type ContainerNPUInfo struct {
    ID        string // 容器ID(Docker或containerd的容器ID)
    Name      string // 容器名称(K8s容器名,如"nginx")
    Namespace string // 容器运行时namespace(如"moby"、"k8s.io")
    PodName   string // K8s Pod名称
    PodNS     string // K8s Pod命名空间
    Devices   []int  // NPU设备ID列表(如[0, 1]表示使用davinci0和davinci1)
}

// DuplicateMountInfo: 重复挂载信息
type DuplicateMountInfo struct {
    DeviceID   int                   // 被重复挂载的设备ID
    Containers []*ContainerNPUInfo   // 挂载该设备的所有容器列表
}

// DetectorConfig: 检测器配置
type DetectorConfig struct {
    CriEndpoint string   // CRI socket端点
    RuntimeType string   // 运行时类型: "docker" 或 "containerd"
}

// ContainerEvent: 容器生命周期事件
type ContainerEvent struct {
    Type        ContainerEventType  // 事件类型: create/destroy
    ContainerID string              // 容器ID
    Namespace   string              // 运行时namespace
    Timestamp   time.Time           // 事件时间戳
}
3.2.12 V2类型定义

文件: kubeclient/types_v2.go

// HccspingMeshItem: PingMesh配置项
// PingMesh是华为集群网络检测组件,用于检测节点间网络连通性
type HccspingMeshItem struct {
    Activate     string `json:"activate"`       // 开关: "on"/"off"
    TaskInterval int    `json:"task_interval"`  // 检测间隔(秒)
}

// ConfigPingMesh: PingMesh配置Map
// key是配置项名称,value是配置内容
type ConfigPingMesh map[string]*HccspingMeshItem
3.2.13 事件与资源类型定义

文件: kubeclient/types.go

// Event: 工作队列中的事件
type Event struct {
    Resource ResourceType  // 资源类型: Pod/ConfigMap
    Key      string        // 资源唯一键: namespace/name
    Type     EventType     // 事件类型: add/update/delete
}

// 事件类型枚举
const (
    EventTypeAdd    EventType = "add"
    EventTypeUpdate EventType = "update"
    EventTypeDelete EventType = "delete"
)

// 资源类型枚举
const (
    PodResource ResourceType = "pod"
    CMResource  ResourceType = "configmap"
)

// 容器运行时类型常量
const (
    DockerRuntime     = "docker"
    ContainerdRuntime = "containerd"
)

四、模块间交互全景图

Docker/Containerd Kubelet K8s API Server ContainerCache containerruntime duplicatedetector kubeclient main.go Docker/Containerd Kubelet K8s API Server ContainerCache containerruntime duplicatedetector kubeclient main.go 启动阶段 每10分钟巡检缓存 运行阶段 loop [事件循环] 运行阶段 - 业务调用 alt [IsApiErr=true] loop [业务请求] NewClientK8s() BuildConfigFromFlags (InClusterConfig) REST Config + Token ClientK8s实例 InitPodInformer() Watch Pods (FieldSelector: spec.nodeName) Pod事件流 PodInformerInspector() GetContainerRuntime() GetNode() Node.Status.NodeInfo.ContainerRuntimeVersion "docker" 或 "containerd" CheckDuplicateDevices(ctx, config) NewClient(config) 检测socket文件 socket路径 Client实例 (Docker或Containerd) ParseAllContainers() List所有容器 容器列表 LoadContainer + Spec OCI Spec (含环境变量/Mounts) map[containerID]*ContainerNPUInfo StoreAllAndFindDuplicates() 构建双向索引 []*DuplicateMountInfo logDuplicate告警 WatchContainerEvents() Subscribe events start/die 事件 ContainerEvent{Create} ParseSingleContainer() ContainerNPUInfo StoreSingleAndFindDuplicates() duplicates logDuplicate告警 ContainerEvent{Destroy} RemoveContainer() GetActivePodListCache() 检查IsApiErr GetAllPodList() PodList refreshPodList() []v1.Pod GetPodsUsedNPUByKlt() HTTPS GET /pods (Bearer Token) PodList (JSON) sets.String (NPU ID集合) WriteDeviceInfoDataIntoCM() UpdateConfigMap() ConfigMap

五、关键设计模式与架构决策总结

5.1 设计模式

模式 应用位置 说明
单例模式 duplicatedetector/manager.go sync.Once 确保 Manager 全局唯一
工厂模式 containerruntime/interface.go NewClient() 根据运行时类型创建不同客户端
策略模式 containerruntime.Client 接口 Docker 和 containerd 两种实现可互换
观察者模式 kubeclient/ResourceEventHandler Informer 事件订阅与处理
缓存模式 kube_cache.go 中的 podCache 多级缓存减少 API Server 压力
组合模式 dockerClient 内嵌 ociClient Docker 客户端复用 containerd 的 Spec 解析能力
降级模式 GetPodsUsedNPUByKlt 降级到 GetPodsUsedNpuByCommon Kubelet 不可用时降级到缓存
双重索引 ContainerCachecontainers + deviceMap 正反向索引加速重复检测

5.2 架构决策

关键架构决策

三级Pod数据获取策略
Kubelet直连 → API缓存 → API直查
确保数据新鲜度与可用性

IsApiErr自恢复机制
API错误标记 → 下次刷新 → 重置标记
无需人工干预

Predicate-Time校正机制
缓存值校正API数据
防止设备重复分配

Pod缓存随机巡检
FNV哈希分散启动时间
避免惊群效应

容器运行时双重设备解析
环境变量优先 → Mounts降级
确保设备信息获取成功率

Docker客户端组合ociClient
Docker API获取列表/事件
containerd API解析Spec

containerd TaskExit过滤
ContainerID == ID判断
排除exec进程退出事件

Socket文件安全检查
拒绝world-writable + 要求root所有
防止权限提升攻击

TLS 1.3 + ECDHE套件
Kubelet通信强制强加密
即使跳过证书验证也保证机密性

ConfigMap数据按型号差异化
A5/A3/默认三种数据结构
适配不同硬件能力

重置信息隔离错误检查
IsolateError + 设备列表
阻止错误状态下继续操作

5.3 数据流总结

输出

缓存层

处理层

输入源

K8s API Server
Node/Pod/CM/Event

Kubelet :10250
/pods 直连

Docker Socket
/run/docker.sock

Containerd Socket
/run/containerd/containerd.sock

环境变量
NODE_NAME/HOST_IP/KUBELET_PORT

kubeclient
K8s交互封装

duplicatedetector
重复检测

podCache
map[UID]*podInfo

nodeServerIp
节点IP缓存

serverUsageLabel
节点标签缓存

nodeDeviceInfoCache
设备信息缓存

ContainerCache
容器NPU缓存

Node Annotation/Status
节点注解与状态更新

Pod Annotation
Pod注解更新

ConfigMap
设备/重置/故障信息

K8s Event
事件资源

告警日志
重复挂载告警


六、文件清单与依赖关系

外部依赖

duplicatedetector 包文件

kubeclient 包文件

kubeclient.go
452行
核心客户端

kube_connect.go
158行
Kubelet直连

kube_cache.go
301行
Pod缓存管理

cur_node_informer.go
70行
Informer初始化

client_server.go
450行
CM写入/Pod注解

kubeclient_v2.go
60行
V2扩展功能

types.go
52行
事件/资源类型

types_v2.go
27行
PingMesh类型

manager.go
184行
检测管理器

cache/container_cache.go
139行
容器缓存

containerruntime/interface.go
134行
运行时接口

containerruntime/docker_client.go
171行
Docker实现

containerruntime/containerd_client.go
143行
containerd实现

types/types.go
66行
检测器类型

k8s.io/client-go
K8s客户端库

github.com/docker/docker
Docker SDK

github.com/containerd/containerd
containerd SDK

Ascend-device-plugin/pkg/common
公共工具包

ascend-common
华为公共库


分析完成。本文档严格基于源码逐行分析,所有代码引用均来自实际源文件。

Logo

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

更多推荐