229
}
230
232
>
shardID := rand.Int31()
233
>
taskID := int64(1)
234
>
235
>
task := s.newRandomReplicationTaskRow(shardID, taskID)
236
>
result, err := s.store.InsertIntoReplicationTasks(newExecutionContext(), []sqlplugin.ReplicationTasksRow{task})
237
>
s.NoError(err)
238
>
rowsAffected, err := result.RowsAffected()
239
>
s.NoError(err)
240
>
s.Equal(1, int(rowsAffected))
241
>
242
>
filter := sqlplugin.ReplicationTasksFilter{
243
>
ShardID: shardID,
244
>
TaskID: taskID,
245
>
}
246
>
result, err = s.store.DeleteFromReplicationTasks(newExecutionContext(), filter)
247
>
s.NoError(err)
248
>
rowsAffected, err = result.RowsAffected()
249
>
s.NoError(err)
250
>
s.Equal(1, int(rowsAffected))
251
>
252
>
rangeFilter := sqlplugin.ReplicationTasksRangeFilter{
253
>
ShardID: shardID,
254
>
InclusiveMinTaskID: taskID,
255
>
ExclusiveMaxTaskID: taskID + 1,
256
>
PageSize: 1,
257
>
}
258
>
rows, err := s.store.RangeSelectFromReplicationTasks(newExecutionContext(), rangeFilter)
259
>
s.NoError(err)
260
>
for index := range rows {
261
rows[index].ShardID = shardID
262
}
264
}
265