769
}
770
771
>
func (a *activities) SeedReplicationQueueWithUserDataEntries(ctx context.Context, params TaskQueueUserDataReplicationParamsWithNamespace) error {
activities.go
772
>
if len(params.Namespace) == 0 {
773
return temporal.NewNonRetryableApplicationError("namespace is required", "InvalidArgument", nil)
774
}
776
params.PageSize = defaultPageSizeForTaskQueueUserDataReplication
777
}
779
params.RPS = defaultRPSForTaskQueueUserDataReplication
780
}
781
782
>
describeResponse, err := a.frontendClient.DescribeNamespace(ctx, &workflowservice.DescribeNamespaceRequest{
activities.go
783
>
Namespace: params.Namespace,
784
>
})
785
>
if err != nil {
786
return err
787
}
788
789
>
rateLimiter := quotas.NewRateLimiter(params.RPS, int(math.Ceil(params.RPS)))
activities.go
790
>
heartbeatDetails := seedReplicationQueueWithUserDataEntriesHeartbeatDetails{}
791
>
792
>
if activity.HasHeartbeatDetails(ctx) {
793
>
if err := activity.GetHeartbeatDetails(ctx, &heartbeatDetails); err != nil {
794
return temporal.NewNonRetryableApplicationError("failed to load previous heartbeat details", "TypeError", err)
795
}
796
}
797
799
>
if err := rateLimiter.Wait(ctx); err != nil {
800
return err
801
}
802
803
>
request := &persistence.ListTaskQueueUserDataEntriesRequest{
activities.go
804
>
NamespaceID: describeResponse.GetNamespaceInfo().Id,
805
>
NextPageToken: heartbeatDetails.NextPageToken,
806
>
PageSize: params.PageSize,
807
>
}
808
>
response, err := a.taskManager.ListTaskQueueUserDataEntries(ctx, request)
809
>
if err != nil {
810
a.Logger.Error("List task queue user data failed", tag.WorkflowNamespaceID(request.NamespaceID), tag.Error(err))
811
return err
812
}
814
>
if heartbeatDetails.IndexInPage > idx {
815
>
continue
816
}
818
>
activity.RecordHeartbeat(ctx, heartbeatDetails)
819
>
err = a.namespaceReplicationQueue.Publish(ctx, &replicationspb.ReplicationTask{
820
>
TaskType: enumsspb.REPLICATION_TASK_TYPE_TASK_QUEUE_USER_DATA,
821
>
Attributes: &replicationspb.ReplicationTask_TaskQueueUserDataAttributes{
822
>
TaskQueueUserDataAttributes: &replicationspb.TaskQueueUserDataAttributes{
823
>
NamespaceId: request.NamespaceID,
824
>
TaskQueueName: entry.TaskQueue,
825
>
UserData: entry.UserData.GetData(),
826
>
},
827
>
},
828
>
})
829
>
if err != nil {
830
>
a.Logger.Error("Inserting into namespace replication queue failed", tag.WorkflowNamespaceID(request.NamespaceID), tag.Error(err))
831
>
return err
832
>
}
833
}
835
>
return nil
836
>
}
837
>
heartbeatDetails.NextPageToken = response.NextPageToken
838
>
heartbeatDetails.IndexInPage = 0
839
>
activity.RecordHeartbeat(ctx, heartbeatDetails)
840
}
841
}