136
}
137
139
>
numCreateBatch := 32
140
>
createBatchSize := 32
141
>
numTasks := int64(createBatchSize * numCreateBatch)
142
>
minTaskID := rand.Int63()
143
>
maxTaskID := minTaskID + numTasks
144
>
145
>
rangeID := rand.Int63()
146
>
taskQueue := s.createTaskQueue(rangeID)
147
>
148
>
var expectedTasks []*persistencespb.AllocatedTaskInfo
149
>
for i := range numCreateBatch {
150
>
var tasks []*persistencespb.AllocatedTaskInfo
151
>
for j := range createBatchSize {
152
>
taskID := minTaskID + int64(i*numCreateBatch+j)
153
>
task := s.randomTask(taskID)
154
>
tasks = append(tasks, task)
155
>
expectedTasks = append(expectedTasks, task)
156
>
}
157
>
_, err := s.taskManager.CreateTasks(s.ctx, &p.CreateTasksRequest{
158
>
TaskQueueInfo: &p.PersistedTaskQueueInfo{
159
>
RangeID: rangeID,
160
>
Data: taskQueue,
161
>
},
162
>
Tasks: tasks,
163
>
})
164
>
s.NoError(err)
165
}
166
168
>
var actualTasks []*persistencespb.AllocatedTaskInfo
169
>
for doContinue := true; doContinue; doContinue = len(token) > 0 {
170
>
resp, err := s.taskManager.GetTasks(s.ctx, &p.GetTasksRequest{
171
>
NamespaceID: s.namespaceID,
172
>
TaskQueue: s.taskQueueName,
173
>
TaskType: s.taskQueueType,
174
>
InclusiveMinTaskID: minTaskID,
175
>
ExclusiveMaxTaskID: maxTaskID + 1,
176
>
PageSize: 1,
177
>
NextPageToken: token,
178
>
})
179
>
s.NoError(err)
180
>
token = resp.NextPageToken
181
>
actualTasks = append(actualTasks, resp.Tasks...)
182
>
}
183
>
protorequire.ProtoSliceEqual(s.T(), expectedTasks, actualTasks)
184
}
185