142
}
143
145
>
numHistoryShards := 5
146
>
ctx := context.Background()
147
>
148
>
namespaceID := "test-namespace"
149
>
workflowID := "test-workflow-id"
150
>
workflowKey := definition.NewWorkflowKey(namespaceID, workflowID, "test-run-id")
151
>
shardID := 2
152
>
assert.Equal(t, int32(shardID), common.WorkflowIDToHistoryShard(namespaceID, workflowID, int32(numHistoryShards)))
153
>
154
>
queueKey := persistencetest.GetQueueKey(t)
155
>
_, err := manager.CreateQueue(ctx, &persistence.CreateQueueRequest{
156
>
QueueKey: queueKey,
157
>
})
158
>
require.NoError(t, err)
159
>
160
>
for i := range 2 {
161
>
task := &tasks.WorkflowTask{
162
>
WorkflowKey: workflowKey,
163
>
TaskID: int64(i + 1),
164
>
}
165
>
res, err := enqueueTask(ctx, manager, queueKey, task)
166
>
require.NoError(t, err)
167
>
assert.Equal(t, int64(persistence.FirstQueueMessageID+i), res.Metadata.ID)
168
>
}
169
171
>
for i := range 3 {
172
>
readRes, err := manager.ReadTasks(ctx, &persistence.ReadTasksRequest{
173
>
QueueKey: queueKey,
174
>
PageSize: 1,
175
>
NextPageToken: nextPageToken,
176
>
})
177
>
require.NoError(t, err)
178
>
179
>
if i < 2 {
180
>
require.Len(t, readRes.Tasks, 1)
181
>
assert.Equal(t, shardID, tasks.GetShardIDForTask(readRes.Tasks[0].Task, numHistoryShards))
182
>
assert.Equal(t, int64(i+1), readRes.Tasks[0].Task.GetTaskID())
183
>
nextPageToken = readRes.NextPageToken
184
>
} else {
185
>
assert.Empty(t, readRes.Tasks)
186
>
assert.Empty(t, readRes.NextPageToken)
187
>
}
188
}
189
}