160
pageSize int,
161
pageToken []byte,
163
>
164
>
replicationTasks, _, ackLevel, token, err := r.readMessagesWithAckLevel(
165
>
ctx,
166
>
sourceCluster,
167
>
lastMessageID,
168
>
pageSize,
169
>
pageToken,
170
>
)
171
>
if err != nil {
172
return nil, err
173
}
174
175
>
taskExecutor, err := r.getOrCreateTaskExecutor(sourceCluster)
dlq_handler.go
176
>
if err != nil {
177
return nil, err
178
}
179
181
>
if err := taskExecutor.Execute(
182
>
ctx,
183
>
task,
184
>
true,
185
>
); err != nil {
186
return nil, err
187
}
188
}
189
190
>
err = r.shard.GetExecutionManager().RangeDeleteReplicationTaskFromDLQ(
dlq_handler.go
191
>
ctx,
192
>
&persistence.RangeDeleteReplicationTaskFromDLQRequest{
193
>
RangeCompleteHistoryTasksRequest: persistence.RangeCompleteHistoryTasksRequest{
194
>
ShardID: r.shard.GetShardID(),
195
>
TaskCategory: tasks.CategoryReplication,
196
>
InclusiveMinTaskKey: tasks.NewImmediateKey(ackLevel + 1),
197
>
ExclusiveMaxTaskKey: tasks.NewImmediateKey(lastMessageID + 1),
198
>
},
199
>
SourceClusterName: sourceCluster,
200
>
},
201
>
)
202
>
if err != nil {
203
return nil, err
204
}
205
207
>
sourceCluster,
208
>
lastMessageID,
209
>
); err != nil {
210
r.logger.Error("Failed to purge history replication message", tag.Error(err))
211
// The update ack level should not block the call. Ignore the error.
212
}
214
}
215