78
// - either all task queues are processed successfully (or)
79
// - Stop() method is called to stop the scavenger
80
>
func NewScavenger(db p.TaskManager, metricsHandler metrics.Handler, logger log.Logger) *Scavenger {
scavenger.go
81
>
stopC := make(chan struct{})
82
>
taskExecutor := executor.NewFixedSizePoolExecutor(
83
>
taskQueueBatchSize, executorMaxDeferredTasks, metricsHandler, metrics.TaskQueueScavengerScope)
84
>
lifecycleCtx, lifecycleCancel := context.WithCancel(
85
>
headers.SetCallerInfo(
86
>
context.Background(),
87
>
headers.SystemBackgroundHighCallerInfo,
88
>
),
89
>
)
90
>
return &Scavenger{
91
>
db: db,
92
>
metricsHandler: metricsHandler.WithTags(metrics.OperationTag(metrics.TaskQueueScavengerScope)),
93
>
logger: logger,
94
>
stopC: stopC,
95
>
executor: taskExecutor,
96
>
lifecycleCtx: lifecycleCtx,
97
>
lifecycleCancel: lifecycleCancel,
98
>
}
99
>
}
100
101
// Start starts the scavenger
103
>
if !atomic.CompareAndSwapInt32(&s.status, common.DaemonStatusInitialized, common.DaemonStatusStarted) {
104
return
105
}
106
>
s.logger.Info("Taskqueue scavenger starting")
scavenger.go
107
>
s.stopWG.Add(1)
108
>
s.executor.Start()
109
>
go s.run()
110
>
metrics.StartedCount.With(s.metricsHandler).Record(1)
111
>
s.logger.Info("Taskqueue scavenger started")
112
}
113
114
// Stop stops the scavenger
116
>
if !atomic.CompareAndSwapInt32(&s.status, common.DaemonStatusStarted, common.DaemonStatusStopped) {
117
return
118
}
119
>
metrics.StoppedCount.With(s.metricsHandler).Record(1)
scavenger.go
120
>
s.logger.Info("Taskqueue scavenger stopping")
121
>
s.lifecycleCancel()
122
>
close(s.stopC)
123
>
s.executor.Stop()
124
>
s.stopWG.Wait()
125
>
s.logger.Info("Taskqueue scavenger stopped")
126
}
127