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
18 changes: 18 additions & 0 deletions controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,24 @@ func NewController(handler *func(*core_v1.Node), dsHandler *func(ops string, ds
(*handler)(node)
}
},
DeleteFunc: func(obj interface{}) {
node, ok := obj.(*core_v1.Node)
if !ok {
tombstone, ok := obj.(cache.DeletedFinalStateUnknown)
if !ok {
return
}
node, ok = tombstone.Obj.(*core_v1.Node)
if !ok {
return
}
}
// Node is gone; tear down its pod watcher so the goroutine,
// informer and cached pods don't leak.
isHandling.Lock()
stopNodeWatcher(node.Name)
isHandling.Unlock()
},
},
cache.Indexers{},
)
Expand Down
46 changes: 33 additions & 13 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,8 +28,9 @@ import (
)

var (
dsList = make(map[string]*v1.DaemonSet)
podStore = make(map[string]*cache.Indexer)
dsList = make(map[string]*v1.DaemonSet)
podStore = make(map[string]*cache.Indexer)
podStopChans = make(map[string]chan struct{})

isHandling = sync.Mutex{}

Expand Down Expand Up @@ -73,6 +74,16 @@ func getClientset() (*kubernetes.Clientset, error) {
return clientset, nil
}

// stopNodeWatcher tears down the per-node pod watcher and forgets the node.
// Caller must hold isHandling. No-op for an untracked node; closes exactly once.
func stopNodeWatcher(nodeName string) {
if ch, ok := podStopChans[nodeName]; ok {
close(ch)
delete(podStopChans, nodeName)
}
delete(podStore, nodeName)
}

func checkDSStatus(node *core_v1.Node, opts config.Ops) (bool, error) {
isHandling.Lock()
defer isHandling.Unlock()
Expand Down Expand Up @@ -260,8 +271,9 @@ func main() {
podHandler := func(pod *core_v1.Pod, node *core_v1.Node, podStopChan chan struct{}) {
// check if the node has required taint, if not, ignore this node
if !taintutils.TaintExists(node.Spec.Taints, notReadyTaint) {
delete(podStore, node.Name)
podStopChan <- struct{}{}
isHandling.Lock()
stopNodeWatcher(node.Name)
isHandling.Unlock()
return
}

Expand Down Expand Up @@ -295,8 +307,9 @@ func main() {
if err != nil {
logrus.Errorf("Failed to remove taint from node: %v", err)
}
delete(podStore, node.Name)
podStopChan <- struct{}{}
isHandling.Lock()
stopNodeWatcher(node.Name)
isHandling.Unlock()
}

hasSynced := false
Expand All @@ -305,25 +318,32 @@ func main() {
if !taintutils.TaintExists(node.Spec.Taints, notReadyTaint) || !hasSynced {
return
}
// there is already a goroutine handling this node
isHandling.Lock()
defer isHandling.Unlock()
if _, ok := podStore[node.Name]; ok {
// Reserve the slot under the lock so a burst of node updates can't
// spawn two watchers (and two informers) for the same node.
if _, ok := podStopChans[node.Name]; ok {
return
}
podStopCh := make(chan struct{})
podStopChans[node.Name] = podStopCh

go func() {
// Handle termination
podStopCh := make(chan struct{})
defer close(podStopCh)
// launch pod watcher
podIndexer, podInformer, err := NewPodInformer(&podHandler, node, podStopCh)
if err != nil {
logrus.Fatalf("Error creating pod watcher: %v", err)
logrus.Errorf("Error creating pod watcher for node %s: %v", node.Name, err)
isHandling.Lock()
stopNodeWatcher(node.Name)
isHandling.Unlock()
return
}
go podInformer.Run(podStopCh)
isHandling.Lock()
podStore[node.Name] = &podIndexer
// Skip if the node was deleted while the informer was starting.
if _, ok := podStopChans[node.Name]; ok {
podStore[node.Name] = &podIndexer
}
isHandling.Unlock()
<-podStopCh
}()
Expand Down