diff --git a/controller.go b/controller.go index 6a3924e..d2c1f78 100644 --- a/controller.go +++ b/controller.go @@ -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{}, ) diff --git a/main.go b/main.go index 1abff1d..2953d58 100644 --- a/main.go +++ b/main.go @@ -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{} @@ -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() @@ -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 } @@ -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 @@ -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 }()