184
}
185
187
>
numCreateBatch := 32
188
>
createBatchSize := 32
189
>
numTasks := int64(createBatchSize * numCreateBatch)
190
>
minTaskID := rand.Int63()
191
>
maxTaskID := minTaskID + numTasks
192
>
193
>
rangeID := rand.Int63()
194
>
taskQueue := s.createTaskQueue(rangeID)
195
>
196
>
for i := range numCreateBatch {
197
>
var tasks []*persistencespb.AllocatedTaskInfo
198
>
for j := range createBatchSize {
199
>
taskID := minTaskID + int64(i*numCreateBatch+j)
200
>
task := s.randomTask(taskID)
201
>
tasks = append(tasks, task)
202
>
}
203
>
_, err := s.taskManager.CreateTasks(s.ctx, &p.CreateTasksRequest{
204
>
TaskQueueInfo: &p.PersistedTaskQueueInfo{
205
>
RangeID: rangeID,
206
>
Data: taskQueue,
207
>
},
208
>
Tasks: tasks,
209
>
})
210
>
s.NoError(err)
211
}
212
213
>
_, err := s.taskManager.CompleteTasksLessThan(s.ctx, &p.CompleteTasksLessThanRequest{
task_queue_task.go
214
>
NamespaceID: s.namespaceID,
215
>
TaskQueueName: s.taskQueueName,
216
>
TaskType: s.taskQueueType,
217
>
ExclusiveMaxTaskID: maxTaskID + 1,
218
>
Limit: int(numTasks),
219
>
})
220
>
s.NoError(err)
221
>
222
>
resp, err := s.taskManager.GetTasks(s.ctx, &p.GetTasksRequest{
223
>
NamespaceID: s.namespaceID,
224
>
TaskQueue: s.taskQueueName,
225
>
TaskType: s.taskQueueType,
226
>
InclusiveMinTaskID: minTaskID,
227
>
ExclusiveMaxTaskID: maxTaskID + 1,
228
>
PageSize: 100,
229
>
NextPageToken: nil,
230
>
})
231
>
s.NoError(err)
232
>
protorequire.ProtoSliceEqual(s.T(), []*persistencespb.AllocatedTaskInfo{}, resp.Tasks)
233
>
s.Nil(resp.NextPageToken)
234
}
235