264
}
265
267
>
numTasks := 20
268
>
pageSize := numTasks * 2
269
>
270
>
shardID := rand.Int31()
271
>
minTaskID := int64(1)
272
>
taskID := minTaskID
273
>
maxTaskID := taskID + int64(numTasks)
274
>
275
>
var tasks []sqlplugin.ReplicationTasksRow
276
>
for range numTasks {
277
>
task := s.newRandomReplicationTaskRow(shardID, taskID)
278
>
taskID++
279
>
tasks = append(tasks, task)
280
>
}
281
>
result, err := s.store.InsertIntoReplicationTasks(newExecutionContext(), tasks)
282
>
s.NoError(err)
283
>
rowsAffected, err := result.RowsAffected()
284
>
s.NoError(err)
285
>
s.Equal(numTasks, int(rowsAffected))
286
>
287
>
filter := sqlplugin.ReplicationTasksRangeFilter{
288
>
ShardID: shardID,
289
>
InclusiveMinTaskID: minTaskID,
290
>
ExclusiveMaxTaskID: maxTaskID,
291
>
PageSize: pageSize,
292
>
}
293
>
result, err = s.store.RangeDeleteFromReplicationTasks(newExecutionContext(), filter)
294
>
s.NoError(err)
295
>
rowsAffected, err = result.RowsAffected()
296
>
s.NoError(err)
297
>
s.Equal(numTasks, int(rowsAffected))
298
>
299
>
rows, err := s.store.RangeSelectFromReplicationTasks(newExecutionContext(), filter)
300
>
s.NoError(err)
301
>
for index := range rows {
302
rows[index].ShardID = shardID
303
}
305
}
306