20
ctx := context.Background()
21
t.Run("HappyPath", func(t *testing.T) {
23
>
24
>
queueKey := persistencetest.GetQueueKey(t, persistencetest.WithQueueType(persistence.QueueTypeHistoryDLQ))
25
>
_, err := manager.CreateQueue(ctx, &persistence.CreateQueueRequest{
26
>
QueueKey: queueKey,
27
>
})
28
>
require.NoError(t, err)
29
>
for range 3 {
30
>
_, err := manager.EnqueueTask(ctx, &persistence.EnqueueTaskRequest{
31
>
QueueType: queueKey.QueueType,
32
>
SourceCluster: queueKey.SourceCluster,
33
>
TargetCluster: queueKey.TargetCluster,
34
>
Task: &tasks.WorkflowTask{},
35
>
SourceShardID: 1,
36
>
})
37
>
require.NoError(t, err)
38
>
}
39
>
_, err = deletedlqtasks.Invoke(ctx, manager, &historyservice.DeleteDLQTasksRequest{
40
>
DlqKey: &commonspb.HistoryDLQKey{
41
>
TaskCategory: int32(queueKey.Category.ID()),
42
>
SourceCluster: queueKey.SourceCluster,
43
>
TargetCluster: queueKey.TargetCluster,
44
>
},
45
>
InclusiveMaxTaskMetadata: &commonspb.HistoryDLQTaskMetadata{
46
>
MessageId: persistence.FirstQueueMessageID + 1,
47
>
},
48
>
}, tasks.NewDefaultTaskCategoryRegistry())
49
>
require.NoError(t, err)
50
>
resp, err := manager.ReadRawTasks(ctx, &persistence.ReadRawTasksRequest{
51
>
QueueKey: queueKey,
52
>
PageSize: 10,
53
>
})
54
>
require.NoError(t, err)
55
>
require.Len(t, resp.Tasks, 1)
56
>
assert.Equal(t, int64(persistence.FirstQueueMessageID+2), resp.Tasks[0].MessageMetadata.ID)
57
})
58
t.Run("QueueDoesNotExist", func(t *testing.T) {