220
}
221
223
>
sourceCluster := shuffle.String(testHistoryReplicationTaskDLQSourceCluster)
224
>
shardID := rand.Int31()
225
>
minTaskID := int64(1)
226
>
maxTaskID := int64(101)
227
>
228
>
filter := sqlplugin.ReplicationDLQTasksRangeFilter{
229
>
ShardID: shardID,
230
>
SourceClusterName: sourceCluster,
231
>
InclusiveMinTaskID: minTaskID,
232
>
ExclusiveMaxTaskID: maxTaskID,
233
>
PageSize: 0,
234
>
}
235
>
result, err := s.store.RangeDeleteFromReplicationDLQTasks(newExecutionContext(), filter)
236
>
s.NoError(err)
237
>
rowsAffected, err := result.RowsAffected()
238
>
s.NoError(err)
239
>
s.Equal(0, int(rowsAffected))
240
>
241
>
rows, err := s.store.RangeSelectFromReplicationDLQTasks(newExecutionContext(), filter)
242
>
s.NoError(err)
243
>
for index := range rows {
244
rows[index].ShardID = shardID
245
rows[index].SourceClusterName = sourceCluster
246
}
248
}
249