285
}
286
288
>
numTasks := 20
289
>
pageSize := numTasks
290
>
291
>
shardID := rand.Int31()
292
>
timestamp := s.now()
293
>
minTimestamp := timestamp
294
>
taskID := int64(1)
295
>
maxTimestamp := timestamp.Add(time.Duration(numTasks) * time.Millisecond)
296
>
297
>
var tasks []sqlplugin.TimerTasksRow
298
>
for range numTasks {
299
>
task := s.newRandomTimerTaskRow(shardID, timestamp, taskID)
300
>
timestamp = timestamp.Add(time.Millisecond)
301
>
taskID++
302
>
tasks = append(tasks, task)
303
>
}
304
>
result, err := s.store.InsertIntoTimerTasks(newExecutionContext(), tasks)
305
>
s.NoError(err)
306
>
rowsAffected, err := result.RowsAffected()
307
>
s.NoError(err)
308
>
s.Equal(numTasks, int(rowsAffected))
309
>
310
>
filter := sqlplugin.TimerTasksRangeFilter{
311
>
ShardID: shardID,
312
>
InclusiveMinVisibilityTimestamp: minTimestamp,
313
>
ExclusiveMaxVisibilityTimestamp: maxTimestamp,
314
>
PageSize: 0,
315
>
}
316
>
result, err = s.store.RangeDeleteFromTimerTasks(newExecutionContext(), filter)
317
>
s.NoError(err)
318
>
rowsAffected, err = result.RowsAffected()
319
>
s.NoError(err)
320
>
s.Equal(numTasks, int(rowsAffected))
321
>
322
>
filter.PageSize = pageSize
323
>
rows, err := s.store.RangeSelectFromTimerTasks(newExecutionContext(), filter)
324
>
s.NoError(err)
325
>
for index := range rows {
326
rows[index].ShardID = shardID
327
}
329
}
330