Part 6: K8s客户端 + 重复检测模块 超深度源码分析

分析范围: pkg/kubeclient/ (9个文件) + pkg/duplicatedetector/ (5个文件)
总代码行数: ~2,473 行
模块定位: Kubernetes API 交互层 + 容器运行时 NPU 设备重复挂载检测层


一、模块定位

1.1 业务职责与功能定位

本模块由两个子模块构成,分别承担不同的业务职责:

子模块 业务职责 功能定位
kubeclient Kubernetes API 客户端封装 为整个 ascended-device-plugin 提供统一的 K8s 资源操作入口,包括 Node/Pod/ConfigMap/Event 的 CRUD 操作、Pod 缓存管理、Informer 监听、Kubelet 直连查询等
duplicatedetector NPU 设备重复挂载检测 运行时检测同一物理 NPU 设备被多个容器同时挂载的异常场景,支持 Docker 和 containerd 两种容器运行时

1.2 在系统中的位置

ascend-device-plugin 系统架构

基础设施层

安全检测层 ← 本模块

K8s交互层 ← 本模块

业务逻辑层

入口层

REST API

HTTPS /pods

gRPC/CLI

启动时调用

K8s API Server

Device Plugin Server
gRPC Interface

Device Manager
设备管理

kubeclient
K8s客户端封装

Fault Manager
故障管理

BurnIn Manager
老化测试

Other Managers

Kubelet

duplicatedetector
重复挂载检测

Container Runtime
Docker/Containerd

核心定位说明

  • kubeclient 是所有上层业务 Manager 与 K8s 交互的唯一通道,封装了 API Server 调用、Kubelet 直连、Informer 缓存三种数据获取方式
  • duplicatedetector 作为安全防护组件,在设备插件启动时初始化,持续监控容器级别的 NPU 设备挂载情况,防止设备重复分配

二、模块整体结构

2.1 kubeclient 子模块结构

Queue中存储

缓存管理

ClientK8s

+Clientset kubernetes.Interface

+NodeName string

+DeviceInfoName string

+IsApiErr bool

+PodInformer cache.SharedIndexInformer

+Queue workqueue.RateLimitingInterface

+KltClient *http.Client

+GetNode()(*v1.Node, error)

+PatchNodeState(curNode, newNode)(*v1.Node, []byte, error)

+AddAnnotation(key, value) : error

+GetPod(pod)(*v1.Pod, error)

+PatchPod(pod, data)(*v1.Pod, error)

+GetActivePodList()([]v1.Pod, error)

+GetAllPodList()(*v1.PodList, error)

+GetContainerRuntime()(string, error)

+CreateConfigMap(cm)(*v1.ConfigMap, error)

+GetConfigMap(cmName, cmNameSpace)(*v1.ConfigMap, error)

+UpdateConfigMap(cm)(*v1.ConfigMap, error)

+ResetDeviceInfo() : void

+ClearResetInfo(taskName, namespace) : error

+CreateEvent(evt)(*v1.Event, error)

+ResourceEventHandler(res, filter) : cache.ResourceEventHandler

+FlushPodCacheNextQuerying() : void

+RemoveOldResource(oldResourceName) : error

+GetActivePodListCache() : []v1.Pod

+GetAllPodListCache() : []v1.Pod

+GetNodeIpCache()(string, error)

+GetServerUsageLabelCache()(string, error)

+GetDeviceInfoCMCache() : *NodeDeviceInfoCache

+WriteDeviceInfoDataIntoCMCache(...) : error

+TryUpdatePodAnnotation(pod, annotation) : error

+TryUpdatePodCacheAnnotation(pod, annotation) : error

+WriteDeviceInfoDataIntoCM(...)(*NodeDeviceInfoCache, error)

+WriteResetInfoDataIntoCM(...)(*v1.ConfigMap, error)

+WriteFaultInfoDataIntoCM(...)(*v1.ConfigMap, error)

+AnnotationReset() : error

+GetPodsUsedNpuByCommon() : sets.String

+GetPodsUsedNPUByKlt() : sets.String

+GetNodeIp()(string, error)

+InitPodInformer() : void

+PodInformerInspector(ctx) : void

+UpdatePodList(newObj, operator) : void

Event

+ResourceType Resource

+Key string

+Type EventType

«enumeration»

EventType

+EventTypeAdd

+EventTypeUpdate

+EventTypeDelete

«enumeration»

ResourceType

+PodResource

+CMResource

HccspingMeshItem

+Activate string

+TaskInterval int

ConfigPingMesh

map[string]*HccspingMeshItem

podInfo

+*v1.Pod Pod

+updateTime time.Time

2.2 duplicatedetector 子模块结构

持有运行时客户端

持有缓存

实现

实现

组合复用

嵌入复用

缓存对象

关联容器

配置初始化

事件处理

Manager

-client containerruntime.Client

-cache *cache.ContainerCache

-isRunning bool

+Start(ctx) : void

+HandleNewContainer(ctx, containerID, namespace) : error

+HandleContainerRemoval(containerID) : void

-scanAllContainers(ctx) : error

-watchContainerEvents(ctx) : void

-logDuplicate(dup) : void

ContainerCache

-containers map[string]*ContainerNPUInfo

-deviceMap map[int][]string

-mutex sync.RWMutex

+StoreAllAndFindDuplicates(infos) : []*DuplicateMountInfo

+StoreSingleAndFindDuplicates(info) : []*DuplicateMountInfo

+RemoveContainer(containerID) : void

-findDuplicates() : []*DuplicateMountInfo

«interface»

Client

+ParseAllContainers(ctx)(map[string]*ContainerNPUInfo, error)

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

+WatchContainerEvents(ctx, handler) : void

ociClient

-client *containerd.Client

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

dockerClient

-client *client.Client

-ociClient *ociClient

+ParseAllContainers(ctx)(map[string]*ContainerNPUInfo, error)

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

+WatchContainerEvents(ctx, handler) : void

containerdClient

-ociClient *ociClient

+ParseAllContainers(ctx)(map[string]*ContainerNPUInfo, error)

+WatchContainerEvents(ctx, handler) : void

-handleEvent(envelope, handler) : void

ContainerNPUInfo

+ID string

+Name string

+Namespace string

+PodName string

+PodNS string

+Devices []int

DuplicateMountInfo

+DeviceID int

+Containers []*ContainerNPUInfo

DetectorConfig

+CriEndpoint string

+RuntimeType string

ContainerEvent

+Type ContainerEventType

+ContainerID string

+Namespace string

+Timestamp time.Time

2.3 核心方法清单

