549
endEventVersion int64,
550
newRunId string, // only verify task should pass this value
552
>
if len(newRunId) != 0 && e.replicationTask.GetTaskType() != enumsspb.REPLICATION_TASK_TYPE_VERIFY_VERSIONED_TRANSITION_TASK {
553
return serviceerror.NewInternal("newRunId should be empty for non verify task")
554
}
555
557
>
item := e.namespace.Load()
558
>
if item != nil {
559
namespaceName = item.(namespace.Name).String()
560
}
562
>
1,
563
>
metrics.OperationTag(e.metricsTag+"BackFill"),
564
>
metrics.NamespaceTag(namespaceName),
565
>
metrics.ServiceRoleTag(metrics.HistoryRoleTagValue),
566
>
)
567
>
startTime := time.Now().UTC()
568
>
defer func() {
569
>
metrics.ReplicationTasksBackFillLatency.With(e.MetricsHandler).Record(
570
>
time.Since(startTime),
571
>
metrics.OperationTag(e.metricsTag+"BackFill"),
572
>
metrics.NamespaceTag(namespaceName),
573
>
metrics.ServiceRoleTag(metrics.HistoryRoleTagValue),
574
>
)
575
>
}()
576
>
shardContext, err := e.ShardController.GetShardByNamespaceWorkflow(
577
>
namespace.ID(workflowKey.NamespaceID),
578
>
workflowKey.WorkflowID,
579
>
)
580
>
if err != nil {
581
return err
582
}
583
585
>
if err != nil {
586
return err
587
}
588
590
>
var newRunEvents []*historypb.HistoryEvent
591
>
var versionHistory []*historyspb.VersionHistoryItem
592
>
const EmptyVersion = int64(-1) // 0 is a valid event version when namespace is local
593
>
var eventsVersion = EmptyVersion
594
>
isLastEvent := false
595
>
if len(newRunId) != 0 {
596
>
iterator := e.RemoteHistoryFetcher.GetSingleWorkflowHistoryPaginatedIteratorInclusive(
597
>
ctx,
598
>
remoteCluster,
599
>
namespace.ID(workflowKey.NamespaceID),
600
>
workflowKey.WorkflowID,
601
>
newRunId,
602
>
1,
603
>
endEventVersion, // continue as new run's first event batch should have the same version as the last event of the old run
604
>
1,
605
>
endEventVersion,
606
>
)
607
>
if !iterator.HasNext() {
608
return serviceerror.NewInternalf("failed to get new run history when backfill")
609
}
611
>
if err != nil {
612
return serviceerror.NewInternalf("failed to get new run history when backfill: %v", err)
613
}
614
>
events, err := e.Serializer.DeserializeEvents(batch.RawEventBatch)
executable_task.go
615
>
if err != nil {
616
return serviceerror.NewInternalf("failed to deserailize run history events when backfill: %v", err)
617
}
619
}
620
622
>
backFillRequest := &historyi.BackfillHistoryEventsRequest{
623
>
WorkflowKey: workflowKey,
624
>
SourceClusterName: e.SourceClusterName(),
625
>
VersionedHistory: e.ReplicationTask().VersionedTransition,
626
>
VersionHistoryItems: versionHistory,
627
>
Events: eventsBatch,
628
>
}
629
>
if isLastEvent && len(newRunId) > 0 && len(newRunEvents) > 0 {
630
>
backFillRequest.NewEvents = newRunEvents
631
>
backFillRequest.NewRunID = newRunId
632
>
}
633
>
err := engine.BackfillHistoryEvents(ctx, backFillRequest)
634
>
if err != nil {
635
return serviceerror.NewInternalf("failed to backfill: %v", err)
636
}
638
>
versionHistory = nil
639
>
eventsVersion = EmptyVersion
640
>
return nil
641
}
642
>
iterator := e.RemoteHistoryFetcher.GetSingleWorkflowHistoryPaginatedIteratorInclusive(
executable_task.go
643
>
ctx,
644
>
remoteCluster,
645
>
namespace.ID(workflowKey.NamespaceID),
646
>
workflowKey.WorkflowID,
647
>
workflowKey.RunID,
648
>
startEventId,
649
>
startEventVersion,
650
>
endEventId,
651
>
endEventVersion,
652
>
)
653
>
for iterator.HasNext() {
654
>
batch, err := iterator.Next()
655
>
if err != nil {
656
return err
657
}
658
>
events, err := e.Serializer.DeserializeEvents(batch.RawEventBatch)
executable_task.go
659
>
if err != nil {
660
return err
661
}
663
return serviceerror.NewInvalidArgument("Empty batch received from remote during resend")
664
}
666
>
if !versionhistory.IsEqualVersionHistoryItems(versionHistory, batch.VersionHistory.Items) ||
667
>
(eventsVersion != EmptyVersion && eventsVersion != events[0].Version) {
668
>
err := applyFn()
669
>
if err != nil {
670
return err
671
}
672
}
673
}
675
>
if events[len(events)-1].GetEventId() == endEventId {
676
>
isLastEvent = true
677
>
}
678
>
versionHistory = batch.VersionHistory.Items
679
>
eventsVersion = events[0].Version
680
>
if len(eventsBatch) >= e.Config.ReplicationResendMaxBatchCount() {
681
err := applyFn()
682
if err != nil {