188
}
189
191
>
sourceCluster := shuffle.String(testHistoryReplicationTaskDLQSourceCluster)
192
>
shardID := rand.Int31()
193
>
taskID := int64(1)
194
>
195
>
filter := sqlplugin.ReplicationDLQTasksFilter{
196
>
ShardID: shardID,
197
>
SourceClusterName: sourceCluster,
198
>
TaskID: taskID,
199
>
}
200
>
result, err := s.store.DeleteFromReplicationDLQTasks(newExecutionContext(), filter)
201
>
s.NoError(err)
202
>
rowsAffected, err := result.RowsAffected()
203
>
s.NoError(err)
204
>
s.Equal(0, int(rowsAffected))
205
>
206
>
rangeFilter := sqlplugin.ReplicationDLQTasksRangeFilter{
207
>
ShardID: shardID,
208
>
SourceClusterName: sourceCluster,
209
>
InclusiveMinTaskID: taskID,
210
>
ExclusiveMaxTaskID: taskID + 1,
211
>
PageSize: 1,
212
>
}
213
>
rows, err := s.store.RangeSelectFromReplicationDLQTasks(newExecutionContext(), rangeFilter)
214
>
s.NoError(err)
215
>
for index := range rows {
216
rows[index].ShardID = shardID
217
rows[index].SourceClusterName = sourceCluster
218
}
220
}
221