-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathStorageQueue.cs
More file actions
119 lines (100 loc) · 3.66 KB
/
StorageQueue.cs
File metadata and controls
119 lines (100 loc) · 3.66 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
using Azure.Storage.Queues;
using Azure.Storage.Queues.Models;
using System;
using System.Collections.Generic;
using System.Collections.Immutable;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
namespace CloudWorker.MessageQueue;
public class StorageQueueMessage : IQueueMessage
{
private readonly QueueClient _client;
private readonly QueueMessage _message;
private readonly TimeSpan _lease;
public StorageQueueMessage(QueueClient client, QueueMessage message, TimeSpan lease)
{
_client = client;
_message = message;
_lease = lease;
}
public string Id => _message.MessageId;
public string Content => _message.MessageText;
public Task RenewLeaseAsync()
{
return _client.UpdateMessageAsync(_message.MessageId, _message.PopReceipt, visibilityTimeout: _lease);
}
public Task ReturnAsync()
{
return _client.UpdateMessageAsync(_message.MessageId, _message.PopReceipt, visibilityTimeout: TimeSpan.Zero);
}
public Task DeleteAsync()
{
return _client.DeleteMessageAsync(_message.MessageId, _message.PopReceipt);
}
}
//NOTE: Storage Queue client has built-in retry policies by QueueClientOptions' properties Retry and RetryPolicy.
//See https://learn.microsoft.com/en-us/dotnet/api/azure.storage.queues.queueclientoptions?view=azure-dotnet
public class StorageQueueOptions : QueueOptions
{
public int? QueryInterval { get; set; } = 500; //In milliseconds.
public static StorageQueueOptions Default { get; } = new StorageQueueOptions();
public override void Merge(QueueOptions? other)
{
base.Merge(other);
if (other is StorageQueueOptions opts)
{
if (opts.QueryInterval != null)
{
QueryInterval = opts.QueryInterval;
}
}
}
public override bool Validate()
{
return base.Validate() && (QueryInterval is null || QueryInterval >= 0);
}
}
public class StorageQueue : IMessageQueue
{
public const string QueueType = "storage";
private readonly StorageQueueOptions _options;
private readonly TimeSpan _messageLease;
private QueueClient _client;
public StorageQueue(StorageQueueOptions options)
{
_options = options;
if (_options.MessageLease is null)
{
throw new ArgumentException();
}
_messageLease = TimeSpan.FromSeconds((double)_options.MessageLease);
_client = new QueueClient(_options.ConnectionString, _options.QueueName);
}
public TimeSpan MessageLease => _messageLease;
public async Task<IQueueMessage> WaitAsync(CancellationToken cancel = default)
{
var messages = await WaitBatchAsync(1, cancel).ConfigureAwait(false);
return messages[0];
}
public async Task<IReadOnlyList<IQueueMessage>> WaitBatchAsync(int batchSize, CancellationToken cancel = default)
{
var delay = _options.QueryInterval ?? StorageQueueOptions.Default.QueryInterval;
while (true)
{
var messages = await _client.ReceiveMessagesAsync(batchSize, MessageLease, cancel).ConfigureAwait(false);
if (messages.Value == null || messages.Value.Length == 0)
{
await Task.Delay(delay!.Value, cancel).ConfigureAwait(false);
}
else
{
return messages.Value.Select(msg => new StorageQueueMessage(_client, msg, MessageLease)).ToImmutableList();
}
}
}
public Task SendAsync(string message, CancellationToken cancel = default)
{
return _client.SendMessageAsync(message, MessageLease, cancellationToken: cancel);
}
}