kubeclient 核心方法
文件 方法 作用
kubeclient.go NewClientK8s() 创建 K8s 客户端实例,初始化 Clientset、NodeName、Kubelet HTTP 客户端
kubeclient.go GetNode() 获取当前节点信息
kubeclient.go PatchNodeState() 补丁更新节点状态
kubeclient.go AddAnnotation() 为节点添加注解(带重试)
kubeclient.go GetPod() 获取指定 Pod
kubeclient.go PatchPod() 补丁更新 Pod 信息
kubeclient.go GetActivePodList() 获取当前节点上活跃的 Pod 列表
kubeclient.go GetAllPodList() 获取当前节点上所有 Pod 列表
kubeclient.go GetContainerRuntime() 获取节点容器运行时类型
kubeclient.go CreateConfigMap() 创建 ConfigMap
kubeclient.go GetConfigMap() 获取 ConfigMap
kubeclient.go UpdateConfigMap() 更新 ConfigMap
kubeclient.go ResetDeviceInfo() 重置设备信息
kubeclient.go ClearResetInfo() 清除重置信息
kubeclient.go CreateEvent() 创建 K8s Event
kubeclient.go ResourceEventHandler() 创建资源事件处理器
kubeclient.go RemoveOldResource() 移除节点上旧的无用资源
kube_connect.go getPodsByKltPort() 通过 Kubelet 端口直接获取 Pod 列表
kube_connect.go createKltPodsReqWithToken() 创建带 Token 认证的 Kubelet HTTP 请求
kube_connect.go getKltPodsURL() 构建 Kubelet /pods URL
kube_cache.go PodInformerInspector() 定期检查 Pod 缓存有效性
kube_cache.go UpdatePodList() 通过 Informer 更新 Pod 缓存
kube_cache.go GetAllPodListCache() 从缓存获取所有 Pod 列表
kube_cache.go GetActivePodListCache() 从缓存获取活跃 Pod 列表
kube_cache.go GetNodeIpCache() 带缓存的节点 IP 获取
kube_cache.go GetServerUsageLabelCache() 带缓存的节点标签获取
kube_cache.go WriteDeviceInfoDataIntoCMCache() 写设备信息到 CM 并缓存
cur_node_informer.go InitPodInformer() 初始化 Pod Informer
client_server.go TryUpdatePodAnnotation() 更新 Pod 注解(带重试)
client_server.go TryUpdatePodCacheAnnotation() 同时更新 Pod 注解和缓存
client_server.go WriteDeviceInfoDataIntoCM() 写设备信息到 ConfigMap
client_server.go WriteResetInfoDataIntoCM() 写重置信息到 ConfigMap
client_server.go WriteFaultInfoDataIntoCM() 写故障信息到 ConfigMap
client_server.go AnnotationReset() 重置节点注解
client_server.go GetPodsUsedNpuByCommon() 从缓存获取已使用 NPU 列表
client_server.go GetPodsUsedNPUByKlt() 从 Kubelet 获取已使用 NPU 列表
client_server.go GetNodeIp() 获取节点 IP
kubeclient_v2.go ParseFaultNetworkInfoCM() 解析故障网络配置 ConfigMap
kubeclient_v2.go IsConfigMapExists() 判断 ConfigMap 是否存在
duplicatedetector 核心方法
文件 方法 作用
manager.go CheckDuplicateDevices() 模块入口:单例初始化并启动检测
manager.go NewManager() 创建重复检测管理器
manager.go Start() 启动检测:全量扫描 + 事件监听
manager.go HandleNewContainer() 处理新容器创建事件
manager.go HandleContainerRemoval() 处理容器销毁事件
cache/container_cache.go StoreAllAndFindDuplicates() 全量存储并查找重复
cache/container_cache.go StoreSingleAndFindDuplicates() 增量存储并查找重复
cache/container_cache.go RemoveContainer() 从缓存移除容器
containerruntime/interface.go NewClient() 根据运行时类型创建客户端
containerruntime/interface.go autoDetectOciEndpoint() 自动探测 OCI socket 路径
containerruntime/interface.go ParseSingleContainer() 解析单个容器的 NPU 设备信息
docker_client.go NewDockerClient() 创建 Docker 客户端
docker_client.go ParseAllContainers() 列出并解析所有 Docker 容器
docker_client.go WatchContainerEvents() 监听 Docker 容器事件
containerd_client.go NewContainerdClient() 创建 containerd 客户端
containerd_client.go ParseAllContainers() 列出并解析所有 containerd 容器
containerd_client.go WatchContainerEvents() 监听 containerd 容器事件

2.4 内部调用关系

cache 包

containerruntime 包

duplicatedetector 包

kubeclient 包

外部调用方

IsApiErr时刷新

失败降级

不存在则创建

Device Manager

Fault Manager

main.go

NewClientK8s

ClientK8s

InitPodInformer

PodInformerInspector

UpdatePodList

GetAllPodListCache

GetActivePodListCache

WriteDeviceInfoDataIntoCM

WriteResetInfoDataIntoCM

WriteFaultInfoDataIntoCM

TryUpdatePodAnnotation

GetPodsUsedNPUByKlt

GetPodsUsedNpuByCommon

getPodsByKltPort

GetConfigMap

UpdateConfigMap

CreateConfigMap

AnnotationReset

RemoveOldResource

CheckDuplicateDevices

Manager

scanAllContainers

watchContainerEvents

HandleNewContainer

HandleContainerRemoval

NewClient

dockerClient

containerdClient

ociClient

ContainerCache

2.5 数据流入流出方式

输出

duplicatedetector 数据流

kubeclient 数据流

数据源

List/Get/Patch/Create/Update

HTTPS GET /pods

Container List/Inspect/Events

Container List/Spec/Events

NODE_NAME/HOST_IP/KUBELET_PORT

GetAllPodListCache

GetActivePodListCache

WriteDeviceInfo

WriteResetInfo

WriteFaultInfo

PatchNodeState

CreateEvent

findDuplicates

K8s API Server

Kubelet :10250

Docker Socket

Containerd Socket

环境变量

REST API 调用

HTTPS /pods

Informer Watch

内存缓存
podCache
nodeServerIp
nodeDeviceInfoCache

WorkQueue

容器解析

事件流

ContainerCache
containers + deviceMap

日志告警

Node Annotation/Patch

Pod Annotation/Patch

ConfigMap Data

Node Status Update

K8s Event

重复挂载告警日志


三、核心业务逻辑深度解析

3.1 kubeclient 模块

3.1.1 客户端初始化流程 — NewClientK8s()

文件: kubeclient.go L72-L106

失败

成功

