266
manager persistence.HistoryTaskQueueManager,
267
blob *commonpb.DataBlob,
269
>
t.Helper()
270
>
271
>
queueType := persistence.QueueTypeHistoryNormal
272
>
queueKey := persistencetest.GetQueueKey(t)
273
>
queueName := queueKey.GetQueueName()
274
>
275
>
_, err := queue.CreateQueue(ctx, &persistence.InternalCreateQueueRequest{
276
>
QueueType: queueType,
277
>
QueueName: queueKey.GetQueueName(),
278
>
})
279
>
require.NoError(t, err)
280
>
historyTask := persistencespb.HistoryTask{
281
>
ShardId: 1,
282
>
Blob: blob,
283
>
}
284
>
historyTaskBytes, _ := historyTask.Marshal()
285
>
_, err = queue.EnqueueMessage(ctx, &persistence.InternalEnqueueMessageRequest{
286
>
QueueType: queueType,
287
>
QueueName: queueName,
288
>
Blob: &commonpb.DataBlob{
289
>
EncodingType: enumspb.ENCODING_TYPE_PROTO3,
290
>
Data: historyTaskBytes,
291
>
},
292
>
})
293
>
require.NoError(t, err)
294
>
295
>
_, err = manager.ReadTasks(ctx, &persistence.ReadTasksRequest{
296
>
QueueKey: queueKey,
297
>
PageSize: 1,
298
>
})
299
>
return err
300
>
}
301
302
func testHistoryTaskQueueManagerDeleteTasksErr(t *testing.T, queue persistence.QueueV2) {