46
limiter DynamicWorkerPoolLimiter,
47
metricsHandler metrics.Handler,
49
>
stopCtx, stopFn := context.WithCancel(context.Background())
50
>
scheduler := &DynamicWorkerPoolScheduler{
51
>
stopCtx: stopCtx,
52
>
stopFn: stopFn,
53
>
limiter: limiter,
54
>
buffer: list.New(),
55
>
56
>
metricsHandler: metricsHandler,
57
>
}
58
>
scheduler.wg.Add(1)
59
>
go scheduler.exportMetricsWorker()
60
>
return scheduler
61
>
}
62
63
// InitiateShutdown aborts all buffered tasks and empties the buffer.
65
>
pool.stopFn()
66
>
pool.mu.Lock()
67
>
// Prevent any running goroutines from picking up already aborted runnables.
68
>
buffer := pool.buffer
69
>
pool.buffer = list.New()
70
>
pool.mu.Unlock()
71
>
for elem := buffer.Front(); elem != nil; elem = elem.Next() {
72
elem.Value.(Runnable).Abort()
73
}