220
pageSize int,
221
pageToken []byte,
222
>
) ([]*replicationspb.ReplicationTask, []*replicationspb.ReplicationTaskInfo, int64, []byte, error) {
dlq_handler.go
223
>
224
>
ackLevel := r.shard.GetReplicatorDLQAckLevel(sourceCluster)
225
>
resp, err := r.shard.GetExecutionManager().GetReplicationTasksFromDLQ(ctx, &persistence.GetReplicationTasksFromDLQRequest{
226
>
GetHistoryTasksRequest: persistence.GetHistoryTasksRequest{
227
>
ShardID: r.shard.GetShardID(),
228
>
TaskCategory: tasks.CategoryReplication,
229
>
InclusiveMinTaskKey: tasks.NewImmediateKey(ackLevel + 1),
230
>
ExclusiveMaxTaskKey: tasks.NewImmediateKey(lastMessageID + 1),
231
>
BatchSize: pageSize,
232
>
NextPageToken: pageToken,
233
>
},
234
>
SourceClusterName: sourceCluster,
235
>
})
236
>
if err != nil {
237
return nil, nil, ackLevel, nil, err
238
}
240
>
241
>
remoteAdminClient, err := r.shard.GetRemoteAdminClient(sourceCluster)
242
>
if err != nil {
243
return nil, nil, ackLevel, nil, err
244
}
245
>
taskInfo := make([]*replicationspb.ReplicationTaskInfo, 0, len(resp.Tasks))
dlq_handler.go
246
>
for _, task := range resp.Tasks {
247
>
switch task := task.(type) {
248
case *tasks.SyncActivityTask:
249
taskInfo = append(taskInfo, &replicationspb.ReplicationTaskInfo{