statusManager
南风未起 2024-06-30 20:36:22 kubelet
# 1.简介
statusManager负责维护状态信息,将pod更新到apiserver。statusManager不会主动监控pod状态的变化,会提供对应接口供其它组件调用。// Manager should be kept up-to-date with the latest v1.PodStatus. It also syncs updates back to the API server. type Manager interface { // 状态获取接口,供外部调用 PodStatusProvider // 同步状态至APIServer主循环 Start() // 设置podStatus,触发sync SetPodStatus(pod *v1.Pod, status v1.PodStatus) // 设置pod.containerStatus就绪 SetContainerReadiness(podUID types.UID, containerID kubecontainer.ContainerID, ready bool) // 设置pod.containerStatus启动 SetContainerStartup(podUID types.UID, containerID kubecontainer.ContainerID, started bool) // 设置pod.container状态为Terminated TerminatePod(pod *v1.Pod) // 删除statusManager缓存的无效数据 RemoveOrphanedStatuses(podUIDs map[types.UID]bool) } // Updates pod statuses in apiserver. Writes only when new status has changed. All methods are thread-safe. type manager struct { kubeClient clientset.Interface // 维护的缓存Pod podManager kubepod.Manager // Pod状态缓存 podStatuses map[types.UID]versionedPodStatus podStatusesLock sync.RWMutex // 接收更新podStatus的管道 podStatusChannel chan podStatusSyncRequest // podStatus的版本号 apiStatusVersions map[kubetypes.MirrorPodUID]uint64 // 删除Pod的接口 podDeletionSafety PodDeletionSafetyProvider } // NewManager returns a functional Manager. func NewManager(kubeClient clientset.Interface, podManager kubepod.Manager, podDeletionSafety PodDeletionSafetyProvider) Manager { return &manager{ kubeClient: kubeClient, podManager: podManager, podStatuses: make(map[types.UID]versionedPodStatus), podStatusChannel: make(chan podStatusSyncRequest, 1000), // Buffer up to 1000 statuses apiStatusVersions: make(map[kubetypes.MirrorPodUID]uint64), // PodDeletionSafetyProvider接口,用Pod数据检查 podDeletionSafety: podDeletionSafety, } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
# 2.statusManager
# 2.1.Start
statusManager.Start()会后台运行协程,间隔10s周期同步一次状态,也会消费状态管道调用syncPod同步状态。func (m *manager) Start() { if m.kubeClient == nil { return } // 间隔10s定时器 syncTicker := time.NewTicker(syncPeriod).C // syncPod and syncBatch share the same go routine to avoid sync races. go wait.Forever(func() { for { select { // 监听podStatusChannel case syncRequest := <-m.podStatusChannel: // 触发同步 m.syncPod(syncRequest.podUID, syncRequest.status) // 定时器 case <-syncTicker: klog.V(5).InfoS("Status Manager: syncing batch") // 批量同步前清理channel信号 for i := len(m.podStatusChannel); i > 0; i-- { <-m.podStatusChannel } // 批量同步 m.syncBatch() } } }, 0) } // syncPod syncs the given status with the API server. The caller must not hold the lock. func (m *manager) syncPod(uid types.UID, status versionedPodStatus) { // pod无需更新 if !m.needsUpdate(uid, status) { return } // 获取Pod pod, err := m.kubeClient.CoreV1().Pods(status.podNamespace).Get(context.TODO(), status.podName, metav1.GetOptions{}) if errors.IsNotFound(err) { return } ... // 获取podManager缓存的poduid translatedUID := m.podManager.TranslatePodUID(pod.UID) // 缓存和当前uid不一致,说明pod重建 if len(translatedUID) > 0 && translatedUID != kubetypes.ResolvedPodUID(uid) { // 清理旧的状态 m.deletePodStatus(uid) return } // 合并pod状态,保留部分关键字 mergedStatus := mergePodStatus(pod.Status, status.status, m.podDeletionSafety.PodCouldHaveRunningContainers(pod)) // 生成patch触发更新 newPod, patchBytes, unchanged, err := statusutil.PatchPodStatus(m.kubeClient, pod.Namespace, pod.Name, pod.UID, pod.Status, mergedStatus) ... if !unchanged { // 记录更新的pod状态 pod = newPod } // 记录pod状态版本 m.apiStatusVersions[kubetypes.MirrorPodUID(pod.UID)] = status.version // pod可以删除 if m.canBeDeleted(pod, status.status) { ... m.kubeClient.CoreV1().Pods(pod.Namespace).Delete(context.TODO(), pod.Name, deleteOptions) ... // 清理pod状态 m.deletePodStatus(uid) } } // syncBatch syncs pods statuses with the apiserver. func (m *manager) syncBatch() { ... // 获取mirrorPod映射 podToMirror, mirrorToPod := m.podManager.GetUIDTranslations() func() { // Critical section m.podStatusesLock.RLock() defer m.podStatusesLock.RUnlock() // 清理无效的状态版本 for uid := range m.apiStatusVersions { _, hasPod := m.podStatuses[types.UID(uid)] _, hasMirror := mirrorToPod[uid] if !hasPod && !hasMirror { delete(m.apiStatusVersions, uid) } } // 整理待同步的pod for uid, status := range m.podStatuses { syncedUID := kubetypes.MirrorPodUID(uid) // 基于mirrorPod身份同步 if mirrorUID, ok := podToMirror[kubetypes.ResolvedPodUID(uid)]; ok { if mirrorUID == "" { continue } syncedUID = mirrorUID } // 需要更新 if m.needsUpdate(types.UID(syncedUID), status) { updatedStatuses = append(updatedStatuses, podStatusSyncRequest{uid, status}) } else if m.needsReconcile(uid, status.status) { // 状态不一致,清理版本再更新 delete(m.apiStatusVersions, syncedUID) updatedStatuses = append(updatedStatuses, podStatusSyncRequest{uid, status}) } } }() for _, update := range updatedStatuses { // patch状态至apiserver m.syncPod(update.podUID, update.status) } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
# 2.2.SetPodStatus
SetPodStatus()用于设置pod status subResource,设置时会触发同步操作,执行一次状态更新。func (m *manager) SetPodStatus(pod *v1.Pod, status v1.PodStatus) { m.podStatusesLock.Lock() defer m.podStatusesLock.Unlock() // Make sure we're caching a deep copy. status = *status.DeepCopy() // 缓存状态及触发更新 m.updateStatusInternal(pod, status, pod.DeletionTimestamp != nil) } // updateStatusInternal updates the internal status cache, and queues an update to the api server if necessary. func (m *manager) updateStatusInternal(pod *v1.Pod, status v1.PodStatus, forceUpdate bool) bool { ... // 获取缓存旧状态 cachedStatus, isCached := m.podStatuses[pod.UID] if isCached { oldStatus = cachedStatus.status } else if mirrorPod, ok := m.podManager.GetMirrorPodByPod(pod); ok { oldStatus = mirrorPod.Status } else { oldStatus = pod.Status } // 检查container状态合法性(状态非法转移) checkContainerStateTransition(oldStatus.ContainerStatuses, status.ContainerStatuses, pod.Spec.RestartPolicy) ... checkContainerStateTransition(oldStatus.InitContainerStatuses, status.InitContainerStatuses, pod.Spec.RestartPolicy) ... // 更新不同condition的LastTransitionTime updateLastTransitionTime(&status, &oldStatus, v1.ContainersReady) updateLastTransitionTime(&status, &oldStatus, v1.PodReady) updateLastTransitionTime(&status, &oldStatus, v1.PodInitialized) updateLastTransitionTime(&status, &oldStatus, v1.PodScheduled) // 保留startTime(确保仅设置一次) if oldStatus.StartTime != nil && !oldStatus.StartTime.IsZero() { status.StartTime = oldStatus.StartTime } else if status.StartTime.IsZero() { // if the status has no start time, we need to set an initial time now := metav1.Now() status.StartTime = &now } // 标准化status normalizeStatus(pod, &status) ... // 状态无变化及非强制更新 if isCached && isPodStatusByKubeletEqual(&cachedStatus.status, &status) && !forceUpdate { return false // No new status. } // 缓存更新的状态 newStatus := versionedPodStatus{ status: status, version: cachedStatus.version + 1, podName: pod.Name, podNamespace: pod.Namespace, } m.podStatuses[pod.UID] = newStatus select { // 向channel发送更新请求 case m.podStatusChannel <- podStatusSyncRequest{pod.UID, newStatus}: return true default: return false } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
# 2.3.SetContainerReadiness
SetContainerReadiness()用于设置pod/container的ready状态,向podStatusChannel发送事件驱动状态同步。func (m *manager) SetContainerReadiness(podUID types.UID, containerID kubecontainer.ContainerID, ready bool) { m.podStatusesLock.Lock() defer m.podStatusesLock.Unlock() // 获取缓存的pod对象 pod, ok := m.podManager.GetPodByUID(podUID) ... // 获取缓存的pod状态 oldStatus, found := m.podStatuses[pod.UID] ... // 获取container状态 containerStatus, _, ok := findContainerStatus(&oldStatus.status, containerID.String()) ... // ready无需更新 if containerStatus.Ready == ready { return } // 拷贝更新container状态为ready status := *oldStatus.status.DeepCopy() containerStatus, _, _ = findContainerStatus(&status, containerID.String()) containerStatus.Ready = ready ... // 更新PodReadyCondition updateConditionFunc(v1.PodReady, GeneratePodReadyCondition(&pod.Spec, status.Conditions, status.ContainerStatuses, status.Phase)) // 更新ContainerReady updateConditionFunc(v1.ContainersReady, GenerateContainersReadyCondition(&pod.Spec, status.ContainerStatuses, status.Phase)) // 更新状态及驱动同步 m.updateStatusInternal(pod, status, false) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
注意
SetContainerReadiness()用于probeManager的就绪探针。
# 2.4.SetContainerStartup
SetContainerStartup()用于设置pod/container的started状态,向podStatusChannel发送事件驱动状态同步。func (m *manager) SetContainerStartup(podUID types.UID, containerID kubecontainer.ContainerID, started bool) { m.podStatusesLock.Lock() defer m.podStatusesLock.Unlock() // 获取缓存的pod pod, ok := m.podManager.GetPodByUID(podUID) ... // 获取缓存的pod状态 oldStatus, found := m.podStatuses[pod.UID] ... // 获取container的状态 containerStatus, _, ok := findContainerStatus(&oldStatus.status, containerID.String()) ... // 已设置的无需更新 if containerStatus.Started != nil && *containerStatus.Started == started { return } // 拷贝更新container的状态 status := *oldStatus.status.DeepCopy() containerStatus, _, _ = findContainerStatus(&status, containerID.String()) containerStatus.Started = &started // 更新状态及驱动同步 m.updateStatusInternal(pod, status, false) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
注意
SetContainerStartup()用于probeManager的就绪探针。