失败

成功

失败

成功

NewClientK8s 入口

BuildConfigFromFlags
从K8s配置构建REST配置

日志: build client config err

返回 nil, err

kubernetes.NewForConfig
创建Clientset

日志: get client err

返回 nil, err

GetNodeNameFromEnv
从环境变量获取节点名

返回 nil, err

创建HTTP Transport
TLS配置: 跳过验证
限定加密套件
最低TLS 1.3

创建Kubelet HTTP Client

构造ClientK8s对象

设置DeviceInfoName
= 前缀 + 节点名

创建RateLimitingQueue

IsApiErr = false

返回ClientK8s实例

逐行解析

// L72-73: 构建K8s配置,两个空字符串参数表示使用InClusterConfig
// 这是K8s推荐的方式:Pod内运行的服务自动使用ServiceAccount的凭据
clientCfg, err := clientcmd.BuildConfigFromFlags("", "")
// L77: 使用配置创建正式的K8s Clientset
// Clientset是K8s client-go提供的类型安全客户端,可操作所有K8s原生资源
client, err := kubernetes.NewForConfig(clientCfg)
// L82: 从环境变量 NODE_NAME 获取当前节点名
// 这是Device Plugin部署时通过K8s downward API注入的
nodeName, err := GetNodeNameFromEnv()
// L84-91: 构建Kubelet HTTP客户端的Transport
// 关键安全配置:
// - InsecureSkipVerify: 跳过证书验证(Kubelet通常使用自签证书)
// - CipherSuites: 限定使用ECDHE系列强加密套件,禁用弱算法
// - MinVersion: 强制TLS 1.3,拒绝旧版本协议
transport := &http.Transport{
    TLSClientConfig: &tls.Config{
        InsecureSkipVerify: true,
        CipherSuites:       defaultSafeCipherSuites,
        MinVersion:         tls.VersionTLS13,
    },
}
// L93: 创建Kubelet专用的HTTP客户端,后续用于直连Kubelet /pods接口
kltClient := &http.Client{Transport: transport}
// L95-103: 组装ClientK8s结构体
// - DeviceInfoName: 设备信息ConfigMap的命名规则 = "mindx-dl-deviceinfo-" + nodeName
// - Queue: 速率限制工作队列,用于Informer事件去重和背压控制
return &ClientK8s{
    Clientset:      client,
    NodeName:       nodeName,
    DeviceInfoName: common.DeviceInfoCMNamePrefix + nodeName,
    Queue:          workqueue.NewRateLimitingQueue(workqueue.DefaultControllerRateLimiter()),
    IsApiErr:       false,
    KltClient:      kltClient,
}, nil

安全加密套件设计意图

// L40-47: defaultSafeCipherSuites 限定了6种ECDHE套件
// 设计目的: 在InsecureSkipVerify=true的前提下,仍保证传输层加密强度
var defaultSafeCipherSuites = []uint16{
    tls.TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256,   // ECDSA + AES-128-GCM
    tls.TLS_ECDHE_ECDSA_WITH_AES_256_GCM_SHA384,   // ECDSA + AES-256-GCM
    tls.TLS_ECDHE_ECDSA_WITH_CHACHA20_POLY1305_SHA256, // ECDSA + ChaCha20
    tls.TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256,     // RSA + AES-128-GCM
    tls.TLS_ECDHE_RSA_WITH_AES_256_GCM_SHA384,     // RSA + AES-256-GCM
    tls.TLS_ECDHE_RSA_WITH_CHACHA20_POLY1305_SHA256, // RSA + ChaCha20
}
3.1.2 节点名验证 — GetNodeNameFromEnv() + checkNodeName()

文件: kubeclient.go L263-L285

// 从环境变量 NODE_NAME 获取节点名
func GetNodeNameFromEnv() (string, error) {
    nodeName := os.Getenv(api.NodeNameEnv)  // 读取环境变量 NODE_NAME
    if err := checkNodeName(nodeName); err != nil {
        return "", fmt.Errorf("check node name failed: %v", err)
    }
    return nodeName, nil
}
// 节点名合法性校验:长度、正则
func checkNodeName(nodeName string) error {
    // 检查1: 非空 — NODE_NAME 环境变量必须设置
    if len(nodeName) == 0 {
        return fmt.Errorf("the env variable whose key is NODE_NAME must be set")
    }
    // 检查2: 长度上限 230 字符 (KubeEnvMaxLength)
    // 防止超长节点名导致K8s资源名截断或API拒绝
    if len(nodeName) > common.KubeEnvMaxLength {
        return fmt.Errorf("node name length %d is bigger than %d", len(nodeName), common.KubeEnvMaxLength)
    }
    // 检查3: 正则匹配 — 确保节点名只包含合法字符
    // pattern 来自 common.GetPattern()["nodeName"],符合K8s节点名规范
    pattern := common.GetPattern()["nodeName"]
    if match := pattern.MatchString(nodeName); !match {
        return fmt.Errorf("node name %s is illegal", nodeName)
    }
    return nil
}
3.1.3 Pod 列表获取 — 三级数据获取策略

本模块设计了三级 Pod 数据获取策略,体现了从"最新数据"到"缓存降级"的容错设计:

方式1: Kubelet直连

方式2: API Server缓存

方式3: API Server直查

成功

失败

true: API有错误

false: 缓存有效

需要获取Pod列表

选择获取方式

GetPodsUsedNPUByKlt

GetActivePodListCache

GetActivePodList

构建HTTPS URL
HOST_IP:KUBELET_PORT/pods

带Bearer Token的GET请求

响应成功?

JSON解析PodList

降级到 GetPodsUsedNpuByCommon

遍历Pod提取NPU注解

IsApiErr标志?

refreshPodList
全量刷新缓存

直接读缓存

遍历podCache
过滤活跃Pod

返回Pod列表

List API调用
FieldSelector: spec.nodeName + 非Failed/Succeeded

checkPodList校验

返回Pod列表

方式1详解 — Kubelet直连 getPodsByKltPort():

文件: kube_connect.go L82-L112

