Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions pkg/device/pods.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
42 changes: 21 additions & 21 deletions pkg/scheduler/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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()
Expand Down Expand Up @@ -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()
Expand All @@ -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) {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -520,18 +515,25 @@ 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
}
Comment thread
DSFans2014 marked this conversation as resolved.

// 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 {
Expand Down Expand Up @@ -615,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 {
Expand All @@ -626,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) {
Expand Down Expand Up @@ -760,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
Expand Down
2 changes: 1 addition & 1 deletion pkg/scheduler/scheduler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down
Loading