From 662242d1e57819df5052867c0acfbeed9d589def Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 28 Sep 2026 23:27:13 +0000 Subject: [PATCH] fix: stop SpscRingBuffer.Count reporting a count the buffer never held [patch] Count read tail and then head. A thread that was neither producer nor consumer could see the consumer move head past the stale tail between the two reads, and the masked difference then wrapped to a value near Capacity on an idle buffer. Count now reads head, then tail, then head again, and retries until head did not move, so the tail it uses was observed while head held that value. Fixes ktsu-dev/Containers#81 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01EdrcmpMTNHSnACRgnNGs8v --- Containers.Test/SpscRingBufferTests.cs | 46 ++++++++++++++++++++++++++ Containers/SpscRingBuffer.cs | 21 +++++++++--- 2 files changed, 63 insertions(+), 4 deletions(-) diff --git a/Containers.Test/SpscRingBufferTests.cs b/Containers.Test/SpscRingBufferTests.cs index 6595a12..d254b18 100644 --- a/Containers.Test/SpscRingBufferTests.cs +++ b/Containers.Test/SpscRingBufferTests.cs @@ -2,6 +2,7 @@ namespace ktsu.Containers.Tests; +using System.Diagnostics; using System.Threading; using System.Threading.Tasks; using Microsoft.VisualStudio.TestTools.UnitTesting; @@ -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 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 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() { diff --git a/Containers/SpscRingBuffer.cs b/Containers/SpscRingBuffer.cs index 40a0f1c..adc8e13 100644 --- a/Containers/SpscRingBuffer.cs +++ b/Containers/SpscRingBuffer.cs @@ -107,15 +107,28 @@ public SpscRingBuffer(int capacity) /// /// /// 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. /// 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; + } + } } }