func (ki *ClientK8s) getPodsByKltPort() (*v1.PodList, error) {
    // 步骤1: 创建带Token认证的HTTP请求
    // 内部调用 getKltPodsURL() 构建URL:
    //   - 从 HOST_IP 环境变量获取主机IP
    //   - 从 KUBELET_PORT 环境变量获取端口(默认10250)
    //   - IPv6地址加方括号
    //   - 路径固定为 /pods
    // 然后调用 createKltPodsReqWithToken():
    //   - 使用 rest.InClusterConfig() 获取ServiceAccount Token
    //   - 添加 Authorization: Bearer <token> 头
    req, err := createKltPodsReqWithToken()
    if err != nil {
        return nil, err
    }
    
    // 步骤2: 发送HTTP请求(使用NewClientK8s中创建的TLS 1.3客户端)
    resp, err := ki.KltClient.Do(req)
    if err != nil {
        return nil, err
    }
    defer resp.Body.Close()
    
    // 步骤3: 检查HTTP状态码
    if resp.StatusCode != http.StatusOK {
        return nil, fmt.Errorf("get kubelet http response failed, resp status code: %v is not %v",
            resp.StatusCode, http.StatusOK)
    }
    
    // 步骤4: 读取响应体(使用LimitReader防止响应过大导致OOM)
    // limiter.DefaultDataLimit 限制最大读取字节数
    body, err := io.ReadAll(io.LimitReader(resp.Body, limiter.DefaultDataLimit))
    if err != nil {
        if err == io.EOF {
            // 响应体超过限制,可能是异常的大响应
            return nil, fmt.Errorf("response size exceeds limit")
        }
        return nil, err
    }
    
    // 步骤5: JSON反序列化为PodList
    var podList v1.PodList
    if err = json.Unmarshal(body, &podList); err != nil {
        return nil, err
    }
    return &podList, nil
}

Kubelet URL构建详解 — getKltPodsURL():

文件: kube_connect.go L43-L67

func getKltPodsURL() (string, error) {
    // 从环境变量获取主机IP
    envHostIP := os.Getenv(HostIPEnv)  // "HOST_IP"
    parseIP := net.ParseIP(envHostIP)
    if parseIP == nil {
        return "", fmt.Errorf("host ip is invalid")  // IP格式非法
    }
    hostIP := parseIP.String()
    
    // 从环境变量获取Kubelet端口,默认10250
    kubeletPort := os.Getenv(KubeletPortEnv)  // "KUBELET_PORT"
    if err := isValidPort(kubeletPort); err != nil {
        // 端口无效时降级为默认端口
        kubeletPort = DefaultKubeletPort  // "10250"
    }
    
    // IPv6地址处理:需要加方括号
    // 例如: [::1]:10250 vs 192.168.1.1:10250
    if net2.IsIPv6(parseIP) {
        hostIP = "[" + hostIP + "]"
    }
    
    // 构建完整URL: https://<hostIP>:<kubeletPort>/pods
    host := fmt.Sprintf("%s:%s", hostIP, kubeletPort)
    kltPodsURL := &url.URL{
        Scheme: "https",
        Host:   host,
        Path:   "/pods",
    }
    return kltPodsURL.String(), nil
}

Token认证详解 — createKltPodsReqWithToken():

文件: kube_connect.go L70-L80

func createKltPodsReqWithToken() (*http.Request, error) {
    // 1. 构建Kubelet URL
    kltPodsURL, err := getKltPodsURL()
    if err != nil {
        return nil, err
    }
    
    // 2. 首次调用时记录URL(使用sync.Once确保只打印一次)
    logUrlOnce.Do(func() {
        hwlog.RunLog.Info("get kubelet pods url: ", kltPodsURL)
    })
    
    // 3. 创建HTTP GET请求
    req, err := http.NewRequest("GET", kltPodsURL, nil)
    if err != nil {
        return nil, err
    }
    
    // 4. 首次调用时初始化InClusterConfig获取BearerToken
    // 使用sync.Once确保只初始化一次,避免重复创建配置
    kubeConfigOnce.Do(func() {
        kubeConfig, err = rest.InClusterConfig()
        // InClusterConfig自动读取:
        //   - Token: /var/run/secrets/kubernetes.io/serviceaccount/token
        //   - CA证书: /var/run/secrets/kubernetes.io/serviceaccount/ca.crt
    })
    if kubeConfig == nil {
        return nil, fmt.Errorf("kubeConfig is nil")
    }
    
    // 5. 添加Bearer Token认证头
    // Kubelet验证Token以确认调用者是集群内的合法Pod
    req.Header.Add("Authorization", "Bearer "+kubeConfig.BearerToken)
    return req, nil
}

方式2详解 — API Server缓存读取 GetActivePodListCache():

文件: kube_cache.go L193-L213

func (ki *ClientK8s) GetActivePodListCache() []v1.Pod {
    // 关键: 检查IsApiErr标志
    // 当API Server调用出错时,IsApiErr被置为true
    // 下次读取缓存时会触发全量刷新,保证缓存数据不会过时太久
    if ki.IsApiErr {
        ki.refreshPodList()  // 重新从API Server全量拉取
    }
    
    newPodList := make([]v1.Pod, 0, common.GeneralMapSize)
    lock.Lock()
    defer lock.Unlock()
    
    // 遍历内存缓存 podCache (map[types.UID]*podInfo)
    for _, pi := range podCache {
        // 安全检查1: Pod名称合法性
        if err := common.CheckPodNameAndSpace(pi.GetName(), common.PodNameMaxLength); err != nil {
            continue  // 跳过非法名称的Pod
        }
        // 安全检查2: Pod命名空间合法性
        if err := common.CheckPodNameAndSpace(pi.GetNamespace(), common.PodNameSpaceMaxLength); err != nil {
            continue
        }
        // 过滤条件: 排除已结束的Pod(Failed和Succeeded状态)
        if pi.Status.Phase == v1.PodFailed || pi.Status.Phase == v1.PodSucceeded {
            continue
        }
        newPodList = append(newPodList, *pi.Pod)
    }
    return newPodList
}

缓存刷新机制 — refreshPodList():

文件: kube_cache.go L155-L169

