351
}
352
353
>
if err = ValidateTasksHaveSamePriority(priority, messages.ReplicationTasks...); err != nil {
stream_receiver.go
354
// This should not happen because source side is sending task 1 by 1. Validate here just in case.
355
return NewStreamError("ReplicationTask priority check failed", err)
356
}
357
359
>
clusterName,
360
>
r.clientShardKey,
361
>
r.serverShardKey,
362
>
messages.ReplicationTasks...,
363
>
)
364
>
exclusiveHighWatermark := messages.ExclusiveHighWatermark
365
>
exclusiveHighWatermarkTime := timestamp.TimeValue(messages.ExclusiveHighWatermarkTime)
366
>
taskTracker, err := r.getTaskTracker(priority)
367
>
if err != nil {
368
// Todo: Change to write Tasks to DLQ. As resend task will not help here
369
return NewStreamError("ReplicationTask wrong priority", err)
370
}
371
372
>
submissionThreshold := r.Config.ReplicationReceiverSlowSubmissionLatencyThreshold()
stream_receiver.go
373
>
374
>
for _, task := range taskTracker.TrackTasks(WatermarkInfo{
375
>
Watermark: exclusiveHighWatermark,
376
>
Timestamp: exclusiveHighWatermarkTime,
377
>
}, convertedTasks...) {
378
schedulerPriority, err := r.getTaskSchedulerPriority(priority, task)
379
if err != nil {