-
Notifications
You must be signed in to change notification settings - Fork 1
/
AsyncSingleProducerSingleConsumerQueue.cs
81 lines (62 loc) · 2.58 KB
/
AsyncSingleProducerSingleConsumerQueue.cs
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
#if NET6_0_OR_GREATER
using System;
using System.Threading;
using System.Threading.Tasks;
using DotNext.Threading.Tasks;
namespace Nucs.Collections;
/// <summary>
/// Provides a producer/consumer queue safe to be used by only one producer and one consumer concurrently, allowing reader to wait for new items asyncronously.
/// </summary>
/// <typeparam name="T">Specifies the type of data contained in the queue.</typeparam>
public class AsyncSingleProducerSingleConsumerQueue<T> : IDisposable {
private readonly SingleProducerSingleConsumerQueue<T> _queues;
private readonly ValueTaskCompletionSource _notifier;
private volatile int _count;
private volatile int _notificationWaiters;
public int Count => _count;
public bool IsEmpty => _count == 0;
public AsyncSingleProducerSingleConsumerQueue(int initialCapacity = 512, bool runContinuationsAsynchronously = true) {
_queues = new SingleProducerSingleConsumerQueue<T>(initialCapacity);
_notifier = new ValueTaskCompletionSource(runContinuationsAsynchronously);
}
#region Writer
public void Enqueue(T item) {
_queues.Enqueue(ref item);
if (Interlocked.Increment(ref _count) == 1 && _notificationWaiters > 0)
_notifier.TrySetResult();
}
public void Enqueue(ref T item) {
_queues.Enqueue(ref item);
if (Interlocked.Increment(ref _count) == 1 && _notificationWaiters > 0)
_notifier.TrySetResult();
}
#endregion
#region Reader
public bool TryDequeue(out T result) {
if (!_queues.TryDequeue(out result))
return false;
if (Interlocked.Decrement(ref _count) == 0 && _notifier.IsCompleted)
_notifier.Reset();
return true;
}
public async ValueTask WaitForReadAsync(CancellationToken cancellationToken = default) {
if (_count > 0)
return; //completed
Interlocked.Increment(ref _notificationWaiters);
await _notifier.CreateTask(Timeout.InfiniteTimeSpan, cancellationToken).ConfigureAwait(false);
Interlocked.Decrement(ref _notificationWaiters);
}
public async ValueTask WaitForReadAsync(TimeSpan timeout, CancellationToken cancellationToken = default) {
if (_count > 0)
return; //completed
Interlocked.Increment(ref _notificationWaiters);
await _notifier.CreateTask(timeout, cancellationToken).ConfigureAwait(false);
Interlocked.Decrement(ref _notificationWaiters);
}
#endregion
public void Dispose() {
_notifier.TrySetCanceled(default);
_queues.Clear();
}
}
#endif