From d076d09317d290f3eabf59864d1b96f58b44b4b4 Mon Sep 17 00:00:00 2001 From: prezaei Date: Sat, 26 Nov 2022 21:36:56 -0800 Subject: [PATCH 1/7] Enhancing the subscriber --- src/Interprocess/Queue/Subscriber.cs | 30 +++++++++++++++++----------- 1 file changed, 18 insertions(+), 12 deletions(-) diff --git a/src/Interprocess/Queue/Subscriber.cs b/src/Interprocess/Queue/Subscriber.cs index c018df2..e841a83 100644 --- a/src/Interprocess/Queue/Subscriber.cs +++ b/src/Interprocess/Queue/Subscriber.cs @@ -102,22 +102,28 @@ private unsafe bool TryDequeueImpl( cancellationSource.ThrowIfCancellationRequested(cancellation); message = ReadOnlyMemory.Empty; - var header = *Header; - // is this an empty queue? - if (header.IsEmpty()) - return false; + while (true) + { + var header = *Header; + + // is this an empty queue? + if (header.IsEmpty()) + return false; - var readLockTimestamp = header.ReadLockTimestamp; - var now = DateTime.UtcNow.Ticks; + var readLockTimestamp = header.ReadLockTimestamp; + var now = DateTime.UtcNow.Ticks; - // is there already a read-lock or has the previous lock timed out meaning that a subscriber crashed? - if (now - readLockTimestamp < TicksForTenSeconds) - return false; + // is there already a read-lock or has the previous lock timed out meaning that a subscriber crashed? + if (now - readLockTimestamp < TicksForTenSeconds) + return false; + + // take a read-lock so no other thread can read a message + if (Interlocked.CompareExchange(ref Header->ReadLockTimestamp, now, readLockTimestamp) == readLockTimestamp) + break; - // take a read-lock so no other thread can read a message - if (Interlocked.CompareExchange(ref Header->ReadLockTimestamp, now, readLockTimestamp) != readLockTimestamp) - return false; + Thread.Yield(); + } try { From 03dbbdd784d4e4837f3a376e8bf93818509e65b4 Mon Sep 17 00:00:00 2001 From: prezaei Date: Mon, 28 Nov 2022 12:41:48 -0800 Subject: [PATCH 2/7] faster subscriber --- src/Interprocess/Queue/Subscriber.cs | 24 ++++++++++-------------- 1 file changed, 10 insertions(+), 14 deletions(-) diff --git a/src/Interprocess/Queue/Subscriber.cs b/src/Interprocess/Queue/Subscriber.cs index e841a83..bd4c337 100644 --- a/src/Interprocess/Queue/Subscriber.cs +++ b/src/Interprocess/Queue/Subscriber.cs @@ -6,7 +6,7 @@ namespace Cloudtoid.Interprocess { internal sealed class Subscriber : Queue, ISubscriber { - private static readonly long TicksForTenSeconds = TimeSpan.FromSeconds(10).Ticks; + private static readonly long TenSeconds = TimeSpan.FromSeconds(10).Ticks; private readonly CancellationTokenSource cancellationSource = new(); private readonly CountdownEvent countdownEvent = new(1); private readonly IInterprocessSemaphoreWaiter signal; @@ -74,18 +74,13 @@ private ReadOnlyMemory DequeueCore(Memory? resultBuffer, Cancellatio try { - int i = -5; while (true) { if (TryDequeueImpl(resultBuffer, cancellation, out var message)) return message; - if (i > 10) - signal.Wait(millisecondsTimeout: 10); - else if (i++ > 0) - signal.Wait(millisecondsTimeout: i); - else - Thread.Yield(); + signal.Wait(millisecondsTimeout: 5); + cancellationSource.ThrowIfCancellationRequested(cancellation); } } finally @@ -99,8 +94,6 @@ private unsafe bool TryDequeueImpl( CancellationToken cancellation, out ReadOnlyMemory message) { - cancellationSource.ThrowIfCancellationRequested(cancellation); - message = ReadOnlyMemory.Empty; while (true) @@ -115,14 +108,16 @@ private unsafe bool TryDequeueImpl( var now = DateTime.UtcNow.Ticks; // is there already a read-lock or has the previous lock timed out meaning that a subscriber crashed? - if (now - readLockTimestamp < TicksForTenSeconds) - return false; + if (now - readLockTimestamp < TenSeconds) + { + Thread.Yield(); + cancellationSource.ThrowIfCancellationRequested(cancellation); + continue; + } // take a read-lock so no other thread can read a message if (Interlocked.CompareExchange(ref Header->ReadLockTimestamp, now, readLockTimestamp) == readLockTimestamp) break; - - Thread.Yield(); } try @@ -142,6 +137,7 @@ private unsafe bool TryDequeueImpl( MessageHeader.ReadyToBeConsumedState) != MessageHeader.ReadyToBeConsumedState) { Thread.Yield(); + cancellationSource.ThrowIfCancellationRequested(cancellation); } // read the message body from the queue From 17b659cede756fa339a69fe7e324488a711df136 Mon Sep 17 00:00:00 2001 From: prezaei Date: Fri, 9 Dec 2022 15:21:05 -0800 Subject: [PATCH 3/7] Updates to readme.md --- README.md | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/README.md b/README.md index 14e9be1..2950cdf 100644 --- a/README.md +++ b/README.md @@ -12,10 +12,10 @@ - [**Fast**](#performance): It is *extremely* fast. - **Cross-platform**: It supports Windows, and Unix-based operating systems such as Linux, [MacOS][MacOSWiki], and [FreeBSD][FreeBSDOrg]. -- [**API**](#Usage): Provides a simple and intuitive API to enqueue/send and dequeue/receive messages. +- [**API**](#usage): Provides a simple and intuitive API to enqueue/send and dequeue/receive messages. - **Multiple publishers and subscribers**: It supports multiple publishers and subscribers to a shared queue. - [**Efficient**](#performance): Sending and receiving messages is almost heap memory allocation free reducing garbage collections. -- [**Developer**](#Author): Developed by a guy at Microsoft. +- [**Developer**](#author): Developed by a guy at Microsoft. ## NuGet Package @@ -40,7 +40,7 @@ Creating a message queue publisher: ```csharp var options = new QueueOptions( queueName: "my-queue", - bytesCapacity: 1024 * 1024); + capacity: 1024 * 1024); using var publisher = factory.CreatePublisher(options); publisher.TryEnqueue(message); @@ -51,7 +51,7 @@ Creating a message queue subscriber: ```csharp options = new QueueOptions( queueName: "my-queue", - bytesCapacity: 1024 * 1024); + capacity: 1024 * 1024); using var subscriber = factory.CreateSubscriber(options); subscriber.TryDequeue(messageBuffer, cancellationToken, out var message); @@ -72,7 +72,7 @@ Creating a message queue publisher using an instance of `IQueueFactory` retrieve ```csharp var options = new QueueOptions( queueName: "my-queue", - bytesCapacity: 1024 * 1024); + capacity: 1024 * 1024); using var publisher = factory.CreatePublisher(options); publisher.TryEnqueue(message); @@ -83,7 +83,7 @@ Creating a message queue subscriber using an instance of `IQueueFactory` retriev ```csharp var options = new QueueOptions( queueName: "my-queue", - bytesCapacity: 1024 * 1024); + capacity: 1024 * 1024); using var subscriber = factory.CreateSubscriber(options); subscriber.TryDequeue(messageBuffer, cancellationToken, out var message); @@ -229,8 +229,8 @@ Here are a couple of items that we are working on. [MacOSWiki]:https://en.wikipedia.org/wiki/MacOS [FreeBSDOrg]:https://www.freebsd.org/ [Wow64Wiki]:https://en.wikipedia.org/wiki/WoW64 -[WslDoc]:https://docs.microsoft.com/en-us/windows/wsl/about +[WslDoc]:https://learn.microsoft.com/windows/wsl/about [BenchmarkOrg]:https://benchmarkdotnet.org/ -[NamedSemaphoresDoc]:https://docs.microsoft.com/en-us/dotnet/api/system.threading.semaphore#remarks -[SemaphoreDoc]:https://docs.microsoft.com/en-us/dotnet/api/system.threading.semaphore +[NamedSemaphoresDoc]:https://docs.microsoft.com/dotnet/api/system.threading.semaphore#remarks +[SemaphoreDoc]:https://docs.microsoft.com/dotnet/api/system.threading.semaphore [PedramLinkedIn]:https://www.linkedin.com/in/pedramrezaei/ From e6710ba8de7bb579315ce11fedf3bb323c223344 Mon Sep 17 00:00:00 2001 From: prezaei Date: Fri, 9 Dec 2022 15:25:39 -0800 Subject: [PATCH 4/7] One more readme! --- README.md | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/README.md b/README.md index 2950cdf..29fe2bf 100644 --- a/README.md +++ b/README.md @@ -11,7 +11,7 @@ **Cloudtoid Interprocess** is a cross-platform shared memory queue for fast communication between processes ([Interprocess Communication or IPC][IPCWiki]). It uses a shared memory-mapped file for extremely fast and efficient communication between processes and it is used internally by Microsoft. - [**Fast**](#performance): It is *extremely* fast. -- **Cross-platform**: It supports Windows, and Unix-based operating systems such as Linux, [MacOS][MacOSWiki], and [FreeBSD][FreeBSDOrg]. +- **Cross-platform**: It supports Windows, and Unix-based operating systems such as Linux, [macOS][macOSWiki], and [FreeBSD][FreeBSDOrg]. - [**API**](#usage): Provides a simple and intuitive API to enqueue/send and dequeue/receive messages. - **Multiple publishers and subscribers**: It supports multiple publishers and subscribers to a shared queue. - [**Efficient**](#performance): Sending and receiving messages is almost heap memory allocation free reducing garbage collections. @@ -102,7 +102,7 @@ Please note that you can start multiple publishers and subscribers sending and r A lot has gone into optimizing the implementation of this library. For instance, it is mostly heap-memory allocation free, reducing the need for garbage collection induced pauses. -**Summary**: A full enqueue followed by a dequeue takes `~250 ns` on Linux, `~650 ns` on MacOS, and `~300 ns` on Windows. +**Summary**: A full enqueue followed by a dequeue takes `~250 ns` on Linux, `~650 ns` on macOS, and `~300 ns` on Windows. **Details**: To benchmark the performance and memory usage, we use [BenchmarkDotNet][BenchmarkOrg] and perform the following runs: @@ -148,7 +148,7 @@ Results: --- -### On MacOS +### On macOS Host: @@ -194,14 +194,14 @@ Results: This library relies on [Named Semaphores][NamedSemaphoresDoc] To signal the existence of a new message to all message subscribers and to do it across process boundaries. Named semaphores are synchronization constructs accessible across processes. -.NET Core 3.1 and .NET 6/7 do not support named semaphores on Unix-based OSs (Linux, macOS, etc.). Instead we are using P/Invoke and relying on operating system's POSIX semaphore implementation. ([Linux](src/Interprocess/Semaphore/Linux/Interop.cs) and [MacOS](src/Interprocess/Semaphore/MacOS/Interop.cs) implementations). +.NET Core 3.1 and .NET 6/7 do not support named semaphores on Unix-based OSs (Linux, macOS, etc.). Instead we are using P/Invoke and relying on operating system's POSIX semaphore implementation. ([Linux](src/Interprocess/Semaphore/Linux/Interop.cs) and [macOS](src/Interprocess/Semaphore/MacOS/Interop.cs) implementations). This implementation will be replaced with [`System.Threading.Semaphore`][SemaphoreDoc] once .NET adds support for named semaphores on all platforms. ## How to Contribute - Create a branch from `main`. -- Ensure that all tests pass on Windows, Linux, and MacOS. +- Ensure that all tests pass on Windows, Linux, and macOS. - Keep the code coverage number above 80% by adding new tests or modifying the existing tests. - Send a pull request. @@ -226,7 +226,7 @@ Here are a couple of items that we are working on. [DotNetPlatformBadge]:https://img.shields.io/badge/.net-%3E%206.0-blue [NuGet]:https://www.nuget.org/packages/Cloudtoid.Interprocess/ [IPCWiki]:https://en.wikipedia.org/wiki/Inter-process_communication -[MacOSWiki]:https://en.wikipedia.org/wiki/MacOS +[macOSWiki]:https://en.wikipedia.org/wiki/macOS [FreeBSDOrg]:https://www.freebsd.org/ [Wow64Wiki]:https://en.wikipedia.org/wiki/WoW64 [WslDoc]:https://learn.microsoft.com/windows/wsl/about From 03ef4d8f49b2b05491929f4dcb00a2e496169b52 Mon Sep 17 00:00:00 2001 From: Pedram Rezaei Date: Sat, 12 Sep 2026 16:32:44 -0700 Subject: [PATCH 5/7] Bound subscriber retries before waiting and preserve cancellation --- README.md | 8 +- src/Interprocess.Benchmark/Program.cs | 2 +- .../Queue/EnqueueBenchmark.cs | 2 +- .../Queue/QueueBenchmark.cs | 2 +- .../Queue/QueueExtendedBenchmark.cs | 2 +- .../Queue/SubscriberBenchmark.cs | 65 +++++++ src/Interprocess.Tests/SubscriberTests.cs | 173 ++++++++++++++++++ src/Interprocess/Queue/Subscriber.cs | 18 +- 8 files changed, 261 insertions(+), 11 deletions(-) create mode 100644 src/Interprocess.Benchmark/Queue/SubscriberBenchmark.cs create mode 100644 src/Interprocess.Tests/SubscriberTests.cs diff --git a/README.md b/README.md index 1efefed..8b4f491 100644 --- a/README.md +++ b/README.md @@ -124,7 +124,13 @@ A lot has gone into optimizing the implementation of this library. For instance, You can replicate the results by running the following command: ```sh -dotnet run Interprocess.Benchmark.csproj -c Release +dotnet run --project src/Interprocess.Benchmark -c Release -- --filter '*QueueBenchmark*' +``` + +To compare throughput with one subscriber versus four concurrent subscribers: + +```sh +dotnet run --project src/Interprocess.Benchmark -c Release -- --filter '*SubscriberBenchmark*' --iterationCount 8 ``` --- diff --git a/src/Interprocess.Benchmark/Program.cs b/src/Interprocess.Benchmark/Program.cs index d466157..cab6510 100644 --- a/src/Interprocess.Benchmark/Program.cs +++ b/src/Interprocess.Benchmark/Program.cs @@ -4,5 +4,5 @@ namespace Cloudtoid.Interprocess.Benchmark; public sealed class Program { - public static void Main() => _ = BenchmarkRunner.Run(typeof(Program).Assembly); + public static void Main(string[] args) => BenchmarkSwitcher.FromAssembly(typeof(Program).Assembly).Run(args); } \ No newline at end of file diff --git a/src/Interprocess.Benchmark/Queue/EnqueueBenchmark.cs b/src/Interprocess.Benchmark/Queue/EnqueueBenchmark.cs index 2164d68..6fee985 100644 --- a/src/Interprocess.Benchmark/Queue/EnqueueBenchmark.cs +++ b/src/Interprocess.Benchmark/Queue/EnqueueBenchmark.cs @@ -3,7 +3,7 @@ namespace Cloudtoid.Interprocess.Benchmark; -[SimpleJob(RuntimeMoniker.Net90)] +[SimpleJob(RuntimeMoniker.Net10_0)] [MemoryDiagnoser] [MarkdownExporterAttribute.GitHub] public class EnqueueBenchmark diff --git a/src/Interprocess.Benchmark/Queue/QueueBenchmark.cs b/src/Interprocess.Benchmark/Queue/QueueBenchmark.cs index a40330f..dec9528 100644 --- a/src/Interprocess.Benchmark/Queue/QueueBenchmark.cs +++ b/src/Interprocess.Benchmark/Queue/QueueBenchmark.cs @@ -3,7 +3,7 @@ namespace Cloudtoid.Interprocess.Benchmark; -[SimpleJob(RuntimeMoniker.Net90)] +[SimpleJob(RuntimeMoniker.Net10_0)] [MemoryDiagnoser] [MarkdownExporterAttribute.GitHub] public class QueueBenchmark diff --git a/src/Interprocess.Benchmark/Queue/QueueExtendedBenchmark.cs b/src/Interprocess.Benchmark/Queue/QueueExtendedBenchmark.cs index bed7be0..89d63ee 100644 --- a/src/Interprocess.Benchmark/Queue/QueueExtendedBenchmark.cs +++ b/src/Interprocess.Benchmark/Queue/QueueExtendedBenchmark.cs @@ -3,7 +3,7 @@ namespace Cloudtoid.Interprocess.Benchmark; -[SimpleJob(RuntimeMoniker.Net90)] +[SimpleJob(RuntimeMoniker.Net10_0)] [MarkdownExporterAttribute.GitHub] public class QueueExtendedBenchmark { diff --git a/src/Interprocess.Benchmark/Queue/SubscriberBenchmark.cs b/src/Interprocess.Benchmark/Queue/SubscriberBenchmark.cs new file mode 100644 index 0000000..c123fc7 --- /dev/null +++ b/src/Interprocess.Benchmark/Queue/SubscriberBenchmark.cs @@ -0,0 +1,65 @@ +using BenchmarkDotNet.Attributes; + +namespace Cloudtoid.Interprocess.Benchmark; + +[ShortRunJob] +[MarkdownExporterAttribute.GitHub] +public class SubscriberBenchmark +{ + private const int MessageCount = 65536; + private readonly QueueFactory factory = new(); + private readonly QueueOptions options = new("subscriber-bench", 65536); + private IPublisher publisher = null!; + private ISubscriber[] subscribers = []; + + [Params(1, 4)] + public int SubscriberCount { get; set; } + + [GlobalSetup] + public void Setup() + { + publisher = factory.CreatePublisher(options); + subscribers = Enumerable.Range(0, SubscriberCount).Select(_ => factory.CreateSubscriber(options)).ToArray(); + } + + [GlobalCleanup] + public void Cleanup() + { + foreach (var subscriber in subscribers) + subscriber.Dispose(); + + publisher.Dispose(); + } + + [Benchmark(OperationsPerInvoke = MessageCount)] + public async Task ReceiveConcurrentlyAsync() + { + using var cancellation = new CancellationTokenSource(TimeSpan.FromSeconds(30)); + var readers = new Task[subscribers.Length]; + for (var reader = 0; reader < subscribers.Length; reader++) + { + var subscriber = subscribers[reader]; + readers[reader] = Task.Factory.StartNew( + () => + { + var buffer = new byte[8]; + for (var i = 0; i < MessageCount / SubscriberCount; i++) + subscriber.Dequeue(buffer, cancellation.Token); + }, + cancellation.Token, + TaskCreationOptions.LongRunning, + TaskScheduler.Default); + } + + for (var i = 0; i < MessageCount; i++) + { + while (!publisher.TryEnqueue("message!"u8)) + { + cancellation.Token.ThrowIfCancellationRequested(); + Thread.Yield(); + } + } + + await Task.WhenAll(readers); + } +} \ No newline at end of file diff --git a/src/Interprocess.Tests/SubscriberTests.cs b/src/Interprocess.Tests/SubscriberTests.cs new file mode 100644 index 0000000..f35a46a --- /dev/null +++ b/src/Interprocess.Tests/SubscriberTests.cs @@ -0,0 +1,173 @@ +namespace Cloudtoid.Interprocess.Tests; + +public sealed class SubscriberTests(UniquePathFixture fixture) : IClassFixture +{ + private readonly QueueOptions options = new(Guid.NewGuid().ToStringInvariant("N")[..16], fixture.Path, 256); + private readonly QueueFactory factory = new(); + + [Fact] + public async Task TryDequeueDoesNotWaitForAnotherSubscriberAsync() + { + using var probe = new QueueProbe(options); + using var publisher = factory.CreatePublisher(options); + using var subscriber = factory.CreateSubscriber(options); + publisher.TryEnqueue("message!"u8).Should().BeTrue(); + probe.LockReads(); + try + { + var attempt = Task.Run(() => subscriber.TryDequeue(default, out _)); + (await attempt.WaitAsync(TimeSpan.FromSeconds(1))).Should().BeFalse(); + } + finally + { + probe.UnlockReads(); + } + + subscriber.TryDequeue(default, out _).Should().BeTrue(); + } + + [Fact] + public void ExpiredSubscriberLockCanBeRecovered() + { + using var probe = new QueueProbe(options); + using var publisher = factory.CreatePublisher(options); + using var subscriber = factory.CreateSubscriber(options); + publisher.TryEnqueue("message!"u8).Should().BeTrue(); + probe.AbandonReadLock(); + subscriber.TryDequeue(default, out var message).Should().BeTrue(); + message.ToArray().Should().Equal("message!"u8.ToArray()); + probe.ReadsAreLocked.Should().BeFalse(); + } + + [Fact] + public async Task BlockingDequeueRetriesContendedReadAsync() + { + using var probe = new QueueProbe(options); + using var publisher = factory.CreatePublisher(options); + using var subscriber = factory.CreateSubscriber(options); + using var cancellation = new CancellationTokenSource(TimeSpan.FromSeconds(5)); + publisher.TryEnqueue("message!"u8).Should().BeTrue(); + probe.LockReads(); + var read = Task.Run(() => subscriber.Dequeue(cancellation.Token)); + try + { + await Task.Delay(20); + read.IsCompleted.Should().BeFalse(); + } + finally + { + probe.UnlockReads(); + } + + (await read.WaitAsync(TimeSpan.FromSeconds(1))).ToArray().Should().Equal("message!"u8.ToArray()); + } + + [Theory] + [InlineData(true)] + [InlineData(false)] + public async Task CancellationInterruptsUnfinishedMessageAsync(bool blocking) + { + using var probe = new QueueProbe(options); + using var subscriber = factory.CreateSubscriber(options); + using var cancellation = new CancellationTokenSource(); + probe.ReserveUnfinishedMessage(); + var read = Task.Run(() => + { + if (blocking) + subscriber.Dequeue(cancellation.Token); + else + subscriber.TryDequeue(cancellation.Token, out _); + }); + + try + { + SpinWait.SpinUntil(() => probe.ReadsAreLocked, TimeSpan.FromSeconds(1)).Should().BeTrue(); + await cancellation.CancelAsync(); + await Assert.ThrowsAnyAsync( + async () => await read.WaitAsync(TimeSpan.FromSeconds(1))); + probe.ReadsAreLocked.Should().BeFalse(); + } + finally + { + await cancellation.CancelAsync(); + } + } + + [Fact] + public async Task EmptyBlockingDequeueCanBeCancelledAsync() + { + using var subscriber = factory.CreateSubscriber(options); + using var cancellation = new CancellationTokenSource(); + var read = Task.Run(() => subscriber.Dequeue(cancellation.Token)); + await Task.Delay(20); + await cancellation.CancelAsync(); + await Assert.ThrowsAnyAsync( + async () => await read.WaitAsync(TimeSpan.FromSeconds(1))); + } + + [Theory] + [InlineData(1)] + [InlineData(4)] + public async Task ConcurrentSubscribersReceiveEveryMessageExactlyOnceAsync(int subscriberCount) + { + const int count = 4000; + var received = new int[count]; + using var cancellation = new CancellationTokenSource(TimeSpan.FromSeconds(10)); + // Keep a participant alive while the workers start and finish. + using var anchor = factory.CreatePublisher(options); + var readers = new Task[subscriberCount]; + for (var reader = 0; reader < subscriberCount; reader++) + { + readers[reader] = Task.Run(() => + { + using var subscriber = factory.CreateSubscriber(options); + var buffer = new byte[8]; + for (var i = 0; i < count / subscriberCount; i++) + { + var message = subscriber.Dequeue(buffer, cancellation.Token); + message.Length.Should().Be(8); + var id = BitConverter.ToInt32(message.Span); + id.Should().BeInRange(0, count - 1); + Interlocked.Increment(ref received[id]); + } + }); + } + + var writers = new Task[2]; + for (var writer = 0; writer < writers.Length; writer++) + { + var firstId = writer; + writers[writer] = Task.Run(() => + { + using var publisher = factory.CreatePublisher(options); + var buffer = new byte[8]; + for (var id = firstId; id < count; id += 2) + { + BitConverter.TryWriteBytes(buffer, id).Should().BeTrue(); + while (!publisher.TryEnqueue(buffer)) + { + cancellation.Token.ThrowIfCancellationRequested(); + Thread.Yield(); + } + } + }); + } + + await Task.WhenAll(readers.Concat(writers)).WaitAsync(TimeSpan.FromSeconds(15)); + received.Should().OnlyContain(value => value == 1); + } + + private sealed class QueueProbe(QueueOptions options) : Queue(options, NullLoggerFactory.Instance) + { + internal unsafe bool ReadsAreLocked => Interlocked.Read(ref Header->ReadLockTimestamp) != 0; + + internal unsafe void LockReads() => Interlocked.Exchange(ref Header->ReadLockTimestamp, DateTime.UtcNow.Ticks); + + internal unsafe void AbandonReadLock() => + Interlocked.Exchange(ref Header->ReadLockTimestamp, DateTime.UtcNow.Ticks - TimeSpan.FromSeconds(11).Ticks); + + internal unsafe void UnlockReads() => Interlocked.Exchange(ref Header->ReadLockTimestamp, 0); + + internal unsafe void ReserveUnfinishedMessage() => Interlocked.Exchange(ref Header->WriteOffset, 16); + } +} \ No newline at end of file diff --git a/src/Interprocess/Queue/Subscriber.cs b/src/Interprocess/Queue/Subscriber.cs index 13f7c3f..dc2aa71 100644 --- a/src/Interprocess/Queue/Subscriber.cs +++ b/src/Interprocess/Queue/Subscriber.cs @@ -83,18 +83,23 @@ private ReadOnlyMemory DequeueCore(Memory? resultBuffer, Cancellatio try { - int i = -5; + SpinWait spin = default; while (true) { if (TryDequeueImpl(resultBuffer, cancellation, out var message)) return message; - if (i > 10) - signal.Wait(millisecondsTimeout: 10); - else if (i++ > 0) - signal.Wait(millisecondsTimeout: i); + // Retry briefly in user space while another reader finishes. Once spinning + // would yield, wait for a signal instead of burning CPU on an idle queue. + if (spin.NextSpinWillYield) + { + signal.Wait(millisecondsTimeout: 5); + spin.Reset(); + } else - Thread.Yield(); + { + spin.SpinOnce(); + } } } finally @@ -159,6 +164,7 @@ private unsafe bool TryDequeueImpl( Interlocked.Exchange(ref Header->ReadOffset, writeOffset); return false; } + cancellationSource.ThrowIfCancellationRequested(cancellation); Thread.Yield(); } From c7d07559a48045c82b63940da9f7910b888d9e82 Mon Sep 17 00:00:00 2001 From: Pedram Rezaei Date: Sat, 12 Sep 2026 16:40:07 -0700 Subject: [PATCH 6/7] Correct benchmark workloads and measure all supported CI platforms --- .github/workflows/benchmarks.yml | 36 +++++++++++++++++++ .../Queue/EnqueueBenchmark.cs | 21 +++++++---- .../Queue/QueueExtendedBenchmark.cs | 24 ++++++++----- 3 files changed, 66 insertions(+), 15 deletions(-) create mode 100644 .github/workflows/benchmarks.yml diff --git a/.github/workflows/benchmarks.yml b/.github/workflows/benchmarks.yml new file mode 100644 index 0000000..a16426c --- /dev/null +++ b/.github/workflows/benchmarks.yml @@ -0,0 +1,36 @@ +name: benchmarks + +on: + workflow_dispatch: + pull_request: + paths: + - 'src/Interprocess.Benchmark/**' + - '.github/workflows/benchmarks.yml' + +permissions: + contents: read + +jobs: + benchmarks: + name: Benchmarks on ${{ matrix.os }} + runs-on: ${{ matrix.os }} + timeout-minutes: 20 + strategy: + fail-fast: false + matrix: + os: [ubuntu-latest, windows-latest, macos-latest] + steps: + - uses: actions/checkout@v7 + - uses: actions/setup-dotnet@v6 + with: + dotnet-version: 10.0.x + - name: Run all benchmarks + working-directory: src + run: dotnet run --project Interprocess.Benchmark -c Release -- --filter '*' --warmupCount 3 --iterationCount 8 --artifacts BenchmarkDotNet.Artifacts + - name: Upload benchmark reports + if: always() + uses: actions/upload-artifact@v7 + with: + name: benchmarks-${{ matrix.os }} + path: src/BenchmarkDotNet.Artifacts/ + if-no-files-found: error diff --git a/src/Interprocess.Benchmark/Queue/EnqueueBenchmark.cs b/src/Interprocess.Benchmark/Queue/EnqueueBenchmark.cs index 6fee985..4f8ab75 100644 --- a/src/Interprocess.Benchmark/Queue/EnqueueBenchmark.cs +++ b/src/Interprocess.Benchmark/Queue/EnqueueBenchmark.cs @@ -8,6 +8,7 @@ namespace Cloudtoid.Interprocess.Benchmark; [MarkdownExporterAttribute.GitHub] public class EnqueueBenchmark { + private const int MessageCount = 320000; private static readonly byte[] Message = [100, 110, 120]; private static readonly Memory MessageBuffer = new byte[Message.Length]; #pragma warning disable CS8618 @@ -19,8 +20,8 @@ public class EnqueueBenchmark public void Setup() { var queueFactory = new QueueFactory(); - publisher = queueFactory.CreatePublisher(new QueueOptions("qn", Path.GetTempPath(), 5120000)); - subscriber = queueFactory.CreateSubscriber(new QueueOptions("qn", Path.GetTempPath(), 5120000)); + publisher = queueFactory.CreatePublisher(new QueueOptions("qn", Path.GetTempPath(), MessageCount * 16)); + subscriber = queueFactory.CreateSubscriber(new QueueOptions("qn", Path.GetTempPath(), MessageCount * 16)); } [GlobalCleanup] @@ -33,15 +34,21 @@ public void Cleanup() [IterationCleanup] public void DrainQueue() { - for (int i = 8; i < 320000; i++) - subscriber.Dequeue(MessageBuffer, default); + for (var i = 0; i < MessageCount; i++) + { + if (!subscriber.TryDequeue(MessageBuffer, default, out _)) + throw new InvalidOperationException("The benchmark did not enqueue the expected number of messages."); + } } // Expecting that there are NO managed heap allocations. - [Benchmark(Description = "Message enqueue (320,000 times)")] + [Benchmark(Description = "Message enqueue", OperationsPerInvoke = MessageCount)] public void Enqueue() { - for (int i = 8; i < 320000; i++) - publisher.TryEnqueue(Message); + for (var i = 0; i < MessageCount; i++) + { + if (!publisher.TryEnqueue(Message)) + throw new InvalidOperationException("The benchmark queue is full."); + } } } \ No newline at end of file diff --git a/src/Interprocess.Benchmark/Queue/QueueExtendedBenchmark.cs b/src/Interprocess.Benchmark/Queue/QueueExtendedBenchmark.cs index 89d63ee..d38199c 100644 --- a/src/Interprocess.Benchmark/Queue/QueueExtendedBenchmark.cs +++ b/src/Interprocess.Benchmark/Queue/QueueExtendedBenchmark.cs @@ -4,6 +4,7 @@ namespace Cloudtoid.Interprocess.Benchmark; [SimpleJob(RuntimeMoniker.Net10_0)] +[MemoryDiagnoser] [MarkdownExporterAttribute.GitHub] public class QueueExtendedBenchmark { @@ -14,13 +15,11 @@ public class QueueExtendedBenchmark private ISubscriber subscriber; #pragma warning restore CS8618 - [GlobalSetup] - public void Setup() - { - var queueFactory = new QueueFactory(); - publisher = queueFactory.CreatePublisher(new QueueOptions("qn", Path.GetTempPath(), 128)); - subscriber = queueFactory.CreateSubscriber(new QueueOptions("qn", Path.GetTempPath(), 128)); - } + [GlobalSetup(Target = nameof(EnqueueDequeue_LongMessage))] + public void Setup() => SetupQueue(128); + + [GlobalSetup(Target = nameof(EnqueueDequeue_WrappedMessages))] + public void SetupWrapped() => SetupQueue(120); [GlobalCleanup] public void Cleanup() @@ -38,7 +37,9 @@ public ReadOnlyMemory EnqueueDequeue_LongMessage() return subscriber.Dequeue(MessageBuffer, default); } - [Benchmark(Description = "Message enqueue and dequeue - wrapped message in circular buffer")] + // A padded message occupies 64 bytes. A 120-byte ring makes message bodies cross + // the end of the buffer; a 128-byte ring only cycles between aligned slots. + [Benchmark(Description = "Message enqueue and dequeue - ring-wrap workload", OperationsPerInvoke = 2)] public ReadOnlyMemory EnqueueDequeue_WrappedMessages() { if (!publisher.TryEnqueue(Message)) @@ -51,4 +52,11 @@ public ReadOnlyMemory EnqueueDequeue_WrappedMessages() return subscriber.Dequeue(MessageBuffer, default); } + + private void SetupQueue(long capacity) + { + var queueFactory = new QueueFactory(); + publisher = queueFactory.CreatePublisher(new QueueOptions("qn", Path.GetTempPath(), capacity)); + subscriber = queueFactory.CreateSubscriber(new QueueOptions("qn", Path.GetTempPath(), capacity)); + } } \ No newline at end of file From 36345d0cb0091e4a684d46778c8676cf7057d355 Mon Sep 17 00:00:00 2001 From: Pedram Rezaei Date: Sat, 12 Sep 2026 16:50:20 -0700 Subject: [PATCH 7/7] Document native macOS benchmark results --- .github/workflows/benchmarks.yml | 36 ---------- README.md | 39 +++++++---- docs/benchmarks/2026-09-12/README.md | 31 +++++++++ docs/benchmarks/2026-09-12/macos-local.md | 80 +++++++++++++++++++++++ 4 files changed, 139 insertions(+), 47 deletions(-) delete mode 100644 .github/workflows/benchmarks.yml create mode 100644 docs/benchmarks/2026-09-12/README.md create mode 100644 docs/benchmarks/2026-09-12/macos-local.md diff --git a/.github/workflows/benchmarks.yml b/.github/workflows/benchmarks.yml deleted file mode 100644 index a16426c..0000000 --- a/.github/workflows/benchmarks.yml +++ /dev/null @@ -1,36 +0,0 @@ -name: benchmarks - -on: - workflow_dispatch: - pull_request: - paths: - - 'src/Interprocess.Benchmark/**' - - '.github/workflows/benchmarks.yml' - -permissions: - contents: read - -jobs: - benchmarks: - name: Benchmarks on ${{ matrix.os }} - runs-on: ${{ matrix.os }} - timeout-minutes: 20 - strategy: - fail-fast: false - matrix: - os: [ubuntu-latest, windows-latest, macos-latest] - steps: - - uses: actions/checkout@v7 - - uses: actions/setup-dotnet@v6 - with: - dotnet-version: 10.0.x - - name: Run all benchmarks - working-directory: src - run: dotnet run --project Interprocess.Benchmark -c Release -- --filter '*' --warmupCount 3 --iterationCount 8 --artifacts BenchmarkDotNet.Artifacts - - name: Upload benchmark reports - if: always() - uses: actions/upload-artifact@v7 - with: - name: benchmarks-${{ matrix.os }} - path: src/BenchmarkDotNet.Artifacts/ - if-no-files-found: error diff --git a/README.md b/README.md index 8b4f491..048074d 100644 --- a/README.md +++ b/README.md @@ -112,7 +112,7 @@ Please note that you can start multiple publishers and subscribers sending and r A lot has gone into optimizing the implementation of this library. For instance, it is mostly heap-memory allocation free, reducing the need for garbage collection induced pauses. -**Summary**: A full enqueue followed by a dequeue takes `~250 ns` on Linux, `~650 ns` on macOS, and `~300 ns` on Windows. +**Latest native macOS measurement**: a three-byte enqueue/dequeue round trip with a reused buffer averaged **210.0 ns** on an Apple M5 Max. Only the macOS results below were refreshed on September 12, 2026; the Windows and Linux sections retain their historical measurements. **Details**: To benchmark the performance and memory usage, we use [BenchmarkDotNet][BenchmarkOrg] and perform the following runs: @@ -158,20 +158,37 @@ Results: ### On macOS -Host: +Measured **September 12, 2026**, running directly on the Mac: ```text -BenchmarkDotNet v0.14.0, macOS Sequoia 15.2 (24C101) [Darwin 24.2.0] -Apple M3 Max, 1 CPU, 16 logical and 16 physical cores -.NET SDK 9.0.101 - [Host] : .NET 9.0.0 (9.0.24.52809), Arm64 RyuJIT AdvSIMD - .NET 9.0 : .NET 9.0.0 (9.0.24.52809), Arm64 RyuJIT AdvSIMD +BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0] +Apple M5 Max, 1 CPU, 18 logical and 18 physical cores +.NET SDK 10.0.401 +.NET runtime 10.0.12, Arm64 RyuJIT +Release build; 3 warm-up iterations; 8 measured iterations; 1 launch ``` -| Method | Mean (ns) | Error (ns) | StdDev | Gen0 | Allocated | -|-------------------------------------------------- |----------:|-----------:|-------:|---------:|----------:| -| 'Message enqueue and dequeue' | `249.2` | `0.74` | `0.62` | `-` | `-` | -| 'Message enqueue and dequeue - no message buffer' | `252.1` | `4.10` | `3.83` | `0.0038` | `32 B` | +All seven cases completed. Times are means in nanoseconds, normalized per operation. For enqueue/dequeue rows, an operation is one complete round trip. Concurrent-delivery rows report amortized time per delivered message. + +| Workload | Mean (ns) | StdDev (ns) | Allocated per operation | +| --- | ---: | ---: | ---: | +| Enqueue, 3 bytes | 182.3 | 5.01 | 0 B | +| Enqueue + dequeue, 3 bytes, reused buffer | 210.0 | 0.46 | 0 B | +| Enqueue + dequeue, 3 bytes, new result array | 214.9 | 0.81 | 32 B | +| Enqueue + dequeue, 50 bytes, reused buffer | 214.6 | 1.33 | 0 B | +| Enqueue + dequeue, 50 bytes, ring-wrap workload | 223.8 | 1.08 | 0 B | +| Concurrent delivery, 8 bytes, 1 subscriber | 246.5 | 2.20 | Not measured | +| Concurrent delivery, 8 bytes, 4 subscribers | 344.7 | 1.81 | Not measured | + +The enqueue case batches 320,000 messages and drains the queue outside the timed body. The ring-wrap case uses a 120-byte queue so padded 64-byte records repeatedly cross the end of the buffer; two round trips per invocation are normalized to one. Concurrent delivery uses one publisher and dedicated subscriber threads to transfer batches of 65,536 messages, including worker startup and completion in the timing. + +These are in-process microbenchmarks, not end-to-end latency between separate applications. The concurrent cases measure throughput under contention, not individual message latency; their allocations were not measured. See the [complete native Mac reports and methodology](docs/benchmarks/2026-09-12/README.md) for source revision, errors, and reproduction details. + +Run all cases from the repository root: + +```sh +dotnet run --project src/Interprocess.Benchmark -c Release -- --filter '*' --warmupCount 3 --iterationCount 8 --artifacts BenchmarkDotNet.Artifacts +``` --- diff --git a/docs/benchmarks/2026-09-12/README.md b/docs/benchmarks/2026-09-12/README.md new file mode 100644 index 0000000..eb9b21a --- /dev/null +++ b/docs/benchmarks/2026-09-12/README.md @@ -0,0 +1,31 @@ +# Native macOS benchmark measurements — September 12, 2026 + +All seven benchmark cases were run directly on an Apple M5 Max using benchmark and library source at [`c7d0755`](https://github.com/cloudtoid/interprocess/commit/c7d07559a48045c82b63940da9f7910b888d9e82). The later README update does not change the measured code. + +[Complete BenchmarkDotNet reports](macos-local.md) + +## Environment + +- macOS Tahoe 26.6.2 (25G83), Darwin 25.6.0. +- Apple M5 Max, arm64, 18 logical and physical cores reported. +- BenchmarkDotNet 0.15.8; .NET SDK 10.0.401; .NET runtime 10.0.12; Release configuration. +- Three warm-up iterations and eight measured iterations per case, one launch, with BenchmarkDotNet's usual pilot, overhead, and outlier handling. +- These are in-process microbenchmarks. They do not measure end-to-end latency between separate applications, idle CPU usage, or latency percentiles. + +## Workloads + +- **Enqueue:** 320,000 three-byte messages per invocation, normalized to one enqueue. Queue capacity is 5,120,000 bytes. Draining and validation happen outside the timed body. Each enqueue checks that it succeeded. +- **Three-byte round trips:** enqueue followed by dequeue on the calling thread, with either a reused buffer or a newly allocated result array; queue capacity 128 bytes. +- **50-byte round trips:** the same calling-thread pattern with a reused buffer and a 128-byte queue. +- **Ring-wrap workload:** 50-byte messages in a 120-byte queue, so padded 64-byte records repeatedly cross the end of the ring. Each invocation performs two enqueue/dequeue pairs; reported time is normalized to one pair. +- **Concurrent delivery:** one publisher sends 65,536 eight-byte messages to one or four dedicated subscriber threads sharing a 65,536-byte queue. Each subscriber receives an equal share. Time is normalized per delivered message and includes per-batch worker startup and completion. + +Memory diagnostics reported 0 B/op for enqueue and reused-buffer round trips, and 32 B/op for the three-byte result-array case. Allocation diagnostics were not enabled for concurrent delivery, which creates workers and tasks per batch. + +## Reproduce + +From the repository root: + +```sh +dotnet run --project src/Interprocess.Benchmark -c Release -- --filter '*' --warmupCount 3 --iterationCount 8 --artifacts BenchmarkDotNet.Artifacts +``` diff --git a/docs/benchmarks/2026-09-12/macos-local.md b/docs/benchmarks/2026-09-12/macos-local.md new file mode 100644 index 0000000..a083a02 --- /dev/null +++ b/docs/benchmarks/2026-09-12/macos-local.md @@ -0,0 +1,80 @@ +# Local macOS — Apple M5 Max + +Measured September 12, 2026, using source commit [`c7d0755`](https://github.com/cloudtoid/interprocess/commit/c7d07559a48045c82b63940da9f7910b888d9e82). + +See [methodology and reproduction commands](README.md). The tables below are BenchmarkDotNet exports. + +## EnqueueBenchmark + +``` + +BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0] +Apple M5 Max, 1 CPU, 18 logical and 18 physical cores +.NET SDK 10.0.401 + [Host] : .NET 10.0.12 (10.0.12, 10.0.1226.42308), Arm64 RyuJIT armv8.0-a + .NET 10.0 : .NET 10.0.12 (10.0.12, 10.0.1226.42308), Arm64 RyuJIT armv8.0-a + +Job=.NET 10.0 Runtime=.NET 10.0 InvocationCount=1 +IterationCount=8 UnrollFactor=1 WarmupCount=3 + +``` +| Method | Mean | Error | StdDev | Allocated | +|------------------ |---------:|--------:|--------:|----------:| +| 'Message enqueue' | 182.3 ns | 9.57 ns | 5.01 ns | - | + +## QueueBenchmark + +``` + +BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0] +Apple M5 Max, 1 CPU, 18 logical and 18 physical cores +.NET SDK 10.0.401 + [Host] : .NET 10.0.12 (10.0.12, 10.0.1226.42308), Arm64 RyuJIT armv8.0-a + .NET 10.0 : .NET 10.0.12 (10.0.12, 10.0.1226.42308), Arm64 RyuJIT armv8.0-a + +Job=.NET 10.0 Runtime=.NET 10.0 IterationCount=8 +WarmupCount=3 + +``` +| Method | Mean | Error | StdDev | Gen0 | Allocated | +|-------------------------------------------------- |---------:|--------:|--------:|-------:|----------:| +| 'Message enqueue and dequeue - no message buffer' | 214.9 ns | 1.54 ns | 0.81 ns | 0.0038 | 32 B | +| 'Message enqueue and dequeue' | 210.0 ns | 1.03 ns | 0.46 ns | - | - | + +## QueueExtendedBenchmark + +``` + +BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0] +Apple M5 Max, 1 CPU, 18 logical and 18 physical cores +.NET SDK 10.0.401 + [Host] : .NET 10.0.12 (10.0.12, 10.0.1226.42308), Arm64 RyuJIT armv8.0-a + .NET 10.0 : .NET 10.0.12 (10.0.12, 10.0.1226.42308), Arm64 RyuJIT armv8.0-a + +Job=.NET 10.0 Runtime=.NET 10.0 IterationCount=8 +WarmupCount=3 + +``` +| Method | Mean | Error | StdDev | Allocated | +|--------------------------------------------------- |---------:|--------:|--------:|----------:| +| 'Message enqueue and dequeue - long message' | 214.6 ns | 2.55 ns | 1.33 ns | - | +| 'Message enqueue and dequeue - ring-wrap workload' | 223.8 ns | 2.43 ns | 1.08 ns | - | + +## SubscriberBenchmark + +``` + +BenchmarkDotNet v0.15.8, macOS Tahoe 26.6.2 (25G83) [Darwin 25.6.0] +Apple M5 Max, 1 CPU, 18 logical and 18 physical cores +.NET SDK 10.0.401 + [Host] : .NET 10.0.12 (10.0.12, 10.0.1226.42308), Arm64 RyuJIT armv8.0-a + ShortRun : .NET 10.0.12 (10.0.12, 10.0.1226.42308), Arm64 RyuJIT armv8.0-a + +Job=ShortRun IterationCount=8 LaunchCount=1 +WarmupCount=3 + +``` +| Method | SubscriberCount | Mean | Error | StdDev | +|------------------------- |---------------- |---------:|--------:|--------:| +| **ReceiveConcurrentlyAsync** | **1** | **246.5 ns** | **4.21 ns** | **2.20 ns** | +| **ReceiveConcurrentlyAsync** | **4** | **344.7 ns** | **3.45 ns** | **1.81 ns** |