func (ki *ClientK8s) refreshPodList() {
    // 从API Server全量拉取当前节点的Pod列表
    newV1PodList, err := ki.GetAllPodList()
    if err != nil {
        // 拉取失败时保持旧缓存不变,避免错误数据覆盖
        return
    }
    
    // 构建全新的缓存Map
    newPodCache := map[types.UID]*podInfo{}
    for _, pod := range newV1PodList.Items {
        // 关键: Go的for range遍历slice时,循环变量地址固定
        // 直接取 &pod 会导致所有缓存项指向同一个地址
        // 因此使用闭包参数传递,确保每个pod有独立地址
        func(pod v1.Pod) {
            newPodCache[pod.UID] = &podInfo{
                Pod:        &pod,
                updateTime: time.Now(),
            }
        }(pod)
    }
    
    lock.Lock()
    podCache = newPodCache  // 原子替换整个缓存
    lock.Unlock()
    
    ki.IsApiErr = false  // 重置错误标志
}
3.1.4 Pod Informer 初始化与事件处理
WorkQueue podCache PodInformer SharedInformerFactory InitPodInformer main.go WorkQueue podCache PodInformer SharedInformerFactory InitPodInformer main.go 事件处理handler1: Add → UpdatePodList(obj, EventTypeAdd) Update → UpdatePodList(obj, EventTypeUpdate) Delete → UpdatePodList(obj, EventTypeDelete) 事件处理handler2: FilteringResourceEventHandler 过滤条件: checkPod (含HuaweiAscend910注解) 通过过滤的事件入WorkQueue Watch错误处理: 日志记录 + FlushPodCacheNextQuerying InitPodInformer() NewSharedInformerFactoryWithOptions FieldSelector: spec.nodeName=<当前节点> Core().V1().Pods().Informer() AddEventHandler (第一个handler) AddEventHandler (第二个handler) SetWatchErrorHandler Start(make(chan struct{})) 开始List/Watch Add事件 → UpdatePodList Update事件 → UpdatePodList Delete事件 → UpdatePodList 过滤后的事件入队

文件: cur_node_informer.go L33-L56

func (ki *ClientK8s) InitPodInformer() {
    // 创建SharedInformerFactory,关键配置:
    // 1. 使用ki.Clientset作为K8s客户端
    // 2. resyncPeriod=0 表示不做定期全量同步,依赖Watch事件驱动
    // 3. WithTweakListOptions: 设置FieldSelector只关注当前节点的Pod
    //    这大幅减少了API Server返回的数据量和网络开销
    factory := informers.NewSharedInformerFactoryWithOptions(ki.Clientset, 0,
        informers.WithTweakListOptions(func(options *v1.ListOptions) {
            options.FieldSelector = "spec.nodeName=" + ki.NodeName
        }))
    
    // 创建Pod Informer
    podInformer := factory.Core().V1().Pods().Informer()
    
    // 注册第一个EventHandler: 维护podCache
    // 所有Pod变更都会触发UpdatePodList,更新内存缓存
    podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
        AddFunc: func(obj interface{}) {
            ki.UpdatePodList(obj, EventTypeAdd)
        },
        UpdateFunc: func(oldObj, newObj interface{}) {
            // DeepEqual优化: 只有Pod对象实际变更时才处理
            if !reflect.DeepEqual(oldObj, newObj) {
                ki.UpdatePodList(newObj, EventTypeUpdate)
            }
        },
        DeleteFunc: func(obj interface{}) {
            ki.UpdatePodList(obj, EventTypeDelete)
        },
    })
    
    // 注册第二个EventHandler: 过滤后入队
    // 使用FilteringResourceEventHandler,过滤条件为checkPod
    // checkPod: 只关注带有 HuaweiAscend910 注解的Pod
    // 这些是使用了Ascend NPU的Pod,需要特殊处理
    podInformer.AddEventHandler(ki.ResourceEventHandler(PodResource, checkPod))
    
    // 设置Watch错误处理器
    // 当Informer与API Server的Watch连接断开时触发
    podInformer.SetWatchErrorHandler(func(r *cache.Reflector, err error) {
        hwlog.RunLog.Errorf("pod informer watch error: %v", err)
        // 如果配置了DealWatchHandler,则标记缓存需要刷新
        // 下次读取缓存时会从API Server全量拉取
        if common.ParamOption.DealWatchHandler {
            ki.FlushPodCacheNextQuerying()  // 设置IsApiErr=true
        }
    })
    
    // 启动Informer Factory
    factory.Start(make(chan struct{}))
    
    // 保存PodInformer引用到ClientK8s
    ki.PodInformer = podInformer
}

checkPod过滤函数:

// cur_node_informer.go L58-L64
// 过滤条件: Pod必须包含 HuaweiAscend910 注解
// 这意味着只有使用了Ascend 910 NPU的Pod才会被加入WorkQueue
// 其他普通Pod不会触发队列事件,减少不必要处理
func checkPod(obj interface{}) bool {
    pod, ok := obj.(*corev1.Pod)
    if !ok {
        return false
    }
    _, exist := pod.Annotations[api.HuaweiAscend910]
    return exist
}
3.1.5 Pod缓存更新与Predicate Time校正机制

文件: kube_cache.go L83-L130

失败

不存在

存在

UpdatePodList 入口
newObj, operator

类型断言
*v1.Pod?

return

加锁 lock.Lock

从Informer Indexer获取
最新Pod对象

日志错误, return

从podCache删除
触发TriggerUpdate

类型断言 *v1.Pod

调用 updatePodPredicateTime
校正predicate-time注解

更新podCache:
存储Pod + 当前时间

触发 TriggerUpdate
通知其他组件数据变更

updatePodPredicateTime() 详解 — 这是调度纠正机制的核心:

// kube_cache.go L132-L153
func (ki *ClientK8s) updatePodPredicateTime(newPod *v1.Pod) {
    // 从缓存中查找该Pod的旧信息
    cachedPodInfo, exists := podCache[newPod.UID]
    if !exists {
        return  // 缓存中不存在,无需校正
    }
    
    // 检查缓存中Pod的注解是否存在
    if cachedPodInfo.Pod.Annotations == nil {
        return
    }
    
    // 读取缓存中的 predicate-time 值
    // predicate-time 是调度阶段的标记,值为 MaxUint64 表示尚未被分配
    cachedPredicateTime := cachedPodInfo.Pod.Annotations[common.PodPredicateTime]
    if cachedPredicateTime != strconv.FormatUint(math.MaxUint64, common.BaseDec) {
        return  // 只有 MaxUint64 才需要校正
    }
    
    // 检查新Pod的注解是否与缓存一致
    needsUpdate := false
    if newPod.Annotations == nil {
        // 新Pod没有注解,需要补充
        newPod.Annotations = make(map[string]string)
        needsUpdate = true
    }
    if newPod.Annotations[common.PodPredicateTime] != cachedPredicateTime {
        // 新Pod的predicate-time与缓存不一致,需要校正
        // 这种情况可能发生在: API Server返回的Pod数据丢失了annotation
        needsUpdate = true
    }
    
    if needsUpdate {
        // 将缓存中的正确值写回新Pod对象
        newPod.Annotations[common.PodPredicateTime] = cachedPredicateTime
        // 异步更新API Server上的Pod annotation
        // 使用goroutine避免阻塞Informer事件处理循环
        annotation := map[string]string{common.PodPredicateTime: strconv.FormatUint(math.MaxUint64, common.BaseDec)}
        go ki.TryUpdatePodCacheAnnotation(newPod, annotation)
    }
}

设计意图分析predicate-time 是华为自研的调度标记,当值为 MaxUint64 时表示 Pod 已通过调度但尚未完成设备分配。如果 API Server 返回的数据丢失了这个注解,可能导致重复分配。此机制确保缓存中的正确值被同步回 API Server。

3.1.6 Pod缓存定期巡检 — PodInformerInspector()
渲染错误: Mermaid 渲染失败: Parse error on line 2: ...[计算节点Hash值
FNV32(nodeName) % 600] -----------------------^ Expecting 'SQE', 'DOUBLECIRCLEEND', 'PE', '-)', 'STADIUMEND', 'SUBROUTINEEND', 'PIPE', 'CYLINDEREND', 'DIAMOND_STOP', 'TAGEND', 'TRAPEND', 'INVTRAPEND', 'UNICODE_TEXT', 'TEXT', 'TAGSTART', got 'PS'

文件: kube_cache.go L40-L75

func (ki *ClientK8s) PodInformerInspector(ctx context.Context) {
    // 使用FNV哈希将节点名映射到0-599的数字
    // 目的: 让不同节点的Device Plugin在不同的时间点开始巡检
    // 避免所有节点同时向API Server发起大量查询(避免惊群效应)
    hashVal := fnv.New32()
    if _, err := hashVal.Write([]byte(ki.NodeName)); err != nil {
        return
    }
    val := hashVal.Sum32() % periodicForStartCheck  // periodicForStartCheck = 600
    
    // 随机延迟启动
    time.Sleep(time.Duration(val) * time.Second)
    
    // 每10分钟执行一次巡检
    // wait.Until会在ctx.Done()时自动停止
    wait.Until(func() { ki.checkPodInCache(ctx) }, timeIntervalForCheckPod, ctx.Done())
    // timeIntervalForCheckPod = 10 * time.Minute
}
func (ki *ClientK8s) checkPodInCache(ctx context.Context) {
    lock.Lock()
    defer lock.Unlock()
    
    needDelete := make([]types.UID, 0)    // 待删除的Pod UID列表
    needRefresh := make([]types.UID, 0)    // 待刷新的Pod UID列表
    
    for uid, pi := range podCache {
        // 检查1: 更新时间距今是否小于1小时
        // podCacheTimeout = time.Hour
        // 1小时内的Pod视为有效,跳过检查
        if time.Since(pi.updateTime) < podCacheTimeout {
            continue
        }
        
        // 超过1小时未更新的Pod,需要从API Server验证
        pod, err := ki.getPod(ctx, pi.Namespace, pi.Name)
        if err != nil {
            if errors.IsNotFound(err) {
                // Pod已被删除,从缓存中清除
                needDelete = append(needDelete, uid)
                continue
            }
            // 其他错误(如网络超时)跳过,下次巡检再检查
            continue
        }
        
        // 检查2: Pod是否仍然调度在当前节点
        // 可能Pod被重新调度到了其他节点
        if pod.Spec.NodeName != ki.NodeName || pod.UID != uid {
            needDelete = append(needDelete, uid)
            continue
        }
        
        // Pod仍然有效,刷新更新时间
        needRefresh = append(needRefresh, uid)
    }
    
    // 批量执行删除
    for _, uid := range needDelete {
        delete(podCache, uid)
    }
    
    // 批量执行刷新
    for _, uid := range needRefresh {
        podCache[uid].updateTime = time.Now()
    }
}
3.1.7 ConfigMap 写入机制
3.1.7.1 设备信息写入 WriteDeviceInfoDataIntoCM()

文件: client_server.go L156-L207

失败

成功

失败

成功

Ascend910A5

Ascend910A3

其他

NotFound

成功

失败

其他错误

WriteDeviceInfoDataIntoCM
nodeDeviceData, manuallySeparateNPU,
switchInfo, dpuInfo, reasonCm

计算校验码
MakeDataHash

序列化设备数据
MarshalData

返回错误

序列化交换机故障数据

返回错误

构建ConfigMap对象
Name: mindx-dl-deviceinfo-节点名
Namespace: KubeNS
Labels: CIMCMLabelKey

RealCardType?

设置A5特有数据:
DeviceInfo + SwitchInfo +
UpgradeFaultReason +
ManuallySeparateNPU + Description

DPU信息非空?

序列化DPU数据
加入DpuInfoCMDataKey

跳过DPU数据

设置A3特有数据:
DeviceInfo + SwitchInfo +
UpgradeFaultReason +
ManuallySeparateNPU + Description
无DPU

设置默认数据:
DeviceInfo +
ManuallySeparateNPU +
UpgradeFaultReason + Description
无SwitchInfo

调用createOrUpdateDeviceCM

Update成功?

返回nodeDeviceData, nil

调用Create

返回错误

逐行解析

func (ki *ClientK8s) WriteDeviceInfoDataIntoCM(nodeDeviceData *common.NodeDeviceInfoCache,
    manuallySeparateNPU string, switchInfo common.SwitchFaultInfo, dpuInfo common.DpuInfo,
    reasonCm string) (*common.NodeDeviceInfoCache, error) {
    
    // 步骤1: 计算数据校验码
    // MakeDataHash 对 DeviceInfo 内容做哈希,用于数据完整性校验
    // 下游消费者可以通过比较CheckCode判断数据是否变更
    nodeDeviceData.CheckCode = common.MakeDataHash(nodeDeviceData.DeviceInfo)
    
    // 步骤2: 序列化各数据块
    var data, switchData, dpuData []byte
    // 判断DPU功能是否开启:通过DeepEqual比较dpuInfo与零值
    dpuOpen := !reflect.DeepEqual(dpuInfo, common.DpuInfo{})
    
    if data = common.MarshalData(nodeDeviceData); len(data) == 0 {
        return nil, fmt.Errorf("marshal nodeDeviceData failed")
    }
    if switchData = common.MarshalData(switchInfo); len(switchData) == 0 {
        return nil, fmt.Errorf("marshal switchDeviceData failed")
    }
    
    // 步骤3: 构建ConfigMap对象
    deviceInfoCM := &v1.ConfigMap{
        ObjectMeta: metav1.ObjectMeta{
            Name:      ki.DeviceInfoName,  // "mindx-dl-deviceinfo-" + nodeName
            Namespace: api.KubeNS,         // 华为设备插件专用namespace
            Labels:    map[string]string{api.CIMCMLabelKey: common.CmConsumerValue},
        },
    }
    
    // 步骤4: 按硬件型号设置不同数据结构
    switch common.ParamOption.RealCardType {
    case api.Ascend910A5:
        // A5型号: 完整数据,包含交换机故障信息和可选的DPU信息
        deviceInfoCM.Data = map[string]string{
            api.DeviceInfoCMDataKey:                   string(data),
            api.SwitchInfoCMDataKey:                   string(switchData),
            common.DeviceInfoCmUpgradeFaultReasonKey:  reasonCm,
            common.DeviceInfoCMManuallySeparateNPUKey: manuallySeparateNPU,
            common.DescriptionKey:                     common.DescriptionValue,
        }
        if dpuOpen {
            if dpuData = common.MarshalData(dpuInfo); len(dpuData) == 0 {
                return nil, fmt.Errorf("marshal DpuDeviceData failed")
            }
            deviceInfoCM.Data[api.DpuInfoCMDataKey] = string(dpuData)
        }
    case api.Ascend910A3:
        // A3型号: 有交换机信息但无DPU
        deviceInfoCM.Data = map[string]string{
            api.DeviceInfoCMDataKey:                   string(data),
            api.SwitchInfoCMDataKey:                   string(switchData),
            common.DeviceInfoCMManuallySeparateNPUKey: manuallySeparateNPU,
            common.DeviceInfoCmUpgradeFaultReasonKey:  reasonCm,
            common.DescriptionKey:                     common.DescriptionValue,
        }
    default:
        // 其他型号: 只有基础设备信息,无交换机/DPU
        deviceInfoCM.Data = map[string]string{
            api.DeviceInfoCMDataKey:                   string(data),
            common.DeviceInfoCMManuallySeparateNPUKey: manuallySeparateNPU,
            common.DeviceInfoCmUpgradeFaultReasonKey:  reasonCm,
            common.DescriptionKey:                     common.DescriptionValue,
        }
    }
    
    // 步骤5: 写入ConfigMap(先Update,不存在则Create)
    if err := ki.createOrUpdateDeviceCM(deviceInfoCM); err != nil {
        return nil, err
    }
    return nodeDeviceData, nil
}
3.1.7.2 重置信息写入 WriteResetInfoDataIntoCM()

文件: client_server.go L209-L260

失败

成功

数据不存在

存在

成功

失败

WriteResetInfoDataIntoCM
taskName, namespace,
taskInfo, needAddRetry

获取现有重置CM
GetConfigMap reset-config-任务名

返回错误

读取现有重置数据

返回错误

旧数据含隔离错误
且新设备列表非空?

返回错误: 需要重新调度

反序列化旧数据

获取旧重试次数

needAddRetry?

retryTime + 1

保持旧值

构建新TaskInfo
setNewTaskInfoWithHexString

设置UpdateTime = now

设置RetryTime

计算校验码

序列化数据

构建CM对象
保留旧TypeMeta和ObjectMeta

needAddRetry?

设置重置类型 = HotReset

保留旧重置类型

UpdateConfigMap

返回CM

返回错误

关键逻辑解析

// 检查隔离错误 + 新设备列表非空的情况
// 如果旧CM中已包含 IsolateError 标记,且新的任务信息中有设备
// 则拒绝更新,返回"task should be rescheduled"
// 设计意图: 隔离错误状态下不应该继续添加设备,需要重新调度
if strings.Contains(oldResetInfoData, common.IsolateError) && len(taskInfo.RankList) != 0 {
    return nil, fmt.Errorf("task should be rescheduled")
}
// setNewTaskInfoWithHexString: 将错误码转换为十六进制字符串
// 设计意图: 上游消费者可能以字符串形式处理错误码,
// 十六进制表示更紧凑且便于人工阅读
func setNewTaskInfoWithHexString(taskInfo *common.TaskResetInfo) *common.TaskResetInfo {
    var newTaskInfo common.TaskResetInfo
    for _, deviceInfo := range taskInfo.RankList {
        newDeviceInfo := *deviceInfo
        // 转换为十六进制大写
        newDeviceInfo.ErrorCodeHex = strings.ToUpper(common.Int64Tool.ToHexString(newDeviceInfo.ErrorCode))
        newDeviceInfo.ErrorCode = []int64{}  // 清空原始错误码,节省存储空间
        newTaskInfo.RankList = append(newTaskInfo.RankList, &newDeviceInfo)
    }
    if newTaskInfo.RankList == nil {
        newTaskInfo.RankList = make([]*common.TaskDevInfo, 0)
    }
    return &newTaskInfo
}
3.1.8 NPU使用情况查询

本模块提供两条路径查询节点上NPU的使用情况,形成主备容错:

失败

查询NPU使用情况

主路径:
GetPodsUsedNPUByKlt

通过Kubelet /pods获取Pod列表

遍历Pod提取注解
PodAnnotationAscendReal

解析注解中的设备ID列表

返回sets.String

降级路径:
GetPodsUsedNpuByCommon

从内存缓存获取活跃Pod列表
GetActivePodListCache

遍历Pod提取注解
PodAnnotationAscendReal

解析注解中的设备ID列表

返回sets.String

文件: client_server.go L318-L372

// 主路径: 通过Kubelet直接获取(数据更新,无API Server延迟)
func (ki *ClientK8s) GetPodsUsedNPUByKlt() sets.String {
    podList, err := ki.getPodsByKltPort()
    if err != nil {
        // Kubelet连接失败时,使用带限流的日志
        // domainForKubeletConnectErr = "kubeletConnect"
        // 防止Kubelet持续不可用时日志刷屏
        hwlog.RunLog.ErrorfWithLimit(domainForKubeletConnectErr, os.Getenv(KubeletPortEnv),
            "get pods used NPU failed: %v", err)
        // 降级到缓存路径
        return ki.GetPodsUsedNpuByCommon()
    }
    
    usedNPU := make([]string, 0)
    for _, pod := range podList.Items {
        // 安全检查: Pod名称和命名空间合法性
        if err := common.CheckPodNameAndSpace(pod.GetName(), common.PodNameMaxLength); err != nil {
            continue
        }
        if err := common.CheckPodNameAndSpace(pod.GetNamespace(), common.PodNameSpaceMaxLength); err != nil {
            continue
        }
        // 过滤已结束的Pod
        if pod.Status.Phase == v1.PodFailed || pod.Status.Phase == v1.PodSucceeded {
            continue
        }
        
        // 读取 PodAnnotationAscendReal 注解
        // 该注解格式: "Ascend910-0,Ascend910-1" 表示使用了设备0和1
        realAllocTag := fmt.Sprintf("%s", api.PodAnnotationAscendReal)
        tmpNPU, ok := pod.Annotations[realAllocTag]
        if !ok || len(tmpNPU) == 0 || len(tmpNPU) > common.PodAnnotationMaxLength {
            continue
        }
        
        // 分割设备列表
        tmpNPUList := strings.Split(tmpNPU, common.CommaSepDev)
        if len(tmpNPUList) == 0 || len(tmpNPUList) > common.MaxDevicesNum {
            continue
        }
        usedNPU = append(usedNPU, tmpNPUList...)
    }
    
    // 使用sets.String去重
    return sets.NewString(usedNPU...)
}
3.1.9 资源事件处理机制 — ResourceEventHandler()

文件: kubeclient.go L287-L312

func (ki *ClientK8s) ResourceEventHandler(res ResourceType, filter func(obj interface{}) bool) cache.
    ResourceEventHandler {
    
    // enqueue: 内部闭包,将事件加入WorkQueue
    enqueue := func(obj interface{}, event EventType) {
        // 关键优化: Pod的Update事件不入队
        // 因为Pod更新通过第一个EventHandler(UpdatePodList)已处理
        // 这里只需要Add和Delete事件入队
        if res == PodResource && event == EventTypeUpdate {
            return
        }
        
        // 生成对象的唯一Key: namespace/name
        key, err := cache.MetaNamespaceKeyFunc(obj)
        if err != nil {
            return
        }
        
        // 将事件加入速率限制队列
        // RateLimitingQueue提供:
        // 1. 去重: 相同Key的事件会被合并
        // 2. 限流: 失败重试有指数退避
        ki.Queue.AddRateLimited(Event{
            Resource: res,
            Key:      key,
            Type:     event,
        })
    }
    
    // 返回FilteringResourceEventHandler
    // 先通过filter函数过滤,再调用对应的处理函数
    return cache.FilteringResourceEventHandler{
        FilterFunc: filter,  // 过滤函数: checkPod(检查HuaweiAscend910注解)
        Handler: cache.ResourceEventHandlerFuncs{
            AddFunc: func(obj interface{}) {
                enqueue(obj, EventTypeAdd)
            },
            UpdateFunc: func(oldObj, newObj interface{}) {
                // DeepEqual优化: 只有对象实际变更才处理
                if reflect.DeepEqual(oldObj, newObj) {
                    return
                }
                enqueue(newObj, EventTypeUpdate)
            },
            DeleteFunc: func(obj interface{}) {
                enqueue(obj, EventTypeDelete)
            },
        },
    }
}
3.1.10 节点资源清理 — RemoveOldResource()

文件: kubeclient.go L420-L452

失败

成功

失败

成功

RemoveOldResource
oldResourceName

获取当前节点对象

返回错误

节点Capacity中
存在旧资源?

日志: 旧资源不存在, 跳过

返回nil

从Capacity删除旧资源

从Allocatable删除旧资源

UpdateStatus更新节点状态

返回错误

验证Capacity
中已删除?

返回错误: 删除失败

验证Allocatable
中已删除?

返回错误: 删除失败

日志: 清理成功

返回nil

func (c *ClientK8s) RemoveOldResource(oldResourceName string) error {
    // 场景: 设备型号切换时,需要清理旧的资源名称
    // 例如: 从 Ascend910-0 切换到 Ascend910A5-0,旧的资源名需要删除
    
    node, err := c.Clientset.CoreV1().Nodes().Get(
        context.TODO(),
        c.NodeName,
        metav1.GetOptions{ResourceVersion: "0"},
    )
    if err != nil {
        return fmt.Errorf("get node failed: %v", err)
    }
    
    // 检查旧资源是否存在于节点容量中
    if _, exists := node.Status.Capacity[v1.ResourceName(oldResourceName)]; !exists {
        // 不存在则跳过,幂等操作
        return nil
    }
    
    // 从Capacity和Allocatable中删除旧资源
    delete(node.Status.Capacity, v1.ResourceName(oldResourceName))
    delete(node.Status.Allocatable, v1.ResourceName(oldResourceName))
    
    // 更新节点状态到API Server
    updatedNode, err := c.Clientset.CoreV1().Nodes().UpdateStatus(
        context.TODO(),
        node,
        metav1.UpdateOptions{},
    )
    if err != nil {
        return fmt.Errorf("update node status failed: %v", err)
    }
    
    // 验证删除是否成功(双重确认)
    if _, exists := updatedNode.Status.Capacity[v1.ResourceName(oldResourceName)]; exists {
        return errors.New("failed to remove resource from capacity")
    }
    if _, exists := updatedNode.Status.Allocatable[v1.ResourceName(oldResourceName)]; exists {
        return errors.New("failed to remove resource from allocatable")
    }
    
    return nil
}
3.1.11 容器运行时检测 — GetContainerRuntime()

文件: kubeclient.go L190-L207

func (ki *ClientK8s) GetContainerRuntime() (string, error) {
    // 通过节点状态信息获取容器运行时
    node, err := ki.GetNode()
    if err != nil {
        return "", fmt.Errorf("failed to get node: %w", err)
    }
    
    // 读取节点状态中的容器运行时版本字符串
    // 格式: "docker://1.13.1" 或 "containerd://1.6.0"
    runtimeVersion := node.Status.NodeInfo.ContainerRuntimeVersion
    if runtimeVersion == "" {
        return "", fmt.Errorf("container runtime version not found in node status")
    }
    
    // 通过URL Scheme前缀判断运行时类型
    if strings.HasPrefix(runtimeVersion, "docker://") {
        return DockerRuntime, nil  // "docker"
    }
    if strings.HasPrefix(runtimeVersion, "containerd://") {
        return ContainerdRuntime, nil  // "containerd"
    }
    
    return "", fmt.Errorf("unknown container runtime: %s", runtimeVersion)
}

设计意图:此方法为 duplicatedetector 模块提供运行时类型判断,决定创建 Docker 客户端还是 containerd 客户端。

3.1.12 IsApiErr 错误标记机制

是, 包含443端口

Informer Watch断开

外部主动调用

API Server调用

返回错误?

IsApiErr = true

IsApiErr保持不变

下次读取缓存时
触发refreshPodList

刷新成功?

IsApiErr = false

IsApiErr保持true
下次再刷新

WatchErrorHandler

FlushPodCacheNextQuerying

机制说明

  • 任何 K8s API 调用返回包含 “443”(ApiServerPort)的错误时,IsApiErr 被置为 true
  • 443 是 K8s API Server 的默认端口,错误中包含此端口表示 API Server 不可达
  • IsApiErr=true 时,下次读取缓存会触发全量刷新
  • 刷新成功后重置为 false,形成自恢复闭环
Logo

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

更多推荐