518
}
519
521
>
for _, task := range resp.GetReplicationTasks() {
522
>
tasks = append(tasks, task)
523
>
}
524
>
p.maxRxReceivedTaskID = resp.GetLastRetrievedMessageId()
525
>
if len(tasks) == 0 {
526
// Update processed timestamp to the source cluster time when there is no replication task
527
p.maxRxProcessedTimestamp = timestamp.TimeValue(resp.GetSyncShardStatus().GetStatusTime())
528
}
529
531
p.rxTaskBackoff = time.Duration(0)
533
p.rxTaskBackoff = p.config.ReplicationTaskProcessorNoTaskRetryWait(p.sourceShardID)
534
}
536
537
case <-p.shutdownChan: