983
ns *namespace.Namespace,
984
heartbeat func(details replicationTasksHeartbeatDetails),
986
>
start := time.Now()
987
>
progress := false
988
>
defer func() {
989
>
if progress {
990
// Update CheckPoint when there is a progress
991
details.CheckPoint = time.Now()
992
}
993
995
>
a.forceReplicationMetricsHandler.Timer(metrics.VerifyReplicationTasksLatency.Name()).Record(time.Since(start))
996
}()
997
998
>
ctx = metadata.NewOutgoingContext(ctx, metadata.Pairs(interceptor.DCRedirectionContextHeaderName, "false"))
activities.go
999
>
1000
>
for ; details.NextIndex < len(request.Executions); details.NextIndex++ {
1001
we := request.Executions[details.NextIndex]
1002
r, err := a.verifySingleReplicationTask(ctx, request, remotAdminClient, ns, we)