72
})
73
t.Run("DeleteDLQTasks", func(t *testing.T) {
75
>
queueKey := persistencetest.GetQueueKey(t, persistencetest.WithQueueType(persistence.QueueTypeHistoryDLQ))
76
>
_, err := historyTaskQueueManager.CreateQueue(context.Background(), &persistence.CreateQueueRequest{
77
>
QueueKey: queueKey,
78
>
})
79
>
require.NoError(t, err)
80
>
enqueueTasks(t, historyTaskQueueManager, 2, queueKey.SourceCluster, queueKey.TargetCluster)
81
>
dlqKey := &commonspb.HistoryDLQKey{
82
>
TaskCategory: int32(tasks.CategoryTransfer.ID()),
83
>
SourceCluster: queueKey.SourceCluster,
84
>
TargetCluster: queueKey.TargetCluster,
85
>
}
86
>
_, err = client.DeleteDLQTasks(context.Background(), &historyservice.DeleteDLQTasksRequest{
87
>
DlqKey: dlqKey,
88
>
InclusiveMaxTaskMetadata: &commonspb.HistoryDLQTaskMetadata{
89
>
MessageId: persistence.FirstQueueMessageID,
90
>
},
91
>
})
92
>
require.NoError(t, err)
93
>
res, err := client.GetDLQTasks(context.Background(), &historyservice.GetDLQTasksRequest{
94
>
DlqKey: dlqKey,
95
>
PageSize: 10,
96
>
})
97
>
require.NoError(t, err)
98
>
assert.Equal(t, 1, len(res.DlqTasks))
99
>
assert.Equal(t, int64(persistence.FirstQueueMessageID+1), res.DlqTasks[0].Metadata.MessageId)
100
>
})
101
102
t.Cleanup(func() {