125
sourceCluster string,
126
lastMessageID int64,
128
>
129
>
ackLevel := r.shard.GetReplicatorDLQAckLevel(sourceCluster)
130
>
err := r.shard.GetExecutionManager().RangeDeleteReplicationTaskFromDLQ(
131
>
ctx,
132
>
&persistence.RangeDeleteReplicationTaskFromDLQRequest{
133
>
RangeCompleteHistoryTasksRequest: persistence.RangeCompleteHistoryTasksRequest{
134
>
ShardID: r.shard.GetShardID(),
135
>
TaskCategory: tasks.CategoryReplication,
136
>
InclusiveMinTaskKey: tasks.NewImmediateKey(ackLevel + 1),
137
>
ExclusiveMaxTaskKey: tasks.NewImmediateKey(lastMessageID + 1),
138
>
},
139
>
SourceClusterName: sourceCluster,
140
>
},
141
>
)
142
>
if err != nil {
143
return err
144
}
145
147
>
sourceCluster,
148
>
lastMessageID,
149
>
); err != nil {
150
r.logger.Error("Failed to purge history replication message", tag.Error(err))
151
// The update ack level should not block the call. Ignore the error.
152
}
154
}
155