From 594aff5f0ed1ecedf3441a363e56f9d5d764f2cd Mon Sep 17 00:00:00 2001 From: james Date: Tue, 23 Jun 2026 19:02:49 +0800 Subject: [PATCH 1/2] fix: return snapshot for InspectAllNodesUsage Signed-off-by: james --- pkg/device/pods.go | 4 ++-- pkg/scheduler/scheduler.go | 9 ++++++++- 2 files changed, 10 insertions(+), 3 deletions(-) diff --git a/pkg/device/pods.go b/pkg/device/pods.go index 431ea9f4f3..8da5d87706 100644 --- a/pkg/device/pods.go +++ b/pkg/device/pods.go @@ -112,8 +112,8 @@ func (m *PodManager) DelPod(pod *corev1.Pod) { } func (m *PodManager) GetPod(pod *corev1.Pod) (*PodInfo, bool) { - m.mutex.Lock() - defer m.mutex.Unlock() + m.mutex.RLock() + defer m.mutex.RUnlock() pi, ok := m.pods[pod.UID] return pi, ok diff --git a/pkg/scheduler/scheduler.go b/pkg/scheduler/scheduler.go index 92aa988b0e..f56a4314a4 100644 --- a/pkg/scheduler/scheduler.go +++ b/pkg/scheduler/scheduler.go @@ -520,7 +520,14 @@ func (s *Scheduler) WaitForCacheSync(ctx context.Context) bool { // InspectAllNodesUsage is used by metrics monitor. func (s *Scheduler) InspectAllNodesUsage() *map[string]*NodeUsage { - return &s.overviewstatus + s.lock.RLock() + defer s.lock.RUnlock() + + snapshot := make(map[string]*NodeUsage, len(s.overviewstatus)) + for nodeID, usage := range s.overviewstatus { + snapshot[nodeID] = usage.DeepCopy() + } + return &snapshot } // returns all nodes and its device memory usage, and we filter it with nodeSelector, taints, nodeAffinity From 021a6dbffc335cb6ec411e32d6ae8b88c25cd79e Mon Sep 17 00:00:00 2001 From: james Date: Wed, 24 Jun 2026 10:40:21 +0800 Subject: [PATCH 2/2] fix comments Signed-off-by: james --- pkg/scheduler/scheduler.go | 33 +++++++++++++-------------------- pkg/scheduler/scheduler_test.go | 2 +- 2 files changed, 14 insertions(+), 21 deletions(-) diff --git a/pkg/scheduler/scheduler.go b/pkg/scheduler/scheduler.go index f56a4314a4..1eb4801d3e 100644 --- a/pkg/scheduler/scheduler.go +++ b/pkg/scheduler/scheduler.go @@ -70,8 +70,6 @@ type Scheduler struct { nodeLister listerscorev1.NodeLister quotaLister listerscorev1.ResourceQuotaLister leaseLister coordinationv1.LeaseLister - //Node status returned by filter - cachedstatus map[string]*NodeUsage //Node Overview overviewstatus map[string]*NodeUsage eventRecorder record.EventRecorder @@ -84,12 +82,12 @@ type Scheduler struct { func NewScheduler() *Scheduler { klog.InfoS("Initializing HAMi scheduler") s := &Scheduler{ - stopCh: make(chan struct{}), - cachedstatus: make(map[string]*NodeUsage), - nodeNotify: make(chan struct{}, 1), - leaderNotify: make(chan struct{}, 1), - started: 0, - synced: false, + stopCh: make(chan struct{}), + overviewstatus: make(map[string]*NodeUsage), + nodeNotify: make(chan struct{}, 1), + leaderNotify: make(chan struct{}, 1), + started: 0, + synced: false, } s.nodeManager = newNodeManager() s.podManager = device.NewPodManager() @@ -231,7 +229,7 @@ func (s *Scheduler) onDelNode(obj any) { s.cleanupNodeUsage(nodeName) } -// cleanupNodeUsage removes the node from overviewstatus and cachedstatus maps +// cleanupNodeUsage removes the node from overviewstatus maps // to ensure metrics no longer report data for deleted nodes. func (s *Scheduler) cleanupNodeUsage(nodeID string) { s.lock.Lock() @@ -240,10 +238,6 @@ func (s *Scheduler) cleanupNodeUsage(nodeID string) { delete(s.overviewstatus, nodeID) klog.V(4).InfoS("Removed node from overviewstatus", "node", nodeID) } - if _, ok := s.cachedstatus[nodeID]; ok { - delete(s.cachedstatus, nodeID) - klog.V(4).InfoS("Removed node from cachedstatus", "node", nodeID) - } } func (s *Scheduler) onAddQuota(obj any) { @@ -438,11 +432,12 @@ func (s *Scheduler) register(labelSelector labels.Selector, printedLog map[strin } } } - _, _, err = s.getNodesUsage(&nodeNames, nil) + _, overallnodeMap, _, err := s.getNodesUsage(&nodeNames, nil) if err != nil { klog.ErrorS(err, "Failed to get node usage", "nodeNames", nodeNames) return } + s.overviewstatus = *overallnodeMap // Set synced to true only after getNodeUsage() succeeds s.synced = true @@ -532,13 +527,13 @@ func (s *Scheduler) InspectAllNodesUsage() *map[string]*NodeUsage { // returns all nodes and its device memory usage, and we filter it with nodeSelector, taints, nodeAffinity // unschedulerable and nodeName. -func (s *Scheduler) getNodesUsage(nodes *[]string, task *corev1.Pod) (*map[string]*NodeUsage, map[string]string, error) { +func (s *Scheduler) getNodesUsage(nodes *[]string, task *corev1.Pod) (*map[string]*NodeUsage, *map[string]*NodeUsage, map[string]string, error) { overallnodeMap := make(map[string]*NodeUsage) cachenodeMap := make(map[string]*NodeUsage) failedNodes := make(map[string]string) allNodes, err := s.ListNodes() if err != nil { - return &overallnodeMap, failedNodes, err + return &overallnodeMap, &overallnodeMap, failedNodes, err } for _, node := range allNodes { @@ -622,7 +617,6 @@ func (s *Scheduler) getNodesUsage(nodes *[]string, task *corev1.Pod) (*map[strin } klog.V(5).Infof("usage: pod %v assigned %v %v", p.Name, p.NodeID, p.Devices) } - s.overviewstatus = overallnodeMap for _, nodeID := range *nodes { node, err := s.GetNode(nodeID) if err != nil { @@ -633,8 +627,7 @@ func (s *Scheduler) getNodesUsage(nodes *[]string, task *corev1.Pod) (*map[strin } cachenodeMap[node.ID] = overallnodeMap[node.ID] } - s.cachedstatus = cachenodeMap - return &cachenodeMap, failedNodes, nil + return &cachenodeMap, &overallnodeMap, failedNodes, nil } func (s *Scheduler) getPodUsage() (map[string]device.PodUseDeviceStat, error) { @@ -767,7 +760,7 @@ func (s *Scheduler) Filter(args extenderv1.ExtenderArgs) (*extenderv1.ExtenderFi if pi, ok := s.podManager.TakeAndDeletePod(args.Pod); ok { s.quotaManager.RmUsage(args.Pod, pi.Devices) } - nodeUsage, failedNodes, err := s.getNodesUsage(args.NodeNames, args.Pod) + nodeUsage, _, failedNodes, err := s.getNodesUsage(args.NodeNames, args.Pod) if err != nil { s.recordScheduleFilterResultEvent(args.Pod, EventReasonFilteringFailed, "", err) return nil, err diff --git a/pkg/scheduler/scheduler_test.go b/pkg/scheduler/scheduler_test.go index f7576b1d98..953f13a72e 100644 --- a/pkg/scheduler/scheduler_test.go +++ b/pkg/scheduler/scheduler_test.go @@ -114,7 +114,7 @@ func Test_getNodesUsage(t *testing.T) { } nodes := make([]string, 0) nodes = append(nodes, "node1") - cachenodeMap, _, err := s.getNodesUsage(&nodes, nil) + cachenodeMap, _, _, err := s.getNodesUsage(&nodes, nil) if err != nil { t.Fatal(err) }