Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 46 additions & 0 deletions Containers.Test/SpscRingBufferTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

namespace ktsu.Containers.Tests;

using System.Diagnostics;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.VisualStudio.TestTools.UnitTesting;
Expand Down Expand Up @@ -149,6 +150,51 @@ public async Task ConcurrentProducerConsumer_TransfersAllItemsInOrder()
Assert.IsTrue(buffer.IsEmpty);
}

[TestMethod]
public async Task Count_ReadFromThirdThread_NeverExceedsTrueCount()
{
// The producer only enqueues into an empty buffer, so the true count is always 0 or 1.
SpscRingBuffer<int> buffer = new(1023);
TimeSpan duration = TimeSpan.FromSeconds(2);
int stop = 0;

Task producer = Task.Run(() =>
{
int next = 0;
while (Volatile.Read(ref stop) == 0)
{
if (buffer.IsEmpty)
{
buffer.TryEnqueue(next++);
}
}
}, TestContext.CancellationToken);

Task consumer = Task.Run(() =>
{
while (Volatile.Read(ref stop) == 0)
{
buffer.TryDequeue(out _);
}
}, TestContext.CancellationToken);

Task<int> observer = Task.Run(() =>
{
int maxObserved = 0;
Stopwatch stopwatch = Stopwatch.StartNew();
while (stopwatch.Elapsed < duration)
{
maxObserved = Math.Max(maxObserved, buffer.Count);
}

Volatile.Write(ref stop, 1);
return maxObserved;
}, TestContext.CancellationToken);

await Task.WhenAll(producer, consumer, observer).ConfigureAwait(false);
Assert.IsLessThanOrEqualTo(1, await observer.ConfigureAwait(false), "Count must never report more elements than the buffer held.");
}

[TestMethod]
public void Dequeue_ReferenceType_ReleasesReference()
{
Expand Down
21 changes: 17 additions & 4 deletions Containers/SpscRingBuffer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -107,15 +107,28 @@ public SpscRingBuffer(int capacity)
/// </summary>
/// <remarks>
/// This is a best-effort snapshot intended for diagnostics/metering. Under concurrent access
/// the true count may have changed by the time the value is returned.
/// the true count may have changed by the time the value is returned, but it is always a count
/// the buffer actually held while the property was being read, even when read from a thread
/// that is neither the producer nor the consumer.
/// </remarks>
public int Count
{
get
{
int currentTail = Volatile.Read(ref tail.Value);
int currentHead = Volatile.Read(ref head.Value);
return (currentTail - currentHead) & mask;
// Head and tail are read one after the other, so between the two reads the producer and
// consumer can both move on. If head passed the tail that was read, the masked difference
// would wrap to a value near Capacity that the buffer never held. Reading head on both
// sides of tail, and retrying until it did not move, gives a tail observed while head
// held that value, so the difference is a count the buffer really had.
while (true)
{
int currentHead = Volatile.Read(ref head.Value);
int currentTail = Volatile.Read(ref tail.Value);
if (Volatile.Read(ref head.Value) == currentHead)
{
return (currentTail - currentHead) & mask;
}
}
}
}

Expand Down
Loading