150
}
151
153
>
numTasks := 20
154
>
pageSize := numTasks * 2
155
>
156
>
sourceCluster := shuffle.String(testHistoryReplicationTaskDLQSourceCluster)
157
>
shardID := rand.Int31()
158
>
minTaskID := int64(1)
159
>
taskID := minTaskID
160
>
maxTaskID := taskID + int64(numTasks)
161
>
162
>
var tasks []sqlplugin.ReplicationDLQTasksRow
163
>
for range numTasks {
164
>
task := s.newRandomReplicationTasksDLQRow(sourceCluster, shardID, taskID)
165
>
taskID++
166
>
tasks = append(tasks, task)
167
>
}
168
>
result, err := s.store.InsertIntoReplicationDLQTasks(newExecutionContext(), tasks)
169
>
s.NoError(err)
170
>
rowsAffected, err := result.RowsAffected()
171
>
s.NoError(err)
172
>
s.Equal(numTasks, int(rowsAffected))
173
>
174
>
filter := sqlplugin.ReplicationDLQTasksRangeFilter{
175
>
ShardID: shardID,
176
>
SourceClusterName: sourceCluster,
177
>
InclusiveMinTaskID: minTaskID,
178
>
ExclusiveMaxTaskID: maxTaskID,
179
>
PageSize: pageSize,
180
>
}
181
>
rows, err := s.store.RangeSelectFromReplicationDLQTasks(newExecutionContext(), filter)
182
>
s.NoError(err)
183
>
for index := range rows {
184
>
rows[index].ShardID = shardID
185
>
rows[index].SourceClusterName = sourceCluster
186
>
}
187
>
s.Equal(tasks, rows)
188
}
189