253
}
254
255
>
func (s *queueMessageSuite) TestInsertDeleteSelect_Multiple() {
queue_message.go
256
>
numMessages := 20
257
>
pageSize := numMessages
258
>
259
>
queueType := persistence.NamespaceReplicationQueueType
260
>
minMessageID := rand.Int63()
261
>
messageID := minMessageID + 1
262
>
maxMessageID := messageID + int64(numMessages)
263
>
264
>
var messages []sqlplugin.QueueMessageRow
265
>
for range numMessages {
266
>
message := s.newRandomQueueMessageRow(queueType, messageID)
267
>
messageID++
268
>
messages = append(messages, message)
269
>
}
270
>
result, err := s.store.InsertIntoMessages(newExecutionContext(), messages)
271
>
s.NoError(err)
272
>
rowsAffected, err := result.RowsAffected()
273
>
s.NoError(err)
274
>
s.Equal(numMessages, int(rowsAffected))
275
>
276
>
filter := sqlplugin.QueueMessagesRangeFilter{
277
>
QueueType: queueType,
278
>
MinMessageID: minMessageID,
279
>
MaxMessageID: maxMessageID,
280
>
PageSize: 0,
281
>
}
282
>
result, err = s.store.RangeDeleteFromMessages(newExecutionContext(), filter)
283
>
s.NoError(err)
284
>
rowsAffected, err = result.RowsAffected()
285
>
s.NoError(err)
286
>
s.Equal(numMessages, int(rowsAffected))
287
>
288
>
filter.PageSize = pageSize
289
>
rows, err := s.store.RangeSelectFromMessages(newExecutionContext(), filter)
290
>
s.NoError(err)
291
>
for index := range rows {
292
rows[index].QueueType = queueType
293
}
295
}
296