575
}
576
if s.config.ReplicationEnableRateLimit() && task.Priority == enumsspb.TASK_PRIORITY_LOW {
577
>
nsName, err := s.shardContext.GetNamespaceRegistry().GetNamespaceName(
stream_sender.go
578
>
namespace.ID(item.GetNamespaceID()),
579
>
)
580
>
if err != nil {
581
// if there is error, then blindly send the task, better safe than sorry
582
nsName = namespace.EmptyName
583
}
585
>
if err := s.ssRateLimiter.Wait(s.server.Context(), quotas.NewRequest(
586
>
task.TaskType.String(),
587
>
taskSchedulerToken,
588
>
nsName.String(),
589
>
headers.SystemPreemptableCallerInfo.CallerType,
590
>
0,
591
>
"",
592
>
)); err != nil {
593
return s.recordRetry(item, attempt, fmt.Errorf("rate_limit: %w", err))
594
}
595
>
metrics.ReplicationRateLimitLatency.With(s.metrics).Record(time.Since(rlStartTime), metrics.OperationTag(TaskOperationTag(task)))
stream_sender.go
596
}
597
if s.config.EmitReplicationLifecycleEvents() {