192
}
193
194
>
mitigator := newMitigator(readerGroup, monitor, logger, metricsHandler, options.MaxReaderCount, grouper)
queue_base.go
195
>
196
>
return &queueBase{
197
>
shard: shard,
198
>
199
>
status: common.DaemonStatusInitialized,
200
>
shutdownCh: make(chan struct{}),
201
>
202
>
category: category,
203
>
options: options,
204
>
scheduler: scheduler,
205
>
rescheduler: rescheduler,
206
>
timeSource: shard.GetTimeSource(),
207
>
monitor: monitor,
208
>
mitigator: mitigator,
209
>
grouper: grouper,
210
>
logger: logger,
211
>
metricsHandler: metricsHandler,
212
>
213
>
paginationFnProvider: paginationFnProvider,
214
>
executableFactory: executableFactory,
215
>
216
>
lastRangeID: -1, // start from an invalid rangeID
217
>
exclusiveDeletionHighWatermark: exclusiveDeletionHighWatermark,
218
>
nonReadableScope: NewScope(
219
>
NewRange(exclusiveReaderHighWatermark, tasks.MaximumKey),
220
>
predicates.Universal[tasks.Task](),
221
>
),
222
>
readerRateLimiter: readerRateLimiter,
223
>
readerGroup: readerGroup,
224
>
225
>
// pollTimer and checkpointTimer are initialized on Start()
226
>
checkpointRetrier: backoff.NewRetrier(
227
>
createCheckpointRetryPolicy(),
228
>
clock.NewRealTimeSource(),
229
>
),
230
>
231
>
alertCh: monitor.AlertCh(),
232
>
}
233
}
234