248
}
249
251
>
sourceCluster := shuffle.String(testHistoryReplicationTaskDLQSourceCluster)
252
>
shardID := rand.Int31()
253
>
taskID := int64(1)
254
>
255
>
task := s.newRandomReplicationTasksDLQRow(sourceCluster, shardID, taskID)
256
>
result, err := s.store.InsertIntoReplicationDLQTasks(newExecutionContext(), []sqlplugin.ReplicationDLQTasksRow{task})
257
>
s.NoError(err)
258
>
rowsAffected, err := result.RowsAffected()
259
>
s.NoError(err)
260
>
s.Equal(1, int(rowsAffected))
261
>
262
>
filter := sqlplugin.ReplicationDLQTasksFilter{
263
>
ShardID: shardID,
264
>
SourceClusterName: sourceCluster,
265
>
TaskID: taskID,
266
>
}
267
>
result, err = s.store.DeleteFromReplicationDLQTasks(newExecutionContext(), filter)
268
>
s.NoError(err)
269
>
rowsAffected, err = result.RowsAffected()
270
>
s.NoError(err)
271
>
s.Equal(1, int(rowsAffected))
272
>
273
>
rangeFilter := sqlplugin.ReplicationDLQTasksRangeFilter{
274
>
ShardID: shardID,
275
>
SourceClusterName: sourceCluster,
276
>
InclusiveMinTaskID: taskID,
277
>
ExclusiveMaxTaskID: taskID + 1,
278
>
PageSize: 1,
279
>
}
280
>
rows, err := s.store.RangeSelectFromReplicationDLQTasks(newExecutionContext(), rangeFilter)
281
>
s.NoError(err)
282
>
for index := range rows {
283
rows[index].ShardID = shardID
284
rows[index].SourceClusterName = sourceCluster
285
}
287
}
288