17
// TestInvoke is a library test function intended to be invoked from a persistence test suite. It works by
18
// enqueueing a task into the DLQ and then calling [getdlqtasks.Invoke] to verify that the right task is returned.
19
>
func TestInvoke(t *testing.T, manager persistence.HistoryTaskQueueManager) {
apitest.go
20
>
ctx := context.Background()
21
>
inTask := &tasks.WorkflowTask{
22
>
TaskID: 42,
23
>
}
24
>
sourceCluster := "test-source-cluster-" + t.Name()
25
>
targetCluster := "test-target-cluster-" + t.Name()
26
>
queueType := persistence.QueueTypeHistoryDLQ
27
>
_, err := manager.CreateQueue(ctx, &persistence.CreateQueueRequest{
28
>
QueueKey: persistence.QueueKey{
29
>
QueueType: queueType,
30
>
Category: inTask.GetCategory(),
31
>
SourceCluster: sourceCluster,
32
>
TargetCluster: targetCluster,
33
>
},
34
>
})
35
>
require.NoError(t, err)
36
>
_, err = manager.EnqueueTask(ctx, &persistence.EnqueueTaskRequest{
37
>
QueueType: queueType,
38
>
SourceCluster: sourceCluster,
39
>
TargetCluster: targetCluster,
40
>
Task: inTask,
41
>
SourceShardID: 1,
42
>
})
43
>
require.NoError(t, err)
44
>
res, err := getdlqtasks.Invoke(
45
>
context.Background(),
46
>
manager,
47
>
tasks.NewDefaultTaskCategoryRegistry(),
48
>
&historyservice.GetDLQTasksRequest{
49
>
DlqKey: &commonspb.HistoryDLQKey{
50
>
TaskCategory: int32(tasks.CategoryTransfer.ID()),
51
>
SourceCluster: sourceCluster,
52
>
TargetCluster: targetCluster,
53
>
},
54
>
PageSize: 1,
55
>
},
56
>
)
57
>
require.NoError(t, err)
58
>
require.Len(t, res.DlqTasks, 1)
59
>
assert.Equal(t, int64(persistence.FirstQueueMessageID), res.DlqTasks[0].Metadata.MessageId)
60
>
assert.Equal(t, 1, int(res.DlqTasks[0].Payload.ShardId))
61
>
serializer := serialization.NewSerializer()
62
>
outTask, err := serializer.DeserializeTask(tasks.CategoryTransfer, res.DlqTasks[0].Payload.Blob)
63
>
require.NoError(t, err)
64
>
assert.Equal(t, inTask, outTask)
65
>
}