72
}
73
75
>
task := persistencespb.HistoryTask{
76
>
ShardId: int32(request.SourceShardID),
77
>
Blob: blob,
78
>
}
79
>
taskBytes, _ := task.Marshal()
80
>
blob = &commonpb.DataBlob{
81
>
EncodingType: enumspb.ENCODING_TYPE_PROTO3,
82
>
Data: taskBytes,
83
>
}
84
>
queueKey := QueueKey{
85
>
QueueType: request.QueueType,
86
>
Category: taskCategory,
87
>
SourceCluster: request.SourceCluster,
88
>
TargetCluster: request.TargetCluster,
89
>
}
90
>
91
>
message, err := m.queue.EnqueueMessage(ctx, &InternalEnqueueMessageRequest{
92
>
QueueType: request.QueueType,
93
>
QueueName: queueKey.GetQueueName(),
94
>
Blob: blob,
95
>
})
96
>
if err != nil {
97
return nil, err
98
}