333
}
334
336
>
newBacklog := make([]byte, read)
337
>
backlogChanged := false
338
>
339
>
// check if we should evaluate drain state
340
>
settings := sm.settings()
341
>
target := scaleState.GetTarget()
342
>
checkDrain := target > 0 &&
343
>
sm.timeSource.Since(time.Unix(0, scaleState.GetTargetVersion())) >= settings.DrainBufferTime
344
>
info := scaleStateToInfo(scaleState, settings)
345
>
var toClear []int32
346
>
347
>
for id := range read {
348
>
callCtx, cancel := context.WithTimeout(ctx, ioTimeout)
349
>
res, err := sm.matchingClient.DescribeTaskQueuePartition(callCtx, sm.describeRequest(id))
350
>
cancel()
351
>
if err != nil {
352
continue
353
}
354
355
// update backlog count
357
>
var prev number.Compact8
358
>
if id < int32(len(prevBacklog)) {
359
prev = prevBacklog[id]
360
}
362
>
backlogChanged = backlogChanged || newBacklog[id] != prev
363
>
364
>
// check drain state for partitions in the draining range
365
>
if checkDrain &&
366
>
id >= target &&
367
>
bitSet(scaleState.BacklogState).get(id) &&
368
>
partitionIsFullyDrained(res, info) {
369
toClear = append(toClear, id)
370
}
371
}
372
374
return
375
}