MCPcopy Create free account
hub / github.com/cloudtty/cloudtty / New

Function New

pkg/workerpool/worker_pool.go:75–110  ·  view source on GitHub ↗
(client client.Client, coreWorkerLimit, maxWorkerLimit int, podInformer informercorev1.PodInformer)

Source from the content-addressed store, hash-verified

73}
74
75func 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
112func (w *WorkerPool) enqueue(obj interface{}) {
113 pod := obj.(*corev1.Pod)

Callers

nothing calls this directly

Calls 5

enqueueMethod · 0.95
NewSchemaFunction · 0.92
newQueueFunction · 0.85
InformerMethod · 0.65
ListerMethod · 0.65

Tested by

no test coverage detected