| 73 | } |
| 74 | |
| 75 | func New(client client.Client, coreWorkerLimit, maxWorkerLimit int, podInformer informercorev1.PodInformer) *WorkerPool { |
| 76 | workerPool := &WorkerPool{ |
| 77 | Client: client, |
| 78 | workerQueue: newQueue(), |
| 79 | requestQueue: newQueue(), |
| 80 | coreWorkerLimit: coreWorkerLimit, |
| 81 | maxWorkerLimit: maxWorkerLimit, |
| 82 | matchRequestSignal: make(chan struct{}), |
| 83 | scheme: gclient.NewSchema(), |
| 84 | scaleInQueueDuration: DefaultScaleInWorkerQueueDuration, |
| 85 | |
| 86 | queue: workqueue.NewRateLimitingQueue( |
| 87 | workqueue.NewItemExponentialFailureRateLimiter(2*time.Second, 5*time.Second), |
| 88 | ), |
| 89 | podInformer: podInformer.Informer(), |
| 90 | podLister: podInformer.Lister(), |
| 91 | } |
| 92 | |
| 93 | if _, err := podInformer.Informer().AddEventHandler( |
| 94 | cache.ResourceEventHandlerFuncs{ |
| 95 | AddFunc: func(obj interface{}) { |
| 96 | workerPool.enqueue(obj) |
| 97 | }, |
| 98 | UpdateFunc: func(_, newObj interface{}) { |
| 99 | workerPool.enqueue(newObj) |
| 100 | }, |
| 101 | DeleteFunc: func(obj interface{}) { |
| 102 | workerPool.enqueue(obj) |
| 103 | }, |
| 104 | }, |
| 105 | ); err != nil { |
| 106 | klog.ErrorS(err, "error when adding event handler to informer") |
| 107 | } |
| 108 | |
| 109 | return workerPool |
| 110 | } |
| 111 | |
| 112 | func (w *WorkerPool) enqueue(obj interface{}) { |
| 113 | pod := obj.(*corev1.Pod) |