736
func (r *TaskGeneratorImpl) GenerateHistoryReplicationTasks(
737
eventBatches [][]*historypb.HistoryEvent,
739
>
if len(eventBatches) == 0 {
740
return nil, nil
741
}
743
>
if len(events) == 0 {
744
return nil, serviceerror.NewInternal("TaskGeneratorImpl encountered empty event batch")
745
}
746
}
747
749
>
firstEvent := firstBatch[0]
750
>
lastBatch := eventBatches[len(eventBatches)-1]
751
>
lastEvent := lastBatch[len(lastBatch)-1]
752
>
version := firstEvent.GetVersion()
753
>
for _, events := range eventBatches {
754
>
if events[0].GetVersion() != version || events[len(events)-1].GetVersion() != version {
755
return nil, serviceerror.NewInternal("TaskGeneratorImpl encountered contradicting versions")
756
}
757
}
758
760
>
&tasks.HistoryReplicationTask{
761
>
// TaskID, VisibilityTimestamp is set by shard
762
>
WorkflowKey: r.mutableState.GetWorkflowKey(),
763
>
FirstEventID: firstEvent.GetEventId(),
764
>
NextEventID: lastEvent.GetEventId() + 1,
765
>
Version: version,
766
>
},
767
>
}, nil
768
}
769