141
}
142
144
>
numTasks := 20
145
>
pageSize := numTasks * 2
146
>
147
>
shardID := rand.Int31()
148
>
minTaskID := int64(1)
149
>
taskID := minTaskID
150
>
maxTaskID := taskID + int64(numTasks)
151
>
152
>
var tasks []sqlplugin.ReplicationTasksRow
153
>
for range numTasks {
154
>
task := s.newRandomReplicationTaskRow(shardID, taskID)
155
>
taskID++
156
>
tasks = append(tasks, task)
157
>
}
158
>
result, err := s.store.InsertIntoReplicationTasks(newExecutionContext(), tasks)
159
>
s.NoError(err)
160
>
rowsAffected, err := result.RowsAffected()
161
>
s.NoError(err)
162
>
s.Equal(numTasks, int(rowsAffected))
163
>
164
>
filter := sqlplugin.ReplicationTasksRangeFilter{
165
>
ShardID: shardID,
166
>
InclusiveMinTaskID: minTaskID,
167
>
ExclusiveMaxTaskID: maxTaskID,
168
>
PageSize: pageSize,
169
>
}
170
>
rows, err := s.store.RangeSelectFromReplicationTasks(newExecutionContext(), filter)
171
>
s.NoError(err)
172
>
for index := range rows {
173
>
rows[index].ShardID = shardID
174
>
}
175
>
s.Equal(tasks, rows)
176
}
177