595
metrics.ReplicationRateLimitLatency.With(s.metrics).Record(time.Since(rlStartTime), metrics.OperationTag(TaskOperationTag(task)))
596
}
598
s.emitReplicationSent(task, item)
599
}
600
>
if err := s.sendToStream(&historyservice.StreamWorkflowReplicationMessagesResponse{
stream_sender.go
601
>
Attributes: &historyservice.StreamWorkflowReplicationMessagesResponse_Messages{
602
>
Messages: &replicationspb.WorkflowReplicationMessages{
603
>
ReplicationTasks: []*replicationspb.ReplicationTask{task},
604
>
ExclusiveHighWatermark: task.SourceTaskId + 1,
605
>
ExclusiveHighWatermarkTime: task.VisibilityTime,
606
>
Priority: priority,
607
>
},
608
>
},
609
>
}); err != nil {
610
return s.recordRetry(item, attempt, fmt.Errorf("send: %w", err))
611
}
613
>
metrics.ReplicationTasksSend.With(s.metrics).Record(
614
>
int64(1),
615
>
metrics.FromClusterIDTag(s.serverShardKey.ClusterID),
616
>
metrics.ToClusterIDTag(s.clientShardKey.ClusterID),
617
>
metrics.OperationTag(TaskOperationTag(task)),
618
>
)
619
>
return nil
620
}
621
622
>
retryPolicy := backoff.NewExponentialRetryPolicy(s.config.ReplicationStreamSenderErrorRetryWait()).
stream_sender.go
623
>
WithBackoffCoefficient(s.config.ReplicationStreamSenderErrorRetryBackoffCoefficient()).
624
>
WithMaximumInterval(s.config.ReplicationStreamSenderErrorRetryMaxInterval()).
625
>
WithMaximumAttempts(s.config.ReplicationStreamSenderErrorRetryMaxAttempts()).
626
>
WithExpirationInterval(s.config.ReplicationStreamSenderErrorRetryExpiration())
627
>
628
>
err = backoff.ThrottleRetry(operation, retryPolicy, isRetryableError)
629
>
metrics.ReplicationTaskSendAttempt.With(s.metrics).Record(
630
>
attempt,
631
>
metrics.FromClusterIDTag(s.serverShardKey.ClusterID),
632
>
metrics.ToClusterIDTag(s.clientShardKey.ClusterID),
633
>
metrics.OperationTag(TaskOperationTagFromTask(item.GetType())),
634
>
metrics.ReplicationTaskPriorityTag(priority),
635
>
)
636
>
metrics.ReplicationTaskSendLatency.With(s.metrics).Record(
637
>
time.Since(item.GetVisibilityTime()),
638
>
metrics.FromClusterIDTag(s.serverShardKey.ClusterID),
639
>
metrics.ToClusterIDTag(s.clientShardKey.ClusterID),
640
>
metrics.OperationTag(TaskOperationTagFromTask(item.GetType())),
641
>
metrics.ReplicationTaskPriorityTag(priority),
642
>
)
643
>
if err != nil {
644
metrics.ReplicationTaskSendError.With(s.metrics).Record(
645
int64(1),