260
}
261
263
>
newMaxKey := p.shard.GetQueueExclusiveHighReadWatermark(p.category)
264
>
265
>
slices := make([]Slice, 0, 1)
266
>
if p.nonReadableScope.CanSplitByRange(newMaxKey) {
267
>
var newReadScope Scope
268
>
newReadScope, p.nonReadableScope = p.nonReadableScope.SplitByRange(newMaxKey)
269
>
slices = append(slices, NewSlice(
270
>
p.paginationFnProvider,
271
>
p.executableFactory,
272
>
p.monitor,
273
>
newReadScope,
274
>
p.grouper,
275
>
p.options.MaxPredicateSize,
276
>
p.options.ShrinkPredicateMaxPendingKeys,
277
>
p.metricsHandler,
278
>
))
279
>
}
280
281
>
reader, ok := p.readerGroup.ReaderByID(DefaultReaderId)
queue_base.go
282
>
if !ok {
283
>
p.readerGroup.NewReader(DefaultReaderId, slices...)
284
>
return
285
>
}
286
287
if now := p.timeSource.Now(); now.After(p.nextForceNewSliceTime) {