287
}
288
290
>
numTasks := 20
291
>
pageSize := numTasks * 2
292
>
293
>
sourceCluster := shuffle.String(testHistoryReplicationTaskDLQSourceCluster)
294
>
shardID := rand.Int31()
295
>
minTaskID := int64(1)
296
>
taskID := minTaskID
297
>
maxTaskID := taskID + int64(numTasks)
298
>
299
>
var tasks []sqlplugin.ReplicationDLQTasksRow
300
>
for range numTasks {
301
>
task := s.newRandomReplicationTasksDLQRow(sourceCluster, shardID, taskID)
302
>
taskID++
303
>
tasks = append(tasks, task)
304
>
}
305
>
result, err := s.store.InsertIntoReplicationDLQTasks(newExecutionContext(), tasks)
306
>
s.NoError(err)
307
>
rowsAffected, err := result.RowsAffected()
308
>
s.NoError(err)
309
>
s.Equal(numTasks, int(rowsAffected))
310
>
311
>
filter := sqlplugin.ReplicationDLQTasksRangeFilter{
312
>
ShardID: shardID,
313
>
SourceClusterName: sourceCluster,
314
>
InclusiveMinTaskID: minTaskID,
315
>
ExclusiveMaxTaskID: maxTaskID,
316
>
PageSize: pageSize,
317
>
}
318
>
result, err = s.store.RangeDeleteFromReplicationDLQTasks(newExecutionContext(), filter)
319
>
s.NoError(err)
320
>
rowsAffected, err = result.RowsAffected()
321
>
s.NoError(err)
322
>
s.Equal(numTasks, int(rowsAffected))
323
>
324
>
rows, err := s.store.RangeSelectFromReplicationDLQTasks(newExecutionContext(), filter)
325
>
s.NoError(err)
326
>
for index := range rows {
327
rows[index].ShardID = shardID
328
rows[index].SourceClusterName = sourceCluster
329
}
331
}
332