425
newTaskNotificationChan <-chan struct{},
426
beginInclusiveWatermark int64,
428
>
syncStatusTimer := time.NewTimer(s.config.ReplicationStreamSendEmptyTaskDuration())
429
>
defer syncStatusTimer.Stop()
430
>
sendTasks := func() error {
431
>
endExclusiveWatermark := s.shardContext.GetQueueExclusiveHighReadWatermark(tasks.CategoryReplication).TaskID
432
>
if err := s.sendTasks(
433
>
priority,
434
>
beginInclusiveWatermark,
435
>
endExclusiveWatermark,
436
>
); err != nil {
437
return err
438
}
440
>
if !syncStatusTimer.Stop() {
441
select {
442
case <-syncStatusTimer.C: