Skip to content
Open
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
3 changes: 1 addition & 2 deletions pkg/scheduler/register_race_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -104,11 +104,10 @@ func Test_register_NodeCacheConcurrency(t *testing.T) {
// indexer.Update never errors for the default store and require/FailNow must not
// run outside the test goroutine, so the error is ignored.
wg.Go(func() {
printed := map[string]bool{}
sel := labels.Everything()
for v := 1; v <= rounds; v++ {
_ = indexer.Update(mkNode(v%2 == 0))
s.register(sel, printed)
s.register(sel)
}
})

Expand Down
25 changes: 16 additions & 9 deletions pkg/scheduler/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -74,8 +74,13 @@ type Scheduler struct {
leaseLister coordinationv1.LeaseLister
//Node Overview
overviewstatus map[string]*NodeUsage
eventRecorder record.EventRecorder
started uint32 // 0 = false, 1 = true
// printedLog records the nodes whose devices have already been logged at
// info level, so a re-registration logs at V(5) instead. It is pruned when
// a node is deleted, both to keep it bounded under node churn and so a node
// that returns under the same name is logged as newly added again.
printedLog map[string]bool
eventRecorder record.EventRecorder
started uint32 // 0 = false, 1 = true

lock sync.RWMutex
synced atomic.Bool
Expand All @@ -93,6 +98,7 @@ func NewScheduler() *Scheduler {
s := &Scheduler{
stopCh: make(chan struct{}),
overviewstatus: make(map[string]*NodeUsage),
printedLog: make(map[string]bool),
nodeNotify: make(chan struct{}, 1),
leaderNotify: make(chan struct{}, 1),
started: 0,
Expand Down Expand Up @@ -306,15 +312,17 @@ func (s *Scheduler) onDelNode(obj any) {
}
}

// cleanupNodeUsage removes the node from overviewstatus maps
// to ensure metrics no longer report data for deleted nodes.
// cleanupNodeUsage removes the node from the overviewstatus and printedLog maps
// to ensure metrics no longer report data for deleted nodes, and that a node
// recreated under the same name is logged as newly added rather than updated.
func (s *Scheduler) cleanupNodeUsage(nodeID string) {
s.lock.Lock()
defer s.lock.Unlock()
if _, ok := s.overviewstatus[nodeID]; ok {
delete(s.overviewstatus, nodeID)
klog.V(4).InfoS("Removed node from overviewstatus", "node", nodeID)
}
delete(s.printedLog, nodeID)
}

func (s *Scheduler) onAddQuota(obj any) {
Expand Down Expand Up @@ -435,7 +443,6 @@ func (s *Scheduler) RegisterFromNodeAnnotations() {

ticker := time.NewTicker(time.Second * 15)
defer ticker.Stop()
printedLog := map[string]bool{}
for {
select {
case <-s.nodeNotify:
Expand All @@ -452,11 +459,11 @@ func (s *Scheduler) RegisterFromNodeAnnotations() {
klog.V(5).InfoS("Scheduler not started yet, skipping ...")
continue
}
s.register(labelSelector, printedLog)
s.register(labelSelector)
}
}

func (s *Scheduler) register(labelSelector labels.Selector, printedLog map[string]bool) {
func (s *Scheduler) register(labelSelector labels.Selector) {
// Lock here to avoid setting s.synced to false, when we lost leadership, while doing register.
// 1. lost leadership before register: synced will set to false in callbacks, and register will be skipped because IsLeader() returns false
// 2. lost leadership during or after register: synced will set to true after finishing register, and callback will set it to false again after lock is acquired by callback
Expand Down Expand Up @@ -551,11 +558,11 @@ func (s *Scheduler) register(labelSelector labels.Selector, printedLog map[strin
s.addNode(val.Name, nodeInfo)
// Log the locally built nodeInfo; reading it back from s.nodes raced with onDelNode->rmNode.
if len(nodeInfo.Devices) > 0 {
if printedLog[val.Name] {
if s.printedLog[val.Name] {
klog.V(5).InfoS("Node device updated", "nodeName", val.Name, "deviceVendor", devhandsk, "nodeInfo", nodeInfo)
} else {
klog.InfoS("Node device added", "nodeName", val.Name, "deviceVendor", devhandsk, "nodeInfo", nodeInfo)
printedLog[val.Name] = true
s.printedLog[val.Name] = true
}
}
}
Expand Down
83 changes: 75 additions & 8 deletions pkg/scheduler/scheduler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1457,7 +1457,7 @@ func TestRegisterSkipsCleanupForUntrackedVendor(t *testing.T) {
})

atomic.StoreUint32(&s.started, 1)
s.register(labels.Everything(), map[string]bool{})
s.register(labels.Everything())

assert.Equal(t, mockDev.nodeCleanedUp, 0)

Expand Down Expand Up @@ -1542,7 +1542,7 @@ func TestRegisterHealthReconciliationOnDiscoveryError_Unhealthy(t *testing.T) {
})

atomic.StoreUint32(&s.started, 1)
s.register(labels.Everything(), map[string]bool{})
s.register(labels.Everything())

assert.Equal(t, 1, mockDev.nodeCleanedUp, "NodeCleanUp should be invoked when device is unhealthy even on discovery error")

Expand Down Expand Up @@ -1625,7 +1625,7 @@ func TestRegisterHealthReconciliationOnDiscoveryError_Healthy(t *testing.T) {
})

atomic.StoreUint32(&s.started, 1)
s.register(labels.Everything(), map[string]bool{})
s.register(labels.Everything())

assert.Equal(t, 0, mockDev.nodeCleanedUp, "NodeCleanUp should NOT be invoked when device is healthy")

Expand Down Expand Up @@ -1752,7 +1752,7 @@ func TestRegisterHealthReconciliationOnDiscoveryError_HeterogeneousNode(t *testi
})

atomic.StoreUint32(&s.started, 1)
s.register(labels.Everything(), map[string]bool{})
s.register(labels.Everything())

assert.Equal(t, 0, devHealthy.nodeCleanedUp)
assert.Equal(t, 1, devUnhealthyErr.nodeCleanedUp)
Expand Down Expand Up @@ -1846,7 +1846,7 @@ func TestRegisterHealthReconciliationOnDiscoveryError_Recovery(t *testing.T) {
atomic.StoreUint32(&s.started, 1)

// Cycle 1: Transient discovery error occurs when fetching node devices.
s.register(labels.Everything(), map[string]bool{})
s.register(labels.Everything())

// Verify Cycle 1 semantics:
// - NodeCleanUp was NOT called because device is healthy.
Expand Down Expand Up @@ -1875,7 +1875,7 @@ func TestRegisterHealthReconciliationOnDiscoveryError_Recovery(t *testing.T) {
mockDev.health = true
mockDev.needUpdate = true

s.register(labels.Everything(), map[string]bool{})
s.register(labels.Everything())

// Verify Cycle 2 recovery semantics:
// - Device state correctly recovers and updates in scheduler cache.
Expand All @@ -1890,7 +1890,7 @@ func TestRegisterHealthReconciliationOnDiscoveryError_Recovery(t *testing.T) {
mockDev.getNodeErr = nil
mockDev.nodeDevices = []*device.DeviceInfo{}

s.register(labels.Everything(), map[string]bool{})
s.register(labels.Everything())

// Verify Cycle 3 zero-device semantics:
// - NodeCleanUp is NOT called because device is healthy.
Expand Down Expand Up @@ -2046,7 +2046,7 @@ func Test_register_StaleDeviceVendorRemoval(t *testing.T) {
// - For node-1: vendor-B (0 devices) is removed from cache (exercises ok == true branch);
// vendor-D (0 devices) is not in cache (exercises ok == false branch).
// - For node-absent: node is absent from scheduler cache (exercises GetNode error branch).
s.register(labels.Everything(), map[string]bool{})
s.register(labels.Everything())

// Expect vendor-B to be removed from node-1 cache, while vendor-A and vendor-C remain
nodeInfo, err = s.GetNode("node-1")
Expand Down Expand Up @@ -3484,3 +3484,70 @@ func TestSchedulerIsSynced(t *testing.T) {

assert.Equal(t, true, s.IsSynced())
}

// Test_register_PrintedLogPrunedOnNodeDelete covers the printedLog bookkeeping
// across a node's full lifecycle. The map used to be a loop-local in
// RegisterFromNodeAnnotations, so onDelNode could not reach it: entries
// accumulated for every node name ever seen, and a node recreated under an old
// name was logged at V(5) as "updated" instead of at info level as "added".
func Test_register_PrintedLogPrunedOnNodeDelete(t *testing.T) {
oldDevicesMap := device.DevicesMap
t.Cleanup(func() { device.DevicesMap = oldDevicesMap })

device.DevicesMap = map[string]device.Devices{
"mock-vendor": &registerMockDevice{
nodeDevices: []*device.DeviceInfo{{
ID: "gpu-1",
DeviceVendor: "mock-vendor",
Health: true,
}},
health: true,
needUpdate: true,
},
}

s := NewScheduler()
s.stopCh = make(chan struct{})
t.Cleanup(func() { close(s.stopCh) })

oldKubeClient := client.KubeClient
client.KubeClient = fake.NewClientset()
t.Cleanup(func() { client.KubeClient = oldKubeClient })
s.kubeClient = client.KubeClient

t.Setenv("POD_NAMESPACE", "default")
t.Setenv("POD_NAME", "scheduler-0")

informerFactory := informers.NewSharedInformerFactoryWithOptions(client.KubeClient, time.Hour)
s.podLister = informerFactory.Core().V1().Pods().Lister()
s.nodeLister = informerFactory.Core().V1().Nodes().Lister()

node := &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "node-1"}}
require.NoError(t, informerFactory.Core().V1().Nodes().Informer().GetIndexer().Add(node))

// A second node stands in for the rest of the cluster: deleting node-1 must
// not disturb it.
s.lock.Lock()
s.printedLog["node-2"] = true
s.lock.Unlock()

// First registration logs the node as added and records it.
s.register(labels.Everything())
s.lock.RLock()
assert.Equal(t, true, s.printedLog["node-1"], "first registration should record the node")
s.lock.RUnlock()

// Deleting the node must drop its entry, and only its entry.
s.onDelNode(node)
s.lock.RLock()
_, stillPresent := s.printedLog["node-1"]
assert.Equal(t, false, stillPresent, "deleting a node should prune its printedLog entry")
assert.Equal(t, true, s.printedLog["node-2"], "deleting a node must not disturb other nodes")
s.lock.RUnlock()

// A node returning under the same name is treated as newly added again.
s.register(labels.Everything())
s.lock.RLock()
assert.Equal(t, true, s.printedLog["node-1"], "a recreated node should be recorded again")
s.lock.RUnlock()
}
Loading