246
pollingCluster string,
247
queryMessageID int64,
249
>
250
>
minTaskID, maxTaskID := p.taskIDsRange(queryMessageID)
251
>
replicationTasks, lastTaskID, err := p.getTasks(
252
>
ctx,
253
>
pollingCluster,
254
>
minTaskID,
255
>
maxTaskID,
256
>
)
257
>
if err != nil {
258
return nil, err
259
}
260
261
// Note this is a very rough indicator of how much the remote DC is behind on this shard.
262
>
metrics.ReplicationTasksLag.With(p.metricsHandler).Record(
ack_manager.go
263
>
maxTaskID-lastTaskID,
264
>
metrics.TargetClusterTag(pollingCluster),
265
>
metrics.OperationTag(metrics.ReplicationTaskFetcherScope),
266
>
)
267
>
268
>
metrics.ReplicationTasksFetched.With(p.metricsHandler).
269
>
Record(int64(len(replicationTasks)))
270
>
271
>
replicationEventTime := timestamppb.New(p.shardContext.GetTimeSource().Now())
272
>
if len(replicationTasks) > 0 {
273
>
replicationEventTime = replicationTasks[len(replicationTasks)-1].GetVisibilityTime()
274
>
}
275
>
return &replicationspb.ReplicationMessages{
276
>
ReplicationTasks: replicationTasks,
277
>
HasMore: lastTaskID < maxTaskID,
278
>
LastRetrievedMessageId: lastTaskID,
279
>
SyncShardStatus: &replicationspb.SyncShardStatus{
280
>
StatusTime: replicationEventTime,
281
>
},
282
>
}, nil
283
}
284