diff --git a/.github/workflows/cd.yml b/.github/workflows/cd.yml index e3d0d1d..8a03acd 100644 --- a/.github/workflows/cd.yml +++ b/.github/workflows/cd.yml @@ -10,7 +10,7 @@ concurrency: env: PROJECT_PATH: 'src/ExtendedThreading/ExtendedThreading.csproj' - SOLUTION_PATH: 'ExtendedThreading.sln' + SOLUTION_PATH: 'ExtendedThreading.slnx' PACKAGE_OUTPUT_DIRECTORY: ${{ github.workspace }}/output NUGET_SOURCE_URL: 'https://api.nuget.org/v3/index.json' @@ -21,12 +21,12 @@ jobs: steps: - name: 'Checkout' - uses: actions/checkout@v4 + uses: actions/checkout@v6 - name: 'Install dotnet' - uses: actions/setup-dotnet@v4 + uses: actions/setup-dotnet@v5 with: - dotnet-version: 8.0.x + dotnet-version: 10.0.x - name: 'Restore packages' run: dotnet restore ${{ env.SOLUTION_PATH }} @@ -42,12 +42,12 @@ jobs: steps: - name: 'Checkout' - uses: actions/checkout@v4 + uses: actions/checkout@v6 - name: 'Install dotnet' - uses: actions/setup-dotnet@v4 + uses: actions/setup-dotnet@v5 with: - dotnet-version: 8.0.x + dotnet-version: 10.0.x - name: 'Restore packages' run: dotnet restore ${{ env.SOLUTION_PATH }} diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index bca44f8..9d82f85 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -12,7 +12,7 @@ concurrency: env: PROJECT_PATH: 'src/ExtendedThreading/ExtendedThreading.csproj' - SOLUTION_PATH: 'ExtendedThreading.sln' + SOLUTION_PATH: 'ExtendedThreading.slnx' PACKAGE_OUTPUT_DIRECTORY: ${{ github.workspace }}\output NUGET_SOURCE_URL: 'https://api.nuget.org/v3/index.json' @@ -23,12 +23,12 @@ jobs: steps: - name: 'Checkout' - uses: actions/checkout@v4 + uses: actions/checkout@v6 - name: 'Install dotnet' - uses: actions/setup-dotnet@v4 + uses: actions/setup-dotnet@v5 with: - dotnet-version: 8.0.x + dotnet-version: 10.0.x - name: 'Restore packages' run: dotnet restore ${{ env.SOLUTION_PATH }} diff --git a/.github/workflows/docfx.yml b/.github/workflows/docfx.yml index eafc5d2..09734b0 100644 --- a/.github/workflows/docfx.yml +++ b/.github/workflows/docfx.yml @@ -27,7 +27,7 @@ jobs: - name: Install .Net uses: actions/setup-dotnet@v5 with: - dotnet-version: 8.0.x + dotnet-version: 10.0.x - name: Install DocFX run: | diff --git a/CHANGELOG.md b/CHANGELOG.md new file mode 100644 index 0000000..5ae8f24 --- /dev/null +++ b/CHANGELOG.md @@ -0,0 +1,17 @@ +# Changelog + +All notable changes to this project will be documented in this file. + +The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). + +## [Unreleased] + +## [1.3.0] - 2026-09-09 + +### Added + +- `AsyncSignal` which mirrors `ThreadSignal` but for async/await with cancellation support rather that threads. + +### Changed + +- Updated to .Net 10 \ No newline at end of file diff --git a/Directory.Build.props b/Directory.Build.props index 061842a..c0822f0 100644 --- a/Directory.Build.props +++ b/Directory.Build.props @@ -1,7 +1,16 @@ - 1.2.1 + 1.3.0 + + net10.0 + enable + true + true + Steffen Skov + Steffen Skov + Steffen Skov 2022 + MIT diff --git a/ExtendedThreading.sln b/ExtendedThreading.sln deleted file mode 100644 index 0c6b914..0000000 --- a/ExtendedThreading.sln +++ /dev/null @@ -1,45 +0,0 @@ - -Microsoft Visual Studio Solution File, Format Version 12.00 -# Visual Studio Version 17 -VisualStudioVersion = 17.5.002.0 -MinimumVisualStudioVersion = 10.0.40219.1 -Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "ExtendedThreading", "src\ExtendedThreading\ExtendedThreading.csproj", "{84394417-21C9-4C81-B049-C7E75ACB60A0}" -EndProject -Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "ExtendedThreading.UnitTests", "tests\ExtendedThreading.UnitTests\ExtendedThreading.UnitTests.csproj", "{63422754-EFA8-4EC7-91E7-024D63A7BD29}" -EndProject -Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Solution Items", "Solution Items", "{E1EC6A95-DC77-4F50-8E24-C01ADA8FB446}" - ProjectSection(SolutionItems) = preProject - .editorconfig = .editorconfig - Directory.Build.props = Directory.Build.props - EndProjectSection -EndProject -Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "src", "src", "{615B6C70-A67D-4901-B3F8-49205B2B61D9}" -EndProject -Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "tests", "tests", "{C28CF6C2-329B-4FF2-9A1C-FF73778BE7AB}" -EndProject -Global - GlobalSection(SolutionConfigurationPlatforms) = preSolution - Debug|Any CPU = Debug|Any CPU - Release|Any CPU = Release|Any CPU - EndGlobalSection - GlobalSection(ProjectConfigurationPlatforms) = postSolution - {84394417-21C9-4C81-B049-C7E75ACB60A0}.Debug|Any CPU.ActiveCfg = Debug|Any CPU - {84394417-21C9-4C81-B049-C7E75ACB60A0}.Debug|Any CPU.Build.0 = Debug|Any CPU - {84394417-21C9-4C81-B049-C7E75ACB60A0}.Release|Any CPU.ActiveCfg = Release|Any CPU - {84394417-21C9-4C81-B049-C7E75ACB60A0}.Release|Any CPU.Build.0 = Release|Any CPU - {63422754-EFA8-4EC7-91E7-024D63A7BD29}.Debug|Any CPU.ActiveCfg = Debug|Any CPU - {63422754-EFA8-4EC7-91E7-024D63A7BD29}.Debug|Any CPU.Build.0 = Debug|Any CPU - {63422754-EFA8-4EC7-91E7-024D63A7BD29}.Release|Any CPU.ActiveCfg = Release|Any CPU - {63422754-EFA8-4EC7-91E7-024D63A7BD29}.Release|Any CPU.Build.0 = Release|Any CPU - EndGlobalSection - GlobalSection(SolutionProperties) = preSolution - HideSolutionNode = FALSE - EndGlobalSection - GlobalSection(ExtensibilityGlobals) = postSolution - SolutionGuid = {852B6D28-02DF-4141-B595-B5A41C97A18D} - EndGlobalSection - GlobalSection(NestedProjects) = preSolution - {84394417-21C9-4C81-B049-C7E75ACB60A0} = {615B6C70-A67D-4901-B3F8-49205B2B61D9} - {63422754-EFA8-4EC7-91E7-024D63A7BD29} = {C28CF6C2-329B-4FF2-9A1C-FF73778BE7AB} - EndGlobalSection -EndGlobal diff --git a/ExtendedThreading.slnx b/ExtendedThreading.slnx new file mode 100644 index 0000000..9be0a23 --- /dev/null +++ b/ExtendedThreading.slnx @@ -0,0 +1,13 @@ + + + + + + + + + + + + + diff --git a/README.md b/README.md index ce8eda6..52e7afa 100644 --- a/README.md +++ b/README.md @@ -4,14 +4,35 @@ This package provides extended Threading functionality, built on top of the buil # Installation -I recommend using the NuGet package: [ExtendedThreading](https://www.nuget.org/packages/ExtendedThreading) however feel free to clone the source instead if that suits your needs -better. +I recommend using the NuGet package: [ExtendedThreading](https://www.nuget.org/packages/ExtendedThreading) however feel free to clone the source instead if that suits your needs better. # Usage +## AsyncSignal + +This is used to simplify signaling between async contexts, e.g. when building the Producer/Consumer pattern purely on async/await. + +``` +public class ProducerConsumer +{ + private AsyncSignal _signal = new(); + + public void Produce(T item){ + // produce + _signal.Pulse(); // Inform consumers that a new item is available + } + + public async Task Consume(CancellationToken cancelationToken) + { + await _signal.WaitAsync(cancellationToken); // Will block until an item becomes available or the token is cancelled + // consume + } +} +``` + ## ThreadSignal -This is used to simplify signalling between threads, e.g. when building the Producer/Consumer pattern: +This is used to simplify signaling between threads, e.g. when building the Producer/Consumer pattern: ``` public class ProducerConsumer @@ -33,8 +54,7 @@ public class ProducerConsumer ## KeyedMutexSynchronizer -This is used to ensure mutual exclusion based on keys. E.g. for an API where you want to grant only a single thread access to do PUT requests on a per-id basis to prevent race -conditions on a per entity basis: +This is used to ensure mutual exclusion based on keys. E.g. for an API where you want to grant only a single thread access to do PUT requests on a per-id basis to prevent race conditions on a per entity basis: ``` public class OrderController @@ -65,8 +85,7 @@ public class OrderController This class only offers one method: `WhenAll`. It functions similarly to the built-in `Task.WhenAll` in .Net, except for how it handles Exceptions. -This version throws an `AggregateException` in case any exceptions occur, to allow you the full picture of all exceptions, instead of just the first one (which is what the -built-in `Task.WhenAll` will throw) +This version throws an `AggregateException` in case any exceptions occur, to allow you the full picture of all exceptions, instead of just the first one (which is what the built-in `Task.WhenAll` will throw) ## AwaitExtensions @@ -84,4 +103,5 @@ var results = await (task1, task2); ``` # Documentation + Auto generated documentation via [DocFx](https://github.com/dotnet/docfx) is available here: https://steffenskov.github.io/ExtendedThreading/ \ No newline at end of file diff --git a/src/ExtendedThreading/AsyncSignal.cs b/src/ExtendedThreading/AsyncSignal.cs new file mode 100644 index 0000000..4b9a656 --- /dev/null +++ b/src/ExtendedThreading/AsyncSignal.cs @@ -0,0 +1,81 @@ +namespace ExtendedThreading; + +public class AsyncSignal +{ + private readonly object _lock = new(); + private readonly List> _waiters = []; + + public void Pulse() + { + lock (_lock) + { + if (_waiters.Count == 0) + { + return; + } + + var completionSource = _waiters[0]; + _waiters.RemoveAt(0); + + completionSource.TrySetResult(true); + } + } + + public void PulseAll() + { + lock (_lock) + { + foreach (var completionSource in _waiters) + { + completionSource.TrySetResult(true); + } + + _waiters.Clear(); + } + } + + public Task WaitAsync(CancellationToken cancellationToken = default) + { + return WaitAsync(Timeout.InfiniteTimeSpan, cancellationToken); + } + + public Task WaitAsync(int millisecondsTimeout, CancellationToken cancellationToken = default) + { + return WaitAsync(TimeSpan.FromMilliseconds(millisecondsTimeout), cancellationToken); + } + + public async Task WaitAsync(TimeSpan timeout, CancellationToken cancellationToken = default) + { + var completionSource = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + lock (_lock) + { + _waiters.Add(completionSource); + } + + using var timeoutSource = timeout == Timeout.InfiniteTimeSpan ? null : new CancellationTokenSource(timeout); + using var linkedSource = timeoutSource is null + ? null + : CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, timeoutSource.Token); + var token = linkedSource?.Token ?? cancellationToken; + + await using var registration = token.Register(static state => + ((TaskCompletionSource)state!).TrySetCanceled(), completionSource); + + try + { + return await completionSource.Task.ConfigureAwait(false); + } + catch (TaskCanceledException) + { + cancellationToken.ThrowIfCancellationRequested(); // real cancellation → rethrow + return false; // otherwise just timed out + } + finally + { + lock (_lock) + { + _waiters.Remove(completionSource); + } + } + } +} \ No newline at end of file diff --git a/src/ExtendedThreading/ExtendedThreading.csproj b/src/ExtendedThreading/ExtendedThreading.csproj index ddbb26f..f7b3f15 100644 --- a/src/ExtendedThreading/ExtendedThreading.csproj +++ b/src/ExtendedThreading/ExtendedThreading.csproj @@ -1,21 +1,13 @@ - - net8.0 - enable - true - true - Steffen Skov - Steffen Skov - A super small library for providing strong typed Ids (as opposed to using primitives) - Steffen Skov 2022 - MIT - https://github.com/steffenskov/ExtendedThreading - Threading - README.md - - - - + + A super small library for providing strong typed Ids (as opposed to using primitives) + https://github.com/steffenskov/ExtendedThreading + Threading + README.md + + + + diff --git a/tests/ExtendedThreading.UnitTests/AsyncSignalTests.cs b/tests/ExtendedThreading.UnitTests/AsyncSignalTests.cs new file mode 100644 index 0000000..680ba7e --- /dev/null +++ b/tests/ExtendedThreading.UnitTests/AsyncSignalTests.cs @@ -0,0 +1,232 @@ +using System.Collections; +using System.Reflection; + +namespace ExtendedThreading.UnitTests; + +public class AsyncSignalTests +{ + [Fact] + public void Pulse_NoWaiters_DoesNotThrow() + { + // Arrange + var signal = new AsyncSignal(); + + // Act + var exception = Record.Exception(signal.Pulse); + + // Assert + Assert.Null(exception); + } + + [Fact] + public void PulseAll_NoWaiters_DoesNotThrow() + { + // Arrange + var signal = new AsyncSignal(); + + // Act + var exception = Record.Exception(signal.PulseAll); + + // Assert + Assert.Null(exception); + } + + [Fact] + public async Task WaitAsync_ParameterlessPulsedBeforeCompletion_CompletesSuccessfully() + { + // Arrange + var signal = new AsyncSignal(); + + // Act + var waitTask = signal.WaitAsync(TestContext.Current.CancellationToken); + await WaitUntilWaiterRegisteredAsync(signal); + signal.Pulse(); + var completedTask = await Task.WhenAny(waitTask, Task.Delay(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken)); + + // Assert + Assert.Same(waitTask, completedTask); + } + + [Fact] + public async Task WaitAsync_MillisecondsTimeoutPulsedBeforeTimeout_ReturnsTrue() + { + // Arrange + var signal = new AsyncSignal(); + var waitTask = signal.WaitAsync(5000, TestContext.Current.CancellationToken); + await WaitUntilWaiterRegisteredAsync(signal); + + // Act + signal.Pulse(); + var result = await waitTask; + + // Assert + Assert.True(result); + } + + [Fact] + public async Task WaitAsync_MillisecondsTimeoutElapses_ReturnsFalse() + { + // Arrange + var signal = new AsyncSignal(); + + // Act + var result = await signal.WaitAsync(50, TestContext.Current.CancellationToken); + + // Assert + Assert.False(result); + } + + [Fact] + public async Task WaitAsync_TimeSpanTimeoutPulsedBeforeTimeout_ReturnsTrue() + { + // Arrange + var signal = new AsyncSignal(); + var waitTask = signal.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + await WaitUntilWaiterRegisteredAsync(signal); + + // Act + signal.Pulse(); + var result = await waitTask; + + // Assert + Assert.True(result); + } + + [Fact] + public async Task WaitAsync_TimeSpanTimeoutElapses_ReturnsFalse() + { + // Arrange + var signal = new AsyncSignal(); + + // Act + var result = await signal.WaitAsync(TimeSpan.FromMilliseconds(50), TestContext.Current.CancellationToken); + + // Assert + Assert.False(result); + } + + [Fact] + public async Task WaitAsync_ExternalTokenCancelledBeforeTimeout_ThrowsOperationCanceledException() + { + // Arrange + var signal = new AsyncSignal(); + using var cts = new CancellationTokenSource(); + + // Act + var waitTask = signal.WaitAsync(TimeSpan.FromSeconds(5), cts.Token); + await WaitUntilWaiterRegisteredAsync(signal); + await cts.CancelAsync(); + + // Assert + await Assert.ThrowsAsync(() => waitTask); + } + + [Fact] + public async Task WaitAsync_ExternalTokenAlreadyCancelled_ThrowsOperationCanceledException() + { + // Arrange + var signal = new AsyncSignal(); + using var cts = new CancellationTokenSource(); + await cts.CancelAsync(); + + // Act & Assert + await Assert.ThrowsAsync(() => signal.WaitAsync(TimeSpan.FromSeconds(5), cts.Token)); + } + + [Fact] + public async Task Pulse_SingleWaiter_WakesThatWaiter() + { + // Arrange + var signal = new AsyncSignal(); + var waitTask = signal.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + await WaitUntilWaiterRegisteredAsync(signal); + + // Act + signal.Pulse(); + var result = await waitTask; + + // Assert + Assert.True(result); + } + + [Fact] + public async Task Pulse_MultipleWaiters_WakesExactlyOne() + { + // Arrange + var signal = new AsyncSignal(); + var waitTask1 = signal.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + var waitTask2 = signal.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + await WaitUntilWaiterRegisteredAsync(signal, 2); + + // Act + signal.Pulse(); + var firstCompleted = await Task.WhenAny(waitTask1, waitTask2); + var stillPending = firstCompleted == waitTask1 ? waitTask2 : waitTask1; + + // Assert + Assert.True(await firstCompleted); + Assert.False(stillPending.IsCompleted); + } + + [Fact] + public async Task PulseAll_MultipleWaiters_WakesAllOfThem() + { + // Arrange + var signal = new AsyncSignal(); + var waitTask1 = signal.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + var waitTask2 = signal.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + var waitTask3 = signal.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + await WaitUntilWaiterRegisteredAsync(signal, 3); + + // Act + signal.PulseAll(); + var results = await Task.WhenAll(waitTask1, waitTask2, waitTask3); + + // Assert + Assert.All(results, Assert.True); + } + + [Fact] + public async Task WaitAsync_TimedOutWaiterFollowedByPulse_DoesNotConsumeStaleWaiter() + { + // Arrange: first waiter times out and should be pruned from the internal list, + // so a later Pulse() must reach the second, still-live waiter instead of + // silently completing against the already-cancelled first one. + var signal = new AsyncSignal(); + var timedOutTask = signal.WaitAsync(50, TestContext.Current.CancellationToken); + await WaitUntilWaiterRegisteredAsync(signal); + var timedOutResult = await timedOutTask; + + var liveWaitTask = signal.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + await WaitUntilWaiterRegisteredAsync(signal); + + // Act + signal.Pulse(); + var liveResult = await liveWaitTask; + + // Assert + Assert.False(timedOutResult); + Assert.True(liveResult); + } + + // Polls the internal waiter count via reflection so tests don't race the + // lock(_lock) { _waiters.Add(tcs); } line inside WaitAsyncCore. + private static async Task WaitUntilWaiterRegisteredAsync(AsyncSignal signal, int expectedCount = 1) + { + var field = typeof(AsyncSignal) + .GetField("_waiters", BindingFlags.NonPublic | BindingFlags.Instance)!; + + for (var i = 0; i < 100; i++) + { + var waiters = (IList)field.GetValue(signal)!; + if (waiters.Count >= expectedCount) + { + return; + } + + await Task.Delay(10); + } + + throw new TimeoutException("Waiter was not registered in time."); + } +} \ No newline at end of file diff --git a/tests/ExtendedThreading.UnitTests/ExtendedThreading.UnitTests.csproj b/tests/ExtendedThreading.UnitTests/ExtendedThreading.UnitTests.csproj index 81e21fc..e3b1d41 100644 --- a/tests/ExtendedThreading.UnitTests/ExtendedThreading.UnitTests.csproj +++ b/tests/ExtendedThreading.UnitTests/ExtendedThreading.UnitTests.csproj @@ -1,7 +1,7 @@ - net8.0 + net10.0 enable enable true