From 2356216c872446997413ce8ffd43bb3634900877 Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Tue, 21 Jul 2026 22:22:32 +0100 Subject: [PATCH 1/4] refactor!: freeze a clean v4 public API surface Now-or-never cleanup before the 4.0.0 contract freezes: - Internalize the ActionTaskWrapper/ItemTaskWrapper structs and demote the abstract processor bases' protected plumbing (TaskWrappers, EnumerableTaskCompletionSources, CancellationToken, constructors, DisposeAsyncCore) to private protected - none of it was usable outside the assembly because Process() is internal abstract. - Seal every leaf processor and builder class. - Move IAsyncEnumerableProcessor to EnumerableAsyncProcessor.Interfaces and make both variants IAsyncDisposable/IDisposable. The six async-enumerable processors now dispose their linked CancellationTokenSource when ExecuteAsync completes (previously leaked a registration on the caller's token) and support explicit disposal. - Honor scheduleOnThreadPool on the bounded parallel paths (it was silently ignored whenever maxConcurrency was set). - Remove the parameterized no-selector IAsyncEnumerable ProcessInParallel overloads whose maxConcurrency/scheduleOnThreadPool did nothing. - Keep binary compatibility with assemblies compiled against v3, notably TUnit.Engine, which this repo's own test runner depends on: restore the parameterless ProcessInParallel(items, ct) collect overload and add ProcessInParallel(int) forwarders on the four enumerable builders. V3BinaryCompatibilityTests pins the exact signatures; this also fixes --treenode-filter discovery, broken since the v4 API consolidation. - Packaging: generate XML docs (IntelliSense was missing from the package), mark AOT-compatible, merge duplicate metadata groups. - Docs: correct stale CLAUDE.md/README claims (RateLimitedParallel strategy, minimumIterationTime model, disposal contract) and expand the migration guide. --- CLAUDE.md | 10 +- .../ProcessInParallelExample.cs | 6 +- .../AsyncEnumerableParallelExtensionsTests.cs | 38 ++----- .../DisposalRegressionTests.cs | 42 +++++++ .../V3BinaryCompatibilityTests.cs | 47 ++++++++ .../Builders/ActionAsyncProcessorBuilder.cs | 11 +- .../Builders/ActionAsyncProcessorBuilder_1.cs | 11 +- ...EnumerableActionAsyncProcessorBuilder_1.cs | 4 +- ...EnumerableActionAsyncProcessorBuilder_2.cs | 4 +- .../AsyncEnumerableAsyncProcessorBuilder.cs | 2 +- .../ExecutionCountAsyncProcessorBuilder.cs | 2 +- .../ItemActionAsyncProcessorBuilder_1.cs | 11 +- .../ItemActionAsyncProcessorBuilder_2.cs | 11 +- .../Builders/ItemAsyncProcessorBuilder.cs | 2 +- .../EnumerableAsyncProcessor.csproj | 22 ++-- .../Extensions/AsyncEnumerableExtensions.cs | 36 ++---- .../Extensions/EnumerableExtensions.cs | 12 -- .../Interfaces/IAsyncEnumerableProcessor.cs | 24 +++- .../PublicAPI.Shipped.txt | 86 +++++---------- .../Abstract/AbstractAsyncProcessor.cs | 6 +- .../Abstract/AbstractAsyncProcessorBase.cs | 8 +- .../Abstract/AbstractAsyncProcessor_1.cs | 6 +- .../AsyncEnumerableBatchProcessor.cs | 53 ++++++--- .../AsyncEnumerableOneAtATimeProcessor.cs | 28 ++++- .../AsyncEnumerableParallelProcessor.cs | 46 +++++--- .../ResultAsyncEnumerableBatchProcessor.cs | 60 ++++++---- ...esultAsyncEnumerableOneAtATimeProcessor.cs | 30 ++++- .../ResultAsyncEnumerableParallelProcessor.cs | 104 +++++++++++------- .../RunnableProcessors/BatchAsyncProcessor.cs | 2 +- .../BatchAsyncProcessor_1.cs | 2 +- .../OneAtATimeAsyncProcessor.cs | 2 +- .../OneAtATimeAsyncProcessor_1.cs | 2 +- .../ParallelAsyncProcessor.cs | 2 +- .../ParallelAsyncProcessor_1.cs | 2 +- .../ResultAbstractAsyncProcessorBase.cs | 8 +- .../ResultAbstractAsyncProcessor_1.cs | 6 +- .../ResultAbstractAsyncProcessor_2.cs | 6 +- .../ResultBatchAsyncProcessor_1.cs | 2 +- .../ResultBatchAsyncProcessor_2.cs | 2 +- .../ResultOneAtATimeAsyncProcessor_1.cs | 2 +- .../ResultOneAtATimeAsyncProcessor_2.cs | 2 +- .../ResultParallelAsyncProcessor_1.cs | 2 +- .../ResultParallelAsyncProcessor_2.cs | 2 +- ...imedRateLimitedParallelAsyncProcessor_1.cs | 2 +- ...imedRateLimitedParallelAsyncProcessor_2.cs | 2 +- .../TimedRateLimitedParallelAsyncProcessor.cs | 2 +- ...imedRateLimitedParallelAsyncProcessor_1.cs | 2 +- EnumerableAsyncProcessor/TaskWrapper.cs | 8 +- README.md | 12 +- 49 files changed, 486 insertions(+), 308 deletions(-) create mode 100644 EnumerableAsyncProcessor.UnitTests/V3BinaryCompatibilityTests.cs diff --git a/CLAUDE.md b/CLAUDE.md index cea6d16..e8024c7 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -29,6 +29,8 @@ dotnet run --project EnumerableAsyncProcessor.UnitTests -f net10.0 -- --treenode TUnit test projects compile to executables; VSTest-style `dotnet test --filter` does not work. See the `tunit-testing` skill for full filter syntax. +**TUnit itself depends on this library.** `TUnit.Engine` is compiled against EnumerableAsyncProcessor v3, and the locally built assembly shadows the copy TUnit shipped with (same assembly identity), so removing or changing a public member that TUnit.Engine binds to crashes test discovery with `MissingMethodException` before any test runs. The exact signatures TUnit needs are pinned by `V3BinaryCompatibilityTests` and marked as binary-compat members in the source — do not remove them until TUnit ships a build compiled against v4. + CI (`.github/workflows/dotnet.yml`) runs the `EnumerableAsyncProcessor.Pipeline` project (a ModularPipelines app, `dotnet run -c Release` from that directory), which builds, tests, packs, and — on `main` — publishes to NuGet. Versioning comes from GitVersion (`GitVersion.yml` pins `next-version: 4.0.0`; keep that file present, its absence makes ModularPipelines generate a Mainline config that crashes on GitHub PR merge commits). ## Architecture @@ -45,17 +47,17 @@ Processor classes vary along three axes, reflected in naming: - **Input**: with items (``) vs. execution-count only (non-generic). - **Output**: `Result*`-prefixed classes (in `RunnableProcessors/ResultProcessors/`) return values via `IAsyncProcessor` (`GetResultsAsync()`, `GetResultsAsyncEnumerable()`, `GetEnumerableTasks()`); unprefixed classes are fire-and-await (`WaitAsync()`). -- **Strategy**: `OneAtATime`, `Batch`, `Parallel`, `RateLimitedParallel`, `TimedRateLimitedParallel`. +- **Strategy**: `OneAtATime`, `Batch`, `Parallel`, `TimedRateLimitedParallel`. -`RunnableProcessors/AsyncEnumerable/` holds parallel variants for `IAsyncEnumerable` sources. File-name suffixes `_1`/`_2` distinguish generic arity (e.g. `BatchAsyncProcessor_1.cs` is `BatchAsyncProcessor`). +`RunnableProcessors/AsyncEnumerable/` holds the `IAsyncEnumerable`-source variants (Parallel, OneAtATime, Batch). File-name suffixes `_1`/`_2` distinguish generic arity (e.g. `BatchAsyncProcessor_1.cs` is `BatchAsyncProcessor`). ### Core mechanics (read these before changing behavior) - **`ProcessorLifecycle.cs`**: owns start/cancel/dispose shared by both base-class hierarchies. `AbstractAsyncProcessorBase` (void) and `ResultAbstractAsyncProcessorBase` (results) cannot share an ancestor because they fan out to differently typed `TaskCompletionSource` lists, so both delegate to this class. Cancellation is registered in `Start`, not the constructor, so a pre-cancelled token can never fire on a partially built instance. `DisposeAsync` waits up to 30 seconds for in-flight tasks; sync `Dispose` cancels without blocking. - **TCS-per-item**: each item gets a `TaskCompletionSource`; `TaskWrapper.Process` never throws — it completes the item's TCS with the failure/cancellation instead, so one failed item cannot kill the run or leave awaiters hanging. -- **`WorkerPool.cs`**: rate-limited processors run a fixed pool of worker loops claiming items via `Interlocked.Increment` (P `Task.Run` tasks total, not N throttled tasks + semaphore). `minimumIterationTime` is how timed rate limiting is implemented: each worker holds its slot for at least that duration per item. +- **`WorkerPool.cs`**: rate-limited processors run a fixed pool of worker loops claiming items via `Interlocked.Increment` (P `Task.Run` tasks total, not N throttled tasks + semaphore). Timed rate limiting is a shared `TokenBucketRateLimiter` (`System.Threading.RateLimiting`): workers acquire a permit before starting each item, so `permitsPerWindow`/`window` bound the start rate independently of `maxConcurrency`. - **Multi-targeting**: `EnumerableExtensions.ToIAsyncEnumerable` uses `Task.WhenEach` on `NET9_0_OR_GREATER` and a completion-order-bucket fallback otherwise. The test project targets `net8.0` specifically to exercise the fallback path — don't drop that TFM. ### Disposal contract -All processors implement `IDisposable`/`IAsyncDisposable`; the README documents the patterns users rely on (`await using`, safe double/early disposal). The builder extension shortcuts on `IAsyncEnumerable` dispose internally; processors returned from the builder pattern are the caller's responsibility. Preserve these semantics — there are dedicated regression tests (`DisposalRegressionTests`, `ExceptionFidelityTests`, `InputEnumerationRegressionTests`). +All processors implement `IDisposable`/`IAsyncDisposable`; the README documents the patterns users rely on (`await using`, safe double/early disposal). `IAsyncEnumerableProcessor` implementations are single-use and additionally dispose their internal linked `CancellationTokenSource` when `ExecuteAsync` completes; `IAsyncProcessor` objects returned from the builder pattern are the caller's responsibility. Preserve these semantics — there are dedicated regression tests (`DisposalRegressionTests`, `ExceptionFidelityTests`, `InputEnumerationRegressionTests`). diff --git a/EnumerableAsyncProcessor.Example/ProcessInParallelExample.cs b/EnumerableAsyncProcessor.Example/ProcessInParallelExample.cs index f47457e..d1fa922 100644 --- a/EnumerableAsyncProcessor.Example/ProcessInParallelExample.cs +++ b/EnumerableAsyncProcessor.Example/ProcessInParallelExample.cs @@ -15,10 +15,10 @@ public static async Task RunExample() Console.WriteLine("ProcessInParallel Extension Examples"); Console.WriteLine("====================================\n"); - // Example 1: Simple parallel processing without transformation - Console.WriteLine("Example 1: Simple parallel processing (no transformation needed!)"); + // Example 1: Simple parallel processing with an identity transformation + Console.WriteLine("Example 1: Simple parallel processing"); var asyncEnumerable1 = GenerateAsyncEnumerable(5); - IEnumerable results1 = await asyncEnumerable1.ProcessInParallel(); // <-- This is the simple extension! + IEnumerable results1 = await asyncEnumerable1.ProcessInParallel(item => Task.FromResult(item)); Console.WriteLine($"Results: {string.Join(", ", results1)}"); // Example 2: Parallel processing with transformation diff --git a/EnumerableAsyncProcessor.UnitTests/AsyncEnumerableParallelExtensionsTests.cs b/EnumerableAsyncProcessor.UnitTests/AsyncEnumerableParallelExtensionsTests.cs index cfa0577..151a0b0 100644 --- a/EnumerableAsyncProcessor.UnitTests/AsyncEnumerableParallelExtensionsTests.cs +++ b/EnumerableAsyncProcessor.UnitTests/AsyncEnumerableParallelExtensionsTests.cs @@ -32,24 +32,13 @@ private static async IAsyncEnumerable GenerateDelayedAsyncEnumerable(int co } } - [Test] - public async Task ProcessInParallel_WithoutTransformation_ReturnsAllItems() - { - var asyncEnumerable = GenerateAsyncEnumerable(10); - - var results = await asyncEnumerable.ProcessInParallel(); - - await Assert.That(results.Count()).IsEqualTo(10); - await Assert.That(results.OrderBy(x => x)).IsEquivalentTo(Enumerable.Range(1, 10)); - } - [Test] public async Task ProcessInParallel_WithMaxConcurrency_ReturnsAllItems() { var asyncEnumerable = GenerateAsyncEnumerable(20); - - var results = await asyncEnumerable.ProcessInParallel(5); - + + var results = await asyncEnumerable.ProcessInParallel(item => Task.FromResult(item), 5); + await Assert.That(results.Count()).IsEqualTo(20); await Assert.That(results.OrderBy(x => x)).IsEquivalentTo(Enumerable.Range(1, 20)); } @@ -94,7 +83,7 @@ public async Task ProcessInParallel_WithCancellation_ThrowsOperationCanceledExce using var cts = new CancellationTokenSource(); var asyncEnumerable = GenerateDelayedAsyncEnumerable(100, 50); - var task = asyncEnumerable.ProcessInParallel(cancellationToken: cts.Token); + var task = asyncEnumerable.ProcessInParallel(item => Task.FromResult(item), cancellationToken: cts.Token); // Cancel after a short delay cts.CancelAfter(100); @@ -160,23 +149,10 @@ public async Task ProcessInParallel_WithMaxConcurrency_LimitsConcurrency() public async Task ProcessInParallel_EmptyEnumerable_ReturnsEmptyResult() { var asyncEnumerable = GenerateAsyncEnumerable(0); - - var results = await asyncEnumerable.ProcessInParallel(); - - await Assert.That(results.Count()).IsEqualTo(0); - } - [Test] - public async Task ProcessInParallel_WithScheduleOnThreadPool_ProcessesAllItems() - { - var asyncEnumerable = GenerateAsyncEnumerable(10); - - var results = await asyncEnumerable.ProcessInParallel( - maxConcurrency: null, - scheduleOnThreadPool: true); - - await Assert.That(results.Count()).IsEqualTo(10); - await Assert.That(results.OrderBy(x => x)).IsEquivalentTo(Enumerable.Range(1, 10)); + var results = await asyncEnumerable.ProcessInParallel(item => Task.FromResult(item)); + + await Assert.That(results.Count()).IsEqualTo(0); } [Test] diff --git a/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs b/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs index 1683625..e49995c 100644 --- a/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs +++ b/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs @@ -1,6 +1,8 @@ using System; +using System.Collections.Generic; using System.Diagnostics; using System.Linq; +using System.Threading; using System.Threading.Tasks; using EnumerableAsyncProcessor.Extensions; @@ -127,4 +129,44 @@ public async Task Disposal_Is_Idempotent_And_Safe_In_Any_Order() await Assert.That(processor.GetEnumerableTasks().Count(x => x.IsCompletedSuccessfully)).IsEqualTo(5); } + + [Test] + public async Task AsyncEnumerable_Processor_Disposal_Is_Idempotent_After_Execution() + { + var processedCount = 0; + + var processor = GenerateAsyncEnumerable(5) + .ForEachAsync(_ => + { + Interlocked.Increment(ref processedCount); + return Task.CompletedTask; + }) + .ProcessInParallel(maxConcurrency: 2); + + await processor.ExecuteAsync(); + + // ExecuteAsync disposes internal resources on completion; explicit disposal stays safe. + await processor.DisposeAsync(); + processor.Dispose(); + + await Assert.That(processedCount).IsEqualTo(5); + } + + [Test] + public async Task AsyncEnumerable_Result_Processor_Supports_Await_Using_Without_Execution() + { + await using (GenerateAsyncEnumerable(3).SelectAsync(i => Task.FromResult(i)).ProcessInParallel(2)) + { + // Never executed - disposal alone must not throw. + } + } + + private static async IAsyncEnumerable GenerateAsyncEnumerable(int count) + { + for (var i = 0; i < count; i++) + { + await Task.Yield(); + yield return i; + } + } } diff --git a/EnumerableAsyncProcessor.UnitTests/V3BinaryCompatibilityTests.cs b/EnumerableAsyncProcessor.UnitTests/V3BinaryCompatibilityTests.cs new file mode 100644 index 0000000..d6759e4 --- /dev/null +++ b/EnumerableAsyncProcessor.UnitTests/V3BinaryCompatibilityTests.cs @@ -0,0 +1,47 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using EnumerableAsyncProcessor.Builders; +using EnumerableAsyncProcessor.Extensions; + +namespace EnumerableAsyncProcessor.UnitTests; + +/// +/// TUnit.Engine (the framework running this suite) is itself compiled against +/// EnumerableAsyncProcessor, and the locally built assembly shadows the version TUnit shipped +/// with because the assembly identity matches. If a member TUnit.Engine binds to disappears, +/// test DISCOVERY crashes with MissingMethodException before any test runs. These are the +/// exact signatures TUnit.Engine references - keep them until TUnit rebuilds against v4. +/// +public class V3BinaryCompatibilityTests +{ + [Test] + public async Task Members_Bound_By_TUnit_Engine_Exist_With_Exact_Signatures() + { + // ItemActionAsyncProcessorBuilder.ProcessInParallel(int) + var builderMethod = typeof(ItemActionAsyncProcessorBuilder<,>).GetMethods() + .SingleOrDefault(m => m.Name == "ProcessInParallel" + && m.GetParameters().Length == 1 + && m.GetParameters()[0].ParameterType == typeof(int)); + await Assert.That(builderMethod).IsNotNull(); + + // AsyncEnumerableExtensions.ProcessInParallel(IAsyncEnumerable, CancellationToken) + var extensionMethod = typeof(AsyncEnumerableExtensions).GetMethods() + .SingleOrDefault(m => m.Name == "ProcessInParallel" + && m.GetGenericArguments().Length == 1 + && m.GetParameters().Length == 2 + && m.GetParameters()[1].ParameterType == typeof(CancellationToken)); + await Assert.That(extensionMethod).IsNotNull(); + + // TUnit.Engine also binds EnumerableExtensions.SelectAsync / SelectManyAsync and + // IAsyncProcessor.GetAwaiter; these exact-signature method groups stop compiling if they drift. + Func, Func>, CancellationToken, ItemActionAsyncProcessorBuilder> selectAsync = + EnumerableExtensions.SelectAsync; + Func, Func>, CancellationToken, IAsyncEnumerable> selectManyAsync = + EnumerableExtensions.SelectManyAsync; + await Assert.That(selectAsync).IsNotNull(); + await Assert.That(selectManyAsync).IsNotNull(); + } +} diff --git a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs index 9c2c1b2..7396783 100644 --- a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs +++ b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs @@ -4,7 +4,7 @@ namespace EnumerableAsyncProcessor.Builders; -public class ActionAsyncProcessorBuilder +public sealed class ActionAsyncProcessorBuilder { private readonly int _count; private readonly Func _taskSelector; @@ -53,6 +53,15 @@ public IAsyncProcessor ProcessInParallel(int? maxConcurrency = null, bool schedu { return new ParallelAsyncProcessor(_count, _taskSelector, _cancellationTokenSource, maxConcurrency, scheduleOnThreadPool).StartProcessing(); } + + /// + /// Processes items in parallel with bounded concurrency. Binary-compatible with assemblies + /// compiled against v3 (equivalent to ProcessInParallel(maxConcurrency: n)). + /// + public IAsyncProcessor ProcessInParallel(int maxConcurrency) + { + return ProcessInParallel((int?)maxConcurrency); + } public IAsyncProcessor ProcessOneAtATime() { diff --git a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs index 65b4ec9..96fc06b 100644 --- a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs @@ -4,7 +4,7 @@ namespace EnumerableAsyncProcessor.Builders; -public class ActionAsyncProcessorBuilder +public sealed class ActionAsyncProcessorBuilder { private readonly int _count; private readonly Func> _taskSelector; @@ -53,6 +53,15 @@ public IAsyncProcessor ProcessInParallel(int? maxConcurrency = null, bo { return new ResultParallelAsyncProcessor(_count, _taskSelector, _cancellationTokenSource, maxConcurrency, scheduleOnThreadPool).StartProcessing(); } + + /// + /// Processes items in parallel with bounded concurrency. Binary-compatible with assemblies + /// compiled against v3 (equivalent to ProcessInParallel(maxConcurrency: n)). + /// + public IAsyncProcessor ProcessInParallel(int maxConcurrency) + { + return ProcessInParallel((int?)maxConcurrency); + } public IAsyncProcessor ProcessOneAtATime() { diff --git a/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_1.cs index a993ac3..83a72a0 100644 --- a/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_1.cs @@ -1,9 +1,9 @@ -using EnumerableAsyncProcessor.Extensions; +using EnumerableAsyncProcessor.Interfaces; using EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable; namespace EnumerableAsyncProcessor.Builders; -public class AsyncEnumerableActionAsyncProcessorBuilder +public sealed class AsyncEnumerableActionAsyncProcessorBuilder { private readonly IAsyncEnumerable _items; private readonly Func _taskSelector; diff --git a/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_2.cs b/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_2.cs index edfd1b1..b40f21b 100644 --- a/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_2.cs +++ b/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_2.cs @@ -1,9 +1,9 @@ -using EnumerableAsyncProcessor.Extensions; +using EnumerableAsyncProcessor.Interfaces; using EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors; namespace EnumerableAsyncProcessor.Builders; -public class AsyncEnumerableActionAsyncProcessorBuilder +public sealed class AsyncEnumerableActionAsyncProcessorBuilder { private readonly IAsyncEnumerable _items; private readonly Func> _taskSelector; diff --git a/EnumerableAsyncProcessor/Builders/AsyncEnumerableAsyncProcessorBuilder.cs b/EnumerableAsyncProcessor/Builders/AsyncEnumerableAsyncProcessorBuilder.cs index 08b60d1..dc24c17 100644 --- a/EnumerableAsyncProcessor/Builders/AsyncEnumerableAsyncProcessorBuilder.cs +++ b/EnumerableAsyncProcessor/Builders/AsyncEnumerableAsyncProcessorBuilder.cs @@ -1,6 +1,6 @@ namespace EnumerableAsyncProcessor.Builders; -public class AsyncEnumerableAsyncProcessorBuilder +public sealed class AsyncEnumerableAsyncProcessorBuilder { private readonly IAsyncEnumerable _items; diff --git a/EnumerableAsyncProcessor/Builders/ExecutionCountAsyncProcessorBuilder.cs b/EnumerableAsyncProcessor/Builders/ExecutionCountAsyncProcessorBuilder.cs index c1c7936..9023740 100644 --- a/EnumerableAsyncProcessor/Builders/ExecutionCountAsyncProcessorBuilder.cs +++ b/EnumerableAsyncProcessor/Builders/ExecutionCountAsyncProcessorBuilder.cs @@ -1,6 +1,6 @@ namespace EnumerableAsyncProcessor.Builders; -public class ExecutionCountAsyncProcessorBuilder +public sealed class ExecutionCountAsyncProcessorBuilder { private readonly int _count; diff --git a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs index 402d3fe..91d0561 100644 --- a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs @@ -4,7 +4,7 @@ namespace EnumerableAsyncProcessor.Builders; -public class ItemActionAsyncProcessorBuilder +public sealed class ItemActionAsyncProcessorBuilder { private readonly IEnumerable _items; private readonly Func _taskSelector; @@ -56,6 +56,15 @@ public IAsyncProcessor ProcessInParallel(int? maxConcurrency = null, bool schedu return new ParallelAsyncProcessor(_items, _taskSelector, _cancellationTokenSource, maxConcurrency, scheduleOnThreadPool) .StartProcessing(); } + + /// + /// Processes items in parallel with bounded concurrency. Binary-compatible with assemblies + /// compiled against v3 (equivalent to ProcessInParallel(maxConcurrency: n)). + /// + public IAsyncProcessor ProcessInParallel(int maxConcurrency) + { + return ProcessInParallel((int?)maxConcurrency); + } public IAsyncProcessor ProcessOneAtATime() { diff --git a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs index ca902a6..a722c85 100644 --- a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs +++ b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs @@ -4,7 +4,7 @@ namespace EnumerableAsyncProcessor.Builders; -public class ItemActionAsyncProcessorBuilder +public sealed class ItemActionAsyncProcessorBuilder { private readonly IEnumerable _items; private readonly Func> _taskSelector; @@ -82,6 +82,15 @@ public IAsyncProcessor ProcessInParallel(int? maxConcurrency = null, bo { return new ResultParallelAsyncProcessor(_items, _taskSelector, _cancellationTokenSource, maxConcurrency, scheduleOnThreadPool).StartProcessing(); } + + /// + /// Processes items in parallel with bounded concurrency. Binary-compatible with assemblies + /// compiled against v3 (equivalent to ProcessInParallel(maxConcurrency: n)). + /// + public IAsyncProcessor ProcessInParallel(int maxConcurrency) + { + return ProcessInParallel((int?)maxConcurrency); + } /// /// Process items one at a time sequentially. diff --git a/EnumerableAsyncProcessor/Builders/ItemAsyncProcessorBuilder.cs b/EnumerableAsyncProcessor/Builders/ItemAsyncProcessorBuilder.cs index 1d2e36a..89e5e9c 100644 --- a/EnumerableAsyncProcessor/Builders/ItemAsyncProcessorBuilder.cs +++ b/EnumerableAsyncProcessor/Builders/ItemAsyncProcessorBuilder.cs @@ -1,6 +1,6 @@ namespace EnumerableAsyncProcessor.Builders; -public class ItemAsyncProcessorBuilder +public sealed class ItemAsyncProcessorBuilder { private readonly IEnumerable _items; diff --git a/EnumerableAsyncProcessor/EnumerableAsyncProcessor.csproj b/EnumerableAsyncProcessor/EnumerableAsyncProcessor.csproj index 92f39e8..3c58fa3 100644 --- a/EnumerableAsyncProcessor/EnumerableAsyncProcessor.csproj +++ b/EnumerableAsyncProcessor/EnumerableAsyncProcessor.csproj @@ -6,15 +6,23 @@ enable latest 99.99.99 + true + true $(WarningsAsErrors);RS0016;RS0017 - - $(NoWarn);RS0026 + + $(NoWarn);RS0026;CS1591 MIT README.md Tom Longhurst + Various Enumerable Async Processors - Batch / Parallel / Rate Limited / One at a time + git + https://github.com/thomhurst/EnumerableAsyncProcessor + https://github.com/thomhurst/EnumerableAsyncProcessor + async enumerable ienumerable linq array list processor delegate task tasks snupkg true true @@ -42,14 +50,4 @@ - - Tom Longhurst - Various Enumerable Async Processors - Batch / Parallel / Rate Limited / One at a time - git - https://github.com/thomhurst/EnumerableAsyncProcessor - https://github.com/thomhurst/EnumerableAsyncProcessor - async enumerable ienumerable linq array list processor delegate task tasks - - - diff --git a/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs b/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs index 0bf1eea..0459eaa 100644 --- a/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs +++ b/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs @@ -151,34 +151,12 @@ public static async IAsyncEnumerable SelectManyAsync( } /// - /// Process items in parallel and return all results as IEnumerable when awaited. + /// Collects every item from the source into a list. This overload has no work to parallelize; + /// it is kept for binary compatibility with assemblies compiled against v3 (notably TUnit). /// public static async Task> ProcessInParallel( this IAsyncEnumerable items, CancellationToken cancellationToken = default) - { - return await items.ProcessInParallel(null, false, cancellationToken).ConfigureAwait(false); - } - - /// - /// Process items in parallel with specified concurrency and return all results as IEnumerable when awaited. - /// - public static async Task> ProcessInParallel( - this IAsyncEnumerable items, - int maxConcurrency, - CancellationToken cancellationToken = default) - { - return await items.ProcessInParallel((int?)maxConcurrency, false, cancellationToken).ConfigureAwait(false); - } - - /// - /// Process items in parallel with optional concurrency and thread pool scheduling, return all results as IEnumerable when awaited. - /// - public static async Task> ProcessInParallel( - this IAsyncEnumerable items, - int? maxConcurrency, - bool scheduleOnThreadPool, - CancellationToken cancellationToken = default) { var results = new List(); @@ -189,7 +167,7 @@ public static async Task> ProcessInParallel( return results; } - + /// /// Process items in parallel with transformation and return all results as IEnumerable when awaited. /// @@ -224,12 +202,16 @@ public static async Task> ProcessInParallel( CancellationToken cancellationToken = default) { var results = new List(); - + if (maxConcurrency.HasValue) { + Func> effectiveSelector = scheduleOnThreadPool + ? item => Task.Run(() => taskSelector(item), cancellationToken) + : taskSelector; + await foreach (var result in AsyncEnumerableWorkerPool.ProcessResultsAsync( items, - taskSelector, + effectiveSelector, maxConcurrency.Value, cancellationToken).ConfigureAwait(false)) { diff --git a/EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs b/EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs index 72bbe0f..edf1af9 100644 --- a/EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs +++ b/EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs @@ -138,18 +138,6 @@ public static async IAsyncEnumerable SelectManyAsync( } } - private static async IAsyncEnumerable ToAsyncEnumerable( - this IEnumerable items, - [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken = default) - { - foreach (var item in items) - { - cancellationToken.ThrowIfCancellationRequested(); - yield return item; - } - await Task.CompletedTask; // Suppress CS1998 warning - } - internal static async IAsyncEnumerable ToIAsyncEnumerable(this IEnumerable> tasks, [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken = default) { #if NET9_0_OR_GREATER diff --git a/EnumerableAsyncProcessor/Interfaces/IAsyncEnumerableProcessor.cs b/EnumerableAsyncProcessor/Interfaces/IAsyncEnumerableProcessor.cs index 9553f78..728824f 100644 --- a/EnumerableAsyncProcessor/Interfaces/IAsyncEnumerableProcessor.cs +++ b/EnumerableAsyncProcessor/Interfaces/IAsyncEnumerableProcessor.cs @@ -1,11 +1,27 @@ -namespace EnumerableAsyncProcessor.Extensions; +namespace EnumerableAsyncProcessor.Interfaces; -public interface IAsyncEnumerableProcessor +/// +/// A single-use processor for an source that performs +/// an operation per item without returning results. +/// +public interface IAsyncEnumerableProcessor : IAsyncDisposable, IDisposable { + /// + /// Processes the source. The processor is single-use; it disposes its internal + /// resources when processing completes. + /// Task ExecuteAsync(); } -public interface IAsyncEnumerableProcessor +/// +/// A single-use processor for an source that streams +/// one result per item. +/// +public interface IAsyncEnumerableProcessor : IAsyncDisposable, IDisposable { + /// + /// Processes the source, streaming results as they become available. The processor is + /// single-use; it disposes its internal resources when enumeration finishes. + /// IAsyncEnumerable ExecuteAsync(); -} \ No newline at end of file +} diff --git a/EnumerableAsyncProcessor/PublicAPI.Shipped.txt b/EnumerableAsyncProcessor/PublicAPI.Shipped.txt index 24982ec..3464632 100644 --- a/EnumerableAsyncProcessor/PublicAPI.Shipped.txt +++ b/EnumerableAsyncProcessor/PublicAPI.Shipped.txt @@ -1,24 +1,16 @@ #nullable enable -abstract EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase.EnumerableTaskCompletionSources.get -> System.Collections.Generic.IReadOnlyList! -abstract EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase.EnumerableTaskCompletionSources.get -> System.Collections.Generic.IReadOnlyList!>! -EnumerableAsyncProcessor.ActionTaskWrapper -EnumerableAsyncProcessor.ActionTaskWrapper.ActionTaskWrapper() -> void -EnumerableAsyncProcessor.ActionTaskWrapper.ActionTaskWrapper(System.Func! taskFactory, System.Threading.Tasks.TaskCompletionSource! taskCompletionSource) -> void -EnumerableAsyncProcessor.ActionTaskWrapper.Process(System.Threading.CancellationToken cancellationToken) -> System.Threading.Tasks.Task! -EnumerableAsyncProcessor.ActionTaskWrapper -EnumerableAsyncProcessor.ActionTaskWrapper.ActionTaskWrapper() -> void -EnumerableAsyncProcessor.ActionTaskWrapper.ActionTaskWrapper(System.Func!>! taskFactory, System.Threading.Tasks.TaskCompletionSource! taskCompletionSource) -> void -EnumerableAsyncProcessor.ActionTaskWrapper.Process(System.Threading.CancellationToken cancellationToken) -> System.Threading.Tasks.Task! EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ActionAsyncProcessorBuilder(int count, System.Func! taskSelector, System.Threading.CancellationToken cancellationToken) -> void EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ActionAsyncProcessorBuilder(int count, System.Func! taskSelector, System.Threading.CancellationToken cancellationToken) -> void EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInBatches(int batchSize) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! +EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInParallel(int maxConcurrency) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInParallel(int maxConcurrency, System.TimeSpan timeSpan) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInParallel(int permitsPerWindow, System.TimeSpan window, int maxConcurrency) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInParallel(int? maxConcurrency = null, bool scheduleOnThreadPool = false) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessOneAtATime() -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInBatches(int batchSize) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! +EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInParallel(int maxConcurrency) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInParallel(int maxConcurrency, System.TimeSpan timeSpan) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInParallel(int permitsPerWindow, System.TimeSpan window, int maxConcurrency) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessInParallel(int? maxConcurrency = null, bool scheduleOnThreadPool = false) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! @@ -26,15 +18,15 @@ EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder.ProcessOn EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.AsyncEnumerableActionAsyncProcessorBuilder(System.Collections.Generic.IAsyncEnumerable! items, System.Func!>! taskSelector, System.Threading.CancellationToken cancellationToken) -> void EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.AsyncEnumerableActionAsyncProcessorBuilder(System.Collections.Generic.IAsyncEnumerable! items, System.Func!>! taskSelector, System.Threading.CancellationToken cancellationToken) -> void -EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessInBatches(int batchSize) -> EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor! -EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessInParallel(int? maxConcurrency = null, bool scheduleOnThreadPool = false) -> EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor! -EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessOneAtATime() -> EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor! +EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessInBatches(int batchSize) -> EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor! +EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessInParallel(int? maxConcurrency = null, bool scheduleOnThreadPool = false) -> EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor! +EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessOneAtATime() -> EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor! EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.AsyncEnumerableActionAsyncProcessorBuilder(System.Collections.Generic.IAsyncEnumerable! items, System.Func! taskSelector, System.Threading.CancellationToken cancellationToken) -> void EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.AsyncEnumerableActionAsyncProcessorBuilder(System.Collections.Generic.IAsyncEnumerable! items, System.Func! taskSelector, System.Threading.CancellationToken cancellationToken) -> void -EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessInBatches(int batchSize) -> EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor! -EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessInParallel(int? maxConcurrency = null, bool scheduleOnThreadPool = false) -> EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor! -EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessOneAtATime() -> EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor! +EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessInBatches(int batchSize) -> EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor! +EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessInParallel(int? maxConcurrency = null, bool scheduleOnThreadPool = false) -> EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor! +EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder.ProcessOneAtATime() -> EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor! EnumerableAsyncProcessor.Builders.AsyncEnumerableAsyncProcessorBuilder EnumerableAsyncProcessor.Builders.AsyncEnumerableAsyncProcessorBuilder.ForEachAsync(System.Func! taskSelector) -> EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder! EnumerableAsyncProcessor.Builders.AsyncEnumerableAsyncProcessorBuilder.ForEachAsync(System.Func! taskSelector, System.Threading.CancellationToken cancellationToken) -> EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder! @@ -56,6 +48,7 @@ EnumerableAsyncProcessor.Builders.ExecutionCountAsyncProcessorBuilder.SelectAsyn EnumerableAsyncProcessor.Builders.ExecutionCountAsyncProcessorBuilder.SelectAsync(System.Func!>! taskSelector, System.Threading.CancellationToken cancellationToken) -> EnumerableAsyncProcessor.Builders.ActionAsyncProcessorBuilder! EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInBatches(int batchSize) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! +EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInParallel(int maxConcurrency) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInParallel(int maxConcurrency, System.TimeSpan timeSpan) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInParallel(int permitsPerWindow, System.TimeSpan window, int maxConcurrency) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInParallel(int? maxConcurrency = null, bool scheduleOnThreadPool = false) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! @@ -64,6 +57,7 @@ EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ItemActionAsyncProcessorBuilder(System.Collections.Generic.IEnumerable! items, System.Func! taskSelector, System.Threading.CancellationToken cancellationToken) -> void EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ItemActionAsyncProcessorBuilder(System.Collections.Generic.IEnumerable! items, System.Func! taskSelector, System.Threading.CancellationToken cancellationToken) -> void EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInBatches(int batchSize) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! +EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInParallel(int maxConcurrency) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInParallel(int maxConcurrency, System.TimeSpan timeSpan) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInParallel(int permitsPerWindow, System.TimeSpan window, int maxConcurrency) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder.ProcessInParallel(int? maxConcurrency = null, bool scheduleOnThreadPool = false) -> EnumerableAsyncProcessor.Interfaces.IAsyncProcessor! @@ -79,11 +73,11 @@ EnumerableAsyncProcessor.Builders.ItemAsyncProcessorBuilder.SelectAsync< EnumerableAsyncProcessor.Builders.ItemAsyncProcessorBuilder.SelectAsync(System.Func!>! taskSelector, System.Threading.CancellationToken cancellationToken) -> EnumerableAsyncProcessor.Builders.ItemActionAsyncProcessorBuilder! EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions EnumerableAsyncProcessor.Extensions.EnumerableExtensions -EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor -EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor.ExecuteAsync() -> System.Threading.Tasks.Task! -EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor -EnumerableAsyncProcessor.Extensions.IAsyncEnumerableProcessor.ExecuteAsync() -> System.Collections.Generic.IAsyncEnumerable! EnumerableAsyncProcessor.Extensions.ParallelExtensions +EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor +EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor.ExecuteAsync() -> System.Threading.Tasks.Task! +EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor +EnumerableAsyncProcessor.Interfaces.IAsyncEnumerableProcessor.ExecuteAsync() -> System.Collections.Generic.IAsyncEnumerable! EnumerableAsyncProcessor.Interfaces.IAsyncProcessor EnumerableAsyncProcessor.Interfaces.IAsyncProcessor.CancelAll() -> void EnumerableAsyncProcessor.Interfaces.IAsyncProcessor.GetAwaiter() -> System.Runtime.CompilerServices.TaskAwaiter @@ -95,20 +89,9 @@ EnumerableAsyncProcessor.Interfaces.IAsyncProcessor.GetAwaiter() -> Sys EnumerableAsyncProcessor.Interfaces.IAsyncProcessor.GetEnumerableTasks() -> System.Collections.Generic.IEnumerable!>! EnumerableAsyncProcessor.Interfaces.IAsyncProcessor.GetResultsAsync() -> System.Threading.Tasks.Task! EnumerableAsyncProcessor.Interfaces.IAsyncProcessor.GetResultsAsyncEnumerable() -> System.Collections.Generic.IAsyncEnumerable! -EnumerableAsyncProcessor.ItemTaskWrapper -EnumerableAsyncProcessor.ItemTaskWrapper.ItemTaskWrapper() -> void -EnumerableAsyncProcessor.ItemTaskWrapper.ItemTaskWrapper(TInput input, System.Func!>! taskFactory, System.Threading.Tasks.TaskCompletionSource! taskCompletionSource) -> void -EnumerableAsyncProcessor.ItemTaskWrapper.Process(System.Threading.CancellationToken cancellationToken) -> System.Threading.Tasks.Task! -EnumerableAsyncProcessor.ItemTaskWrapper -EnumerableAsyncProcessor.ItemTaskWrapper.ItemTaskWrapper() -> void -EnumerableAsyncProcessor.ItemTaskWrapper.ItemTaskWrapper(TInput input, System.Func! taskFactory, System.Threading.Tasks.TaskCompletionSource! taskCompletionSource) -> void -EnumerableAsyncProcessor.ItemTaskWrapper.Process(System.Threading.CancellationToken cancellationToken) -> System.Threading.Tasks.Task! EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessor -EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessor.AbstractAsyncProcessor(int count, System.Func! taskSelector, System.Threading.CancellationTokenSource! cancellationTokenSource) -> void EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessor -EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessor.AbstractAsyncProcessor(System.Collections.Generic.IEnumerable! items, System.Func! taskSelector, System.Threading.CancellationTokenSource! cancellationTokenSource) -> void EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase -EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase.AbstractAsyncProcessorBase(System.Threading.CancellationTokenSource! cancellationTokenSource) -> void EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase.CancelAll() -> void EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase.Dispose() -> void EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase.DisposeAsync() -> System.Threading.Tasks.ValueTask @@ -116,16 +99,28 @@ EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase. EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase.GetEnumerableTasks() -> System.Collections.Generic.IEnumerable! EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase.WaitAsync() -> System.Threading.Tasks.Task! EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableBatchProcessor +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableBatchProcessor.Dispose() -> void +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableBatchProcessor.DisposeAsync() -> System.Threading.Tasks.ValueTask EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableBatchProcessor.ExecuteAsync() -> System.Threading.Tasks.Task! EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableOneAtATimeProcessor +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableOneAtATimeProcessor.Dispose() -> void +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableOneAtATimeProcessor.DisposeAsync() -> System.Threading.Tasks.ValueTask EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableOneAtATimeProcessor.ExecuteAsync() -> System.Threading.Tasks.Task! EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableParallelProcessor +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableParallelProcessor.Dispose() -> void +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableParallelProcessor.DisposeAsync() -> System.Threading.Tasks.ValueTask EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.AsyncEnumerableParallelProcessor.ExecuteAsync() -> System.Threading.Tasks.Task! EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableBatchProcessor +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableBatchProcessor.Dispose() -> void +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableBatchProcessor.DisposeAsync() -> System.Threading.Tasks.ValueTask EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableBatchProcessor.ExecuteAsync() -> System.Collections.Generic.IAsyncEnumerable! EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableOneAtATimeProcessor +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableOneAtATimeProcessor.Dispose() -> void +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableOneAtATimeProcessor.DisposeAsync() -> System.Threading.Tasks.ValueTask EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableOneAtATimeProcessor.ExecuteAsync() -> System.Collections.Generic.IAsyncEnumerable! EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableParallelProcessor +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableParallelProcessor.Dispose() -> void +EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableParallelProcessor.DisposeAsync() -> System.Threading.Tasks.ValueTask EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors.ResultAsyncEnumerableParallelProcessor.ExecuteAsync() -> System.Collections.Generic.IAsyncEnumerable! EnumerableAsyncProcessor.RunnableProcessors.BatchAsyncProcessor EnumerableAsyncProcessor.RunnableProcessors.BatchAsyncProcessor @@ -134,9 +129,7 @@ EnumerableAsyncProcessor.RunnableProcessors.OneAtATimeAsyncProcessor EnumerableAsyncProcessor.RunnableProcessors.ParallelAsyncProcessor EnumerableAsyncProcessor.RunnableProcessors.ParallelAsyncProcessor EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessor -EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessor.ResultAbstractAsyncProcessor(System.Collections.Generic.IEnumerable! items, System.Func!>! taskSelector, System.Threading.CancellationTokenSource! cancellationTokenSource) -> void EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessor -EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessor.ResultAbstractAsyncProcessor(int count, System.Func!>! taskSelector, System.Threading.CancellationTokenSource! cancellationTokenSource) -> void EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase.CancelAll() -> void EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase.Dispose() -> void @@ -145,7 +138,6 @@ EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbst EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase.GetEnumerableTasks() -> System.Collections.Generic.IEnumerable!>! EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase.GetResultsAsync() -> System.Threading.Tasks.Task! EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase.GetResultsAsyncEnumerable() -> System.Collections.Generic.IAsyncEnumerable! -EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase.ResultAbstractAsyncProcessorBase(System.Threading.CancellationTokenSource! cancellationTokenSource) -> void EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.ResultBatchAsyncProcessor EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.ResultBatchAsyncProcessor EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.ResultOneAtATimeAsyncProcessor @@ -156,35 +148,13 @@ EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.ResultTimedRateLimi EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.ResultTimedRateLimitedParallelAsyncProcessor EnumerableAsyncProcessor.RunnableProcessors.TimedRateLimitedParallelAsyncProcessor EnumerableAsyncProcessor.RunnableProcessors.TimedRateLimitedParallelAsyncProcessor -override EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessor.EnumerableTaskCompletionSources.get -> System.Collections.Generic.IReadOnlyList! -override EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessor.EnumerableTaskCompletionSources.get -> System.Collections.Generic.IReadOnlyList! -override EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessor.EnumerableTaskCompletionSources.get -> System.Collections.Generic.IReadOnlyList!>! -override EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessor.EnumerableTaskCompletionSources.get -> System.Collections.Generic.IReadOnlyList!>! -readonly EnumerableAsyncProcessor.ActionTaskWrapper.TaskCompletionSource -> System.Threading.Tasks.TaskCompletionSource! -readonly EnumerableAsyncProcessor.ActionTaskWrapper.TaskFactory -> System.Func! -readonly EnumerableAsyncProcessor.ActionTaskWrapper.TaskCompletionSource -> System.Threading.Tasks.TaskCompletionSource! -readonly EnumerableAsyncProcessor.ActionTaskWrapper.TaskFactory -> System.Func!>! -readonly EnumerableAsyncProcessor.ItemTaskWrapper.Input -> TInput -readonly EnumerableAsyncProcessor.ItemTaskWrapper.TaskCompletionSource -> System.Threading.Tasks.TaskCompletionSource! -readonly EnumerableAsyncProcessor.ItemTaskWrapper.TaskFactory -> System.Func!>! -readonly EnumerableAsyncProcessor.ItemTaskWrapper.Input -> TInput -readonly EnumerableAsyncProcessor.ItemTaskWrapper.TaskCompletionSource -> System.Threading.Tasks.TaskCompletionSource! -readonly EnumerableAsyncProcessor.ItemTaskWrapper.TaskFactory -> System.Func! -readonly EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessor.TaskWrappers -> EnumerableAsyncProcessor.ActionTaskWrapper[]! -readonly EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessor.TaskWrappers -> EnumerableAsyncProcessor.ItemTaskWrapper[]! -readonly EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase.CancellationToken -> System.Threading.CancellationToken -readonly EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessor.TaskWrappers -> EnumerableAsyncProcessor.ItemTaskWrapper[]! -readonly EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessor.TaskWrappers -> EnumerableAsyncProcessor.ActionTaskWrapper[]! -readonly EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase.CancellationToken -> System.Threading.CancellationToken static EnumerableAsyncProcessor.Builders.AsyncProcessorBuilder.WithExecutionCount(int count) -> EnumerableAsyncProcessor.Builders.ExecutionCountAsyncProcessorBuilder! static EnumerableAsyncProcessor.Builders.AsyncProcessorBuilder.WithItems(System.Collections.Generic.IEnumerable! items) -> EnumerableAsyncProcessor.Builders.ItemAsyncProcessorBuilder! static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ForEachAsync(this System.Collections.Generic.IAsyncEnumerable! items, System.Func! taskSelector, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder! static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ForEachAsync(this System.Collections.Generic.IAsyncEnumerable! items, System.Func! taskSelector, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder! +static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ProcessInParallel(this System.Collections.Generic.IAsyncEnumerable! items, System.Func!>! taskSelector, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> System.Threading.Tasks.Task!>! static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ProcessInParallel(this System.Collections.Generic.IAsyncEnumerable! items, System.Func!>! taskSelector, int maxConcurrency, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> System.Threading.Tasks.Task!>! static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ProcessInParallel(this System.Collections.Generic.IAsyncEnumerable! items, System.Func!>! taskSelector, int? maxConcurrency, bool scheduleOnThreadPool, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> System.Threading.Tasks.Task!>! -static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ProcessInParallel(this System.Collections.Generic.IAsyncEnumerable! items, System.Func!>! taskSelector, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> System.Threading.Tasks.Task!>! -static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ProcessInParallel(this System.Collections.Generic.IAsyncEnumerable! items, int maxConcurrency, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> System.Threading.Tasks.Task!>! -static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ProcessInParallel(this System.Collections.Generic.IAsyncEnumerable! items, int? maxConcurrency, bool scheduleOnThreadPool, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> System.Threading.Tasks.Task!>! static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ProcessInParallel(this System.Collections.Generic.IAsyncEnumerable! items, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> System.Threading.Tasks.Task!>! static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.SelectAsync(this System.Collections.Generic.IAsyncEnumerable! items, System.Func!>! taskSelector, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder! static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.SelectAsync(this System.Collections.Generic.IAsyncEnumerable! items, System.Func!>! taskSelector, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> EnumerableAsyncProcessor.Builders.AsyncEnumerableActionAsyncProcessorBuilder! @@ -210,5 +180,3 @@ static EnumerableAsyncProcessor.Extensions.ParallelExtensions.InParallelAsync(this System.Collections.Generic.IEnumerable! source, int levelOfParallelism, System.Action! taskSelector, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> System.Threading.Tasks.Task! static EnumerableAsyncProcessor.Extensions.ParallelExtensions.InParallelAsync(this System.Collections.Generic.IEnumerable! source, int levelOfParallelism, System.Func! taskSelector) -> System.Threading.Tasks.Task! static EnumerableAsyncProcessor.Extensions.ParallelExtensions.InParallelAsync(this System.Collections.Generic.IEnumerable! source, int levelOfParallelism, System.Func! taskSelector, System.Threading.CancellationToken cancellationToken) -> System.Threading.Tasks.Task! -virtual EnumerableAsyncProcessor.RunnableProcessors.Abstract.AbstractAsyncProcessorBase.DisposeAsyncCore() -> System.Threading.Tasks.ValueTask -virtual EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract.ResultAbstractAsyncProcessorBase.DisposeAsyncCore() -> System.Threading.Tasks.ValueTask diff --git a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor.cs index 88ef7b1..d894994 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor.cs @@ -4,13 +4,13 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.Abstract; public abstract class AbstractAsyncProcessor : AbstractAsyncProcessorBase { - protected readonly ActionTaskWrapper[] TaskWrappers; + private protected readonly ActionTaskWrapper[] TaskWrappers; private readonly TaskCompletionSource[] _taskCompletionSources; - protected override IReadOnlyList EnumerableTaskCompletionSources => _taskCompletionSources; + private protected override IReadOnlyList EnumerableTaskCompletionSources => _taskCompletionSources; - protected AbstractAsyncProcessor(int count, Func taskSelector, CancellationTokenSource cancellationTokenSource) : base(cancellationTokenSource) + private protected AbstractAsyncProcessor(int count, Func taskSelector, CancellationTokenSource cancellationTokenSource) : base(cancellationTokenSource) { ValidationHelper.ThrowIfNegative(count); ValidationHelper.ThrowIfNull(taskSelector); diff --git a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs index 87f6d53..ed964f2 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs @@ -5,15 +5,15 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.Abstract; public abstract class AbstractAsyncProcessorBase : IAsyncProcessor, IAsyncDisposable, IDisposable { - protected abstract IReadOnlyList EnumerableTaskCompletionSources { get; } - protected readonly CancellationToken CancellationToken; + private protected abstract IReadOnlyList EnumerableTaskCompletionSources { get; } + private protected readonly CancellationToken CancellationToken; private readonly ProcessorLifecycle _lifecycle; private Task? _overallTask; private Task OverallTask => _overallTask ??= Task.WhenAll(GetEnumerableTasks()); - protected AbstractAsyncProcessorBase(CancellationTokenSource cancellationTokenSource) + private protected AbstractAsyncProcessorBase(CancellationTokenSource cancellationTokenSource) { _lifecycle = new ProcessorLifecycle(cancellationTokenSource, TrySetCanceledAll, TrySetExceptionAll); CancellationToken = _lifecycle.Token; @@ -69,7 +69,7 @@ public async ValueTask DisposeAsync() GC.SuppressFinalize(this); } - protected virtual ValueTask DisposeAsyncCore() + private protected virtual ValueTask DisposeAsyncCore() { return ValueTask.CompletedTask; } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor_1.cs index 6b493b7..6cb7827 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor_1.cs @@ -4,13 +4,13 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.Abstract; public abstract class AbstractAsyncProcessor : AbstractAsyncProcessorBase { - protected readonly ItemTaskWrapper[] TaskWrappers; + private protected readonly ItemTaskWrapper[] TaskWrappers; private readonly TaskCompletionSource[] _taskCompletionSources; - protected override IReadOnlyList EnumerableTaskCompletionSources => _taskCompletionSources; + private protected override IReadOnlyList EnumerableTaskCompletionSources => _taskCompletionSources; - protected AbstractAsyncProcessor(IEnumerable items, Func taskSelector, CancellationTokenSource cancellationTokenSource) : base(cancellationTokenSource) + private protected AbstractAsyncProcessor(IEnumerable items, Func taskSelector, CancellationTokenSource cancellationTokenSource) : base(cancellationTokenSource) { ValidationHelper.ThrowIfNull(items); ValidationHelper.ThrowIfNull(taskSelector); diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs index 5b59347..bd24f1b 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs @@ -1,8 +1,8 @@ -using EnumerableAsyncProcessor.Extensions; +using EnumerableAsyncProcessor.Interfaces; namespace EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable; -public class AsyncEnumerableBatchProcessor : IAsyncEnumerableProcessor +public sealed class AsyncEnumerableBatchProcessor : IAsyncEnumerableProcessor { private readonly IAsyncEnumerable _items; private readonly Func _taskSelector; @@ -24,29 +24,48 @@ internal AsyncEnumerableBatchProcessor( public async Task ExecuteAsync() { var cancellationToken = _cancellationTokenSource.Token; - var batch = new List(_batchSize); - - await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) + + try { - batch.Add(item); - - if (batch.Count >= _batchSize) + var batch = new List(_batchSize); + + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - await ProcessBatch(batch, cancellationToken).ConfigureAwait(false); - batch = new List(_batchSize); + batch.Add(item); + + if (batch.Count >= _batchSize) + { + await ProcessBatch(batch).ConfigureAwait(false); + batch = new List(_batchSize); + } + } + + // Process any remaining items in the final batch + if (batch.Count > 0) + { + await ProcessBatch(batch).ConfigureAwait(false); } } - - // Process any remaining items in the final batch - if (batch.Count > 0) + finally { - await ProcessBatch(batch, cancellationToken).ConfigureAwait(false); + _cancellationTokenSource.Dispose(); } } - - private async Task ProcessBatch(List batch, CancellationToken cancellationToken) + + private async Task ProcessBatch(List batch) { var tasks = batch.Select(item => _taskSelector(item)).ToArray(); await Task.WhenAll(tasks).ConfigureAwait(false); } -} \ No newline at end of file + + public void Dispose() + { + _cancellationTokenSource.Dispose(); + } + + public ValueTask DisposeAsync() + { + Dispose(); + return ValueTask.CompletedTask; + } +} diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs index 107d472..f5fb486 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs @@ -1,11 +1,11 @@ -using EnumerableAsyncProcessor.Extensions; +using EnumerableAsyncProcessor.Interfaces; namespace EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable; /// /// Sequential processor that processes items one at a time from an IAsyncEnumerable. /// -public class AsyncEnumerableOneAtATimeProcessor : IAsyncEnumerableProcessor +public sealed class AsyncEnumerableOneAtATimeProcessor : IAsyncEnumerableProcessor { private readonly IAsyncEnumerable _items; private readonly Func _taskSelector; @@ -25,9 +25,27 @@ public async Task ExecuteAsync() { var cancellationToken = _cancellationTokenSource.Token; - await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) + try { - await _taskSelector(item).ConfigureAwait(false); + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) + { + await _taskSelector(item).ConfigureAwait(false); + } } + finally + { + _cancellationTokenSource.Dispose(); + } + } + + public void Dispose() + { + _cancellationTokenSource.Dispose(); + } + + public ValueTask DisposeAsync() + { + Dispose(); + return ValueTask.CompletedTask; } -} \ No newline at end of file +} diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs index 5e395be..04dd3f1 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs @@ -1,8 +1,8 @@ -using EnumerableAsyncProcessor.Extensions; +using EnumerableAsyncProcessor.Interfaces; namespace EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable; -public class AsyncEnumerableParallelProcessor : IAsyncEnumerableProcessor +public sealed class AsyncEnumerableParallelProcessor : IAsyncEnumerableProcessor { private readonly IAsyncEnumerable _items; private readonly Func _taskSelector; @@ -27,19 +27,24 @@ internal AsyncEnumerableParallelProcessor( public async Task ExecuteAsync() { var cancellationToken = _cancellationTokenSource.Token; - - if (_maxConcurrency.HasValue) - { - await AsyncEnumerableWorkerPool.ProcessAsync( - _items, - _taskSelector, - _maxConcurrency.Value, - cancellationToken).ConfigureAwait(false); - return; - } - else + try { + if (_maxConcurrency.HasValue) + { + Func taskSelector = _scheduleOnThreadPool + ? item => Task.Run(() => _taskSelector(item), cancellationToken) + : _taskSelector; + + await AsyncEnumerableWorkerPool.ProcessAsync( + _items, + taskSelector, + _maxConcurrency.Value, + cancellationToken).ConfigureAwait(false); + + return; + } + // Unbounded parallel processing var tasks = new List(); @@ -70,5 +75,20 @@ await AsyncEnumerableWorkerPool.ProcessAsync( } } } + finally + { + _cancellationTokenSource.Dispose(); + } + } + + public void Dispose() + { + _cancellationTokenSource.Dispose(); + } + + public ValueTask DisposeAsync() + { + Dispose(); + return ValueTask.CompletedTask; } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs index 1af6962..b5e889e 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs @@ -1,8 +1,9 @@ -using EnumerableAsyncProcessor.Extensions; +using System.Runtime.CompilerServices; +using EnumerableAsyncProcessor.Interfaces; namespace EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors; -public class ResultAsyncEnumerableBatchProcessor : IAsyncEnumerableProcessor +public sealed class ResultAsyncEnumerableBatchProcessor : IAsyncEnumerableProcessor { private readonly IAsyncEnumerable _items; private readonly Func> _taskSelector; @@ -24,40 +25,59 @@ internal ResultAsyncEnumerableBatchProcessor( public async IAsyncEnumerable ExecuteAsync() { var cancellationToken = _cancellationTokenSource.Token; - var batch = new List(_batchSize); - - await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) + + try { - batch.Add(item); - - if (batch.Count >= _batchSize) + var batch = new List(_batchSize); + + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) + { + batch.Add(item); + + if (batch.Count >= _batchSize) + { + await foreach (var result in ProcessBatch(batch, cancellationToken).ConfigureAwait(false)) + { + yield return result; + } + batch = new List(_batchSize); + } + } + + // Process any remaining items in the final batch + if (batch.Count > 0) { await foreach (var result in ProcessBatch(batch, cancellationToken).ConfigureAwait(false)) { yield return result; } - batch = new List(_batchSize); } } - - // Process any remaining items in the final batch - if (batch.Count > 0) + finally { - await foreach (var result in ProcessBatch(batch, cancellationToken).ConfigureAwait(false)) - { - yield return result; - } + _cancellationTokenSource.Dispose(); } } - - private async IAsyncEnumerable ProcessBatch(List batch, [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken) + + private async IAsyncEnumerable ProcessBatch(List batch, [EnumeratorCancellation] CancellationToken cancellationToken) { var tasks = batch.Select(item => _taskSelector(item)).ToArray(); var results = await Task.WhenAll(tasks).ConfigureAwait(false); - + foreach (var result in results) { yield return result; } } -} \ No newline at end of file + + public void Dispose() + { + _cancellationTokenSource.Dispose(); + } + + public ValueTask DisposeAsync() + { + Dispose(); + return ValueTask.CompletedTask; + } +} diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs index 7af06b5..7b49e28 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs @@ -1,11 +1,11 @@ -using EnumerableAsyncProcessor.Extensions; +using EnumerableAsyncProcessor.Interfaces; namespace EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors; /// /// Sequential processor that processes items one at a time and returns results in order. /// -public class ResultAsyncEnumerableOneAtATimeProcessor : IAsyncEnumerableProcessor +public sealed class ResultAsyncEnumerableOneAtATimeProcessor : IAsyncEnumerableProcessor { private readonly IAsyncEnumerable _items; private readonly Func> _taskSelector; @@ -25,10 +25,28 @@ public async IAsyncEnumerable ExecuteAsync() { var cancellationToken = _cancellationTokenSource.Token; - await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) + try { - var result = await _taskSelector(item).ConfigureAwait(false); - yield return result; + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) + { + var result = await _taskSelector(item).ConfigureAwait(false); + yield return result; + } } + finally + { + _cancellationTokenSource.Dispose(); + } + } + + public void Dispose() + { + _cancellationTokenSource.Dispose(); + } + + public ValueTask DisposeAsync() + { + Dispose(); + return ValueTask.CompletedTask; } -} \ No newline at end of file +} diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs index 0bbba83..c98af3c 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs @@ -1,8 +1,9 @@ using EnumerableAsyncProcessor.Extensions; +using EnumerableAsyncProcessor.Interfaces; namespace EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable.ResultProcessors; -public class ResultAsyncEnumerableParallelProcessor : IAsyncEnumerableProcessor +public sealed class ResultAsyncEnumerableParallelProcessor : IAsyncEnumerableProcessor { private readonly IAsyncEnumerable _items; private readonly Func> _taskSelector; @@ -27,61 +28,84 @@ internal ResultAsyncEnumerableParallelProcessor( public async IAsyncEnumerable ExecuteAsync() { var cancellationToken = _cancellationTokenSource.Token; - if (_maxConcurrency.HasValue) - { - await foreach (var result in AsyncEnumerableWorkerPool.ProcessResultsAsync( - _items, - _taskSelector, - _maxConcurrency.Value, - cancellationToken).ConfigureAwait(false)) - { - yield return result; - } - - yield break; - } - - var tasks = new List>(); - // Unbounded parallel processing try { - await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) + if (_maxConcurrency.HasValue) { - var capturedItem = item; + Func> taskSelector = _scheduleOnThreadPool + ? item => Task.Run(() => _taskSelector(item), cancellationToken) + : _taskSelector; - Task task; - if (_scheduleOnThreadPool) - { - task = Task.Run(() => _taskSelector(capturedItem), cancellationToken); - } - else + await foreach (var result in AsyncEnumerableWorkerPool.ProcessResultsAsync( + _items, + taskSelector, + _maxConcurrency.Value, + cancellationToken).ConfigureAwait(false)) { - task = _taskSelector(capturedItem); + yield return result; } - tasks.Add(task); + yield break; } - // Yield all results as they complete - await foreach (var result in tasks.ToIAsyncEnumerable(cancellationToken).ConfigureAwait(false)) - { - yield return result; - } - } - finally - { - if (tasks.Count > 0) + var tasks = new List>(); + + // Unbounded parallel processing + try { - try + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) + { + var capturedItem = item; + + Task task; + if (_scheduleOnThreadPool) + { + task = Task.Run(() => _taskSelector(capturedItem), cancellationToken); + } + else + { + task = _taskSelector(capturedItem); + } + + tasks.Add(task); + } + + // Yield all results as they complete + await foreach (var result in tasks.ToIAsyncEnumerable(cancellationToken).ConfigureAwait(false)) { - await Task.WhenAll(tasks).ConfigureAwait(false); + yield return result; } - catch + } + finally + { + if (tasks.Count > 0) { - // Preserve the exception already propagating from enumeration or result consumption. + try + { + await Task.WhenAll(tasks).ConfigureAwait(false); + } + catch + { + // Preserve the exception already propagating from enumeration or result consumption. + } } } } + finally + { + _cancellationTokenSource.Dispose(); + } + } + + public void Dispose() + { + _cancellationTokenSource.Dispose(); + } + + public ValueTask DisposeAsync() + { + Dispose(); + return ValueTask.CompletedTask; } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor.cs index 7fc67f1..6e3d1cf 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; -public class BatchAsyncProcessor : AbstractAsyncProcessor +public sealed class BatchAsyncProcessor : AbstractAsyncProcessor { private readonly int _batchSize; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor_1.cs index b8e91ec..f31ba22 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor_1.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; -public class BatchAsyncProcessor : AbstractAsyncProcessor +public sealed class BatchAsyncProcessor : AbstractAsyncProcessor { private readonly int _batchSize; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor.cs index ee9e816..88f3f71 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor.cs @@ -2,7 +2,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; -public class OneAtATimeAsyncProcessor : AbstractAsyncProcessor +public sealed class OneAtATimeAsyncProcessor : AbstractAsyncProcessor { internal OneAtATimeAsyncProcessor(int count, Func taskSelector, CancellationTokenSource cancellationTokenSource) : base(count, taskSelector, cancellationTokenSource) { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor_1.cs index 4fcb885..3ff8917 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor_1.cs @@ -2,7 +2,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; -public class OneAtATimeAsyncProcessor : AbstractAsyncProcessor +public sealed class OneAtATimeAsyncProcessor : AbstractAsyncProcessor { internal OneAtATimeAsyncProcessor(IEnumerable items, Func taskSelector, CancellationTokenSource cancellationTokenSource) : base(items, taskSelector, cancellationTokenSource) { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs index 8e7f6a2..93e40ae 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; -public class ParallelAsyncProcessor : AbstractAsyncProcessor +public sealed class ParallelAsyncProcessor : AbstractAsyncProcessor { private readonly int? _maxConcurrency; private readonly bool _scheduleOnThreadPool; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs index fd7a364..e06ed0b 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; -public class ParallelAsyncProcessor : AbstractAsyncProcessor +public sealed class ParallelAsyncProcessor : AbstractAsyncProcessor { private readonly int? _maxConcurrency; private readonly bool _scheduleOnThreadPool; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs index 0f21396..8b3959d 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs @@ -6,15 +6,15 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract; public abstract class ResultAbstractAsyncProcessorBase : IAsyncProcessor, IAsyncDisposable, IDisposable { - protected abstract IReadOnlyList> EnumerableTaskCompletionSources { get; } - protected readonly CancellationToken CancellationToken; + private protected abstract IReadOnlyList> EnumerableTaskCompletionSources { get; } + private protected readonly CancellationToken CancellationToken; private readonly ProcessorLifecycle _lifecycle; private Task? _results; private Task Results => _results ??= Task.WhenAll(GetEnumerableTasks()); - protected ResultAbstractAsyncProcessorBase(CancellationTokenSource cancellationTokenSource) + private protected ResultAbstractAsyncProcessorBase(CancellationTokenSource cancellationTokenSource) { _lifecycle = new ProcessorLifecycle(cancellationTokenSource, TrySetCanceledAll, TrySetExceptionAll); CancellationToken = _lifecycle.Token; @@ -75,7 +75,7 @@ public async ValueTask DisposeAsync() GC.SuppressFinalize(this); } - protected virtual ValueTask DisposeAsyncCore() + private protected virtual ValueTask DisposeAsyncCore() { return ValueTask.CompletedTask; } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_1.cs index 85ff76b..8605bbc 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_1.cs @@ -4,13 +4,13 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract; public abstract class ResultAbstractAsyncProcessor : ResultAbstractAsyncProcessorBase { - protected readonly ActionTaskWrapper[] TaskWrappers; + private protected readonly ActionTaskWrapper[] TaskWrappers; private readonly TaskCompletionSource[] _taskCompletionSources; - protected override IReadOnlyList> EnumerableTaskCompletionSources => _taskCompletionSources; + private protected override IReadOnlyList> EnumerableTaskCompletionSources => _taskCompletionSources; - protected ResultAbstractAsyncProcessor(int count, Func> taskSelector, CancellationTokenSource cancellationTokenSource) : base(cancellationTokenSource) + private protected ResultAbstractAsyncProcessor(int count, Func> taskSelector, CancellationTokenSource cancellationTokenSource) : base(cancellationTokenSource) { ValidationHelper.ThrowIfNegative(count); ValidationHelper.ThrowIfNull(taskSelector); diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_2.cs index 615b5f9..32a871f 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_2.cs @@ -4,13 +4,13 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract; public abstract class ResultAbstractAsyncProcessor : ResultAbstractAsyncProcessorBase { - protected readonly ItemTaskWrapper[] TaskWrappers; + private protected readonly ItemTaskWrapper[] TaskWrappers; private readonly TaskCompletionSource[] _taskCompletionSources; - protected override IReadOnlyList> EnumerableTaskCompletionSources => _taskCompletionSources; + private protected override IReadOnlyList> EnumerableTaskCompletionSources => _taskCompletionSources; - protected ResultAbstractAsyncProcessor(IEnumerable items, Func> taskSelector, CancellationTokenSource cancellationTokenSource) : base(cancellationTokenSource) + private protected ResultAbstractAsyncProcessor(IEnumerable items, Func> taskSelector, CancellationTokenSource cancellationTokenSource) : base(cancellationTokenSource) { ValidationHelper.ThrowIfNull(items); ValidationHelper.ThrowIfNull(taskSelector); diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_1.cs index 639da5f..40ba585 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_1.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; -public class ResultBatchAsyncProcessor : ResultAbstractAsyncProcessor +public sealed class ResultBatchAsyncProcessor : ResultAbstractAsyncProcessor { private readonly int _batchSize; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_2.cs index 7ae1d64..6297280 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_2.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; -public class ResultBatchAsyncProcessor : ResultAbstractAsyncProcessor +public sealed class ResultBatchAsyncProcessor : ResultAbstractAsyncProcessor { private readonly int _batchSize; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_1.cs index 39e358d..d66fee0 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_1.cs @@ -2,7 +2,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; -public class ResultOneAtATimeAsyncProcessor : ResultAbstractAsyncProcessor +public sealed class ResultOneAtATimeAsyncProcessor : ResultAbstractAsyncProcessor { internal ResultOneAtATimeAsyncProcessor(int count, Func> taskSelector, CancellationTokenSource cancellationTokenSource) : base(count, taskSelector, cancellationTokenSource) { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_2.cs index 5d90c02..ed40022 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_2.cs @@ -2,7 +2,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; -public class ResultOneAtATimeAsyncProcessor : ResultAbstractAsyncProcessor +public sealed class ResultOneAtATimeAsyncProcessor : ResultAbstractAsyncProcessor { internal ResultOneAtATimeAsyncProcessor(IEnumerable items, Func> taskSelector, CancellationTokenSource cancellationTokenSource) : base(items, taskSelector, cancellationTokenSource) { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs index a25bb9d..12c55a0 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; -public class ResultParallelAsyncProcessor : ResultAbstractAsyncProcessor +public sealed class ResultParallelAsyncProcessor : ResultAbstractAsyncProcessor { private readonly int? _maxConcurrency; private readonly bool _scheduleOnThreadPool; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs index aa7f7e0..2e16e73 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; -public class ResultParallelAsyncProcessor : ResultAbstractAsyncProcessor +public sealed class ResultParallelAsyncProcessor : ResultAbstractAsyncProcessor { private readonly int? _maxConcurrency; private readonly bool _scheduleOnThreadPool; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs index cbb9752..0422c3b 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; -public class ResultTimedRateLimitedParallelAsyncProcessor : ResultAbstractAsyncProcessor +public sealed class ResultTimedRateLimitedParallelAsyncProcessor : ResultAbstractAsyncProcessor { private readonly int _permitsPerWindow; private readonly TimeSpan _window; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs index f680ed9..a2ea234 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; -public class ResultTimedRateLimitedParallelAsyncProcessor : ResultAbstractAsyncProcessor +public sealed class ResultTimedRateLimitedParallelAsyncProcessor : ResultAbstractAsyncProcessor { private readonly int _permitsPerWindow; private readonly TimeSpan _window; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor.cs index 253c47f..d19a1a3 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; -public class TimedRateLimitedParallelAsyncProcessor : AbstractAsyncProcessor +public sealed class TimedRateLimitedParallelAsyncProcessor : AbstractAsyncProcessor { private readonly int _permitsPerWindow; private readonly TimeSpan _window; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs index f9ea993..681d073 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs @@ -3,7 +3,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; -public class TimedRateLimitedParallelAsyncProcessor : AbstractAsyncProcessor +public sealed class TimedRateLimitedParallelAsyncProcessor : AbstractAsyncProcessor { private readonly int _permitsPerWindow; private readonly TimeSpan _window; diff --git a/EnumerableAsyncProcessor/TaskWrapper.cs b/EnumerableAsyncProcessor/TaskWrapper.cs index 0b3acd4..0204386 100644 --- a/EnumerableAsyncProcessor/TaskWrapper.cs +++ b/EnumerableAsyncProcessor/TaskWrapper.cs @@ -53,7 +53,7 @@ internal static void TrySetFromFault(this TaskCompletionSource /// /// A struct wrapper pairing an action task factory with its completion source. /// -public readonly struct ActionTaskWrapper : ITaskWrapper +internal readonly struct ActionTaskWrapper : ITaskWrapper { public readonly Func TaskFactory; public readonly TaskCompletionSource TaskCompletionSource; @@ -89,7 +89,7 @@ public async Task Process(CancellationToken cancellationToken) /// /// A struct wrapper pairing an input item and its task factory with a completion source. /// -public readonly struct ItemTaskWrapper : ITaskWrapper +internal readonly struct ItemTaskWrapper : ITaskWrapper { public readonly TInput Input; public readonly Func TaskFactory; @@ -127,7 +127,7 @@ public async Task Process(CancellationToken cancellationToken) /// /// A struct wrapper pairing an input item and its result-producing task factory with a completion source. /// -public readonly struct ItemTaskWrapper : ITaskWrapper +internal readonly struct ItemTaskWrapper : ITaskWrapper { public readonly TInput Input; public readonly Func> TaskFactory; @@ -164,7 +164,7 @@ public async Task Process(CancellationToken cancellationToken) /// /// A struct wrapper pairing a result-producing task factory with its completion source. /// -public readonly struct ActionTaskWrapper : ITaskWrapper +internal readonly struct ActionTaskWrapper : ITaskWrapper { public readonly Func> TaskFactory; public readonly TaskCompletionSource TaskCompletionSource; diff --git a/README.md b/README.md index 806f640..857f9d0 100644 --- a/README.md +++ b/README.md @@ -25,13 +25,16 @@ Version 4 requires .NET 8 or later. The package targets and tests `net8.0`, `net Version 4 is a major release. Review these source and behaviour changes before upgrading: - **.NET 8 is the minimum runtime.** The `net6.0` and `netstandard2.0` targets, the Polyfill dependency, and the `TaskCompletionSource` compatibility shim were removed. The package targets and tests `net8.0`, `net9.0`, and `net10.0`. -- **Parallel processing has one API.** Use `ProcessInParallel(maxConcurrency: 100)` for bounded concurrency and `ProcessInParallel()` for unbounded concurrency. The old non-timed `RateLimitedParallelAsyncProcessor*` types and duplicate overload family were removed. Rename named `levelOfParallelism` arguments to `maxConcurrency`; replace positional `ProcessInParallel(true)` calls with `ProcessInParallel(scheduleOnThreadPool: true)`. +- **Parallel processing has one API.** Use `ProcessInParallel(maxConcurrency: 100)` for bounded concurrency and `ProcessInParallel()` for unbounded concurrency. The old non-timed `RateLimitedParallelAsyncProcessor*` types and duplicate overload family were removed. Rename named `levelOfParallelism` arguments to `maxConcurrency` on the `ProcessInParallel` builder family (`InParallelAsync` keeps its `levelOfParallelism` parameter name); replace positional `ProcessInParallel(true)` calls with `ProcessInParallel(scheduleOnThreadPool: true)`. +- **The parameterized no-selector `IAsyncEnumerable.ProcessInParallel(maxConcurrency, ...)` overloads were removed.** They only buffered the stream into a list — their `maxConcurrency`/`scheduleOnThreadPool` parameters had no effect. Pass a selector (`items.ProcessInParallel(item => Task.FromResult(item), maxConcurrency)`) instead. The parameterless `ProcessInParallel(cancellationToken)` overload remains for binary compatibility with v3-compiled assemblies (notably TUnit) and is documented as a plain collect. +- **`IAsyncEnumerableProcessor` moved to `EnumerableAsyncProcessor.Interfaces`** (previously `EnumerableAsyncProcessor.Extensions`) and now implements `IAsyncDisposable`/`IDisposable`. Processors returned by the `IAsyncEnumerable` builder path dispose their internal resources automatically when `ExecuteAsync` completes; they are single-use. +- **Processor and builder classes are sealed.** None of them were externally subclassable in practice (their pipelines hinge on internal members); v4 makes that explicit. - **Timed processing is a real start-rate limit.** It now uses a shared token bucket instead of holding each worker slot for at least one window. `ProcessInParallel(permitsPerWindow, window, maxConcurrency)` controls start rate and in-flight concurrency independently. The existing two-argument overload remains and uses its first value for both limits. - **`IEnumerable` input is materialized once when the processor is built.** One-shot enumerables are now supported and side effects run once. Iterator exceptions surface from the terminal builder call, such as `ProcessInParallel(...)`, instead of later from an awaiter. - **Synchronous disposal no longer waits.** `Dispose()` cancels pending work and releases resources without blocking. Use `await DisposeAsync()` or `await using` when shutdown must wait for in-flight work; the async wait is bounded to 30 seconds. - **Arbitrary upper limits were removed.** Task counts and batch sizes may exceed 10,000, and time windows may exceed 24 hours. Validity checks remain: counts, batch sizes, concurrency, and permit counts must be positive; time windows cannot be negative. - **Validation is consistent and eager.** Invalid `maxConcurrency`, parallelism, batch, and timed-rate arguments now throw while the processor is built for both action and result variants. -- **Incidental `TaskWrapper` API was removed.** `IEquatable`, equality operators, `Deconstruct`, and custom `GetHashCode` members were implementation details and are no longer public. Recompile consumers that referenced those members. +- **`TaskWrapper` types and processor plumbing are internal.** The `ActionTaskWrapper`/`ItemTaskWrapper` structs and the previously `protected` members of the abstract processor base classes (`TaskWrappers`, `EnumerableTaskCompletionSources`, `CancellationToken`, constructors) were implementation details that could not be meaningfully used outside the library, and are no longer public. - **Synchronous `InParallelAsync` delegates now run concurrently.** The `Func` and `Action` overloads use thread-pool workers instead of executing delegates sequentially on the caller. Do not depend on their previous thread affinity or execution order. Version 4 also adds cancellation-aware selectors. Use `(item, cancellationToken) => ...` when in-flight work must observe external cancellation, `CancelAll()`, or disposal: @@ -391,14 +394,15 @@ When a processor is disposed: ### Extension Method Disposal -Note that the convenience extension methods like `IAsyncEnumerable.ProcessInParallel()` handle disposal automatically and return the final results directly: +Note that the convenience extension methods like `IAsyncEnumerable.ProcessInParallel(selector)` handle disposal automatically and return the final results directly: ```csharp // These extension methods handle disposal internally -var results = await asyncEnumerable.ProcessInParallel(); var transformedResults = await asyncEnumerable.ProcessInParallel(async x => await TransformAsync(x)); ``` +Processors built from an `IAsyncEnumerable` source (`IAsyncEnumerableProcessor`) also dispose their internal resources automatically when `ExecuteAsync` completes; `await using` is still supported and safe if `ExecuteAsync` is never called. + The disposal guidance above applies when you're working with the processor objects directly (using the builder pattern). ## Quick Reference: Disposal Patterns From 775a6725464f488db8162b5fa56c58b1db5c3903 Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Tue, 21 Jul 2026 22:24:49 +0100 Subject: [PATCH 2/4] refactor: hide binary-compat shims from IntelliSense Mark the no-selector ProcessInParallel collect overload [Obsolete] so source consumers migrate off it, and hide the builder ProcessInParallel(int) forwarders with [EditorBrowsable(Never)]. The forwarders cannot be [Obsolete]: exact-int overload resolution prefers them for the idiomatic ProcessInParallel(5) call, which would warn on every normal bounded use. --- .../Builders/ActionAsyncProcessorBuilder.cs | 1 + .../Builders/ActionAsyncProcessorBuilder_1.cs | 1 + .../Builders/ItemActionAsyncProcessorBuilder_1.cs | 1 + .../Builders/ItemActionAsyncProcessorBuilder_2.cs | 1 + .../Extensions/AsyncEnumerableExtensions.cs | 2 ++ 5 files changed, 6 insertions(+) diff --git a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs index 7396783..47fe659 100644 --- a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs +++ b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs @@ -58,6 +58,7 @@ public IAsyncProcessor ProcessInParallel(int? maxConcurrency = null, bool schedu /// Processes items in parallel with bounded concurrency. Binary-compatible with assemblies /// compiled against v3 (equivalent to ProcessInParallel(maxConcurrency: n)). /// + [System.ComponentModel.EditorBrowsable(System.ComponentModel.EditorBrowsableState.Never)] public IAsyncProcessor ProcessInParallel(int maxConcurrency) { return ProcessInParallel((int?)maxConcurrency); diff --git a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs index 96fc06b..0737a92 100644 --- a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs @@ -58,6 +58,7 @@ public IAsyncProcessor ProcessInParallel(int? maxConcurrency = null, bo /// Processes items in parallel with bounded concurrency. Binary-compatible with assemblies /// compiled against v3 (equivalent to ProcessInParallel(maxConcurrency: n)). /// + [System.ComponentModel.EditorBrowsable(System.ComponentModel.EditorBrowsableState.Never)] public IAsyncProcessor ProcessInParallel(int maxConcurrency) { return ProcessInParallel((int?)maxConcurrency); diff --git a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs index 91d0561..bdb6ebb 100644 --- a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs @@ -61,6 +61,7 @@ public IAsyncProcessor ProcessInParallel(int? maxConcurrency = null, bool schedu /// Processes items in parallel with bounded concurrency. Binary-compatible with assemblies /// compiled against v3 (equivalent to ProcessInParallel(maxConcurrency: n)). /// + [System.ComponentModel.EditorBrowsable(System.ComponentModel.EditorBrowsableState.Never)] public IAsyncProcessor ProcessInParallel(int maxConcurrency) { return ProcessInParallel((int?)maxConcurrency); diff --git a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs index a722c85..d687b80 100644 --- a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs +++ b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs @@ -87,6 +87,7 @@ public IAsyncProcessor ProcessInParallel(int? maxConcurrency = null, bo /// Processes items in parallel with bounded concurrency. Binary-compatible with assemblies /// compiled against v3 (equivalent to ProcessInParallel(maxConcurrency: n)). /// + [System.ComponentModel.EditorBrowsable(System.ComponentModel.EditorBrowsableState.Never)] public IAsyncProcessor ProcessInParallel(int maxConcurrency) { return ProcessInParallel((int?)maxConcurrency); diff --git a/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs b/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs index 0459eaa..82e211b 100644 --- a/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs +++ b/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs @@ -154,6 +154,8 @@ public static async IAsyncEnumerable SelectManyAsync( /// Collects every item from the source into a list. This overload has no work to parallelize; /// it is kept for binary compatibility with assemblies compiled against v3 (notably TUnit). /// + [Obsolete("This overload only collects the stream and parallelizes nothing. Pass a selector, or enumerate the stream directly. It exists for binary compatibility with assemblies compiled against v3.")] + [System.ComponentModel.EditorBrowsable(System.ComponentModel.EditorBrowsableState.Never)] public static async Task> ProcessInParallel( this IAsyncEnumerable items, CancellationToken cancellationToken = default) From c563673293e986d11be2fb7eae7abaa3a200554a Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Tue, 21 Jul 2026 22:31:21 +0100 Subject: [PATCH 3/4] fix: cancel in-flight work when an async-enumerable processor is disposed Explicit Dispose/DisposeAsync on the six IAsyncEnumerableProcessor implementations previously disposed the linked CancellationTokenSource without cancelling it, so an active ExecuteAsync kept running and later caller-token cancellation no longer reached the pipeline (disposal had severed the link). Disposal now cancels first; the completion-path self-dispose skips the cancel since nothing is left in flight. Both paths race-safely share an Interlocked guard. Builders capture the processor token up front so late selector invocations never touch a disposed source. Regression test: disposing mid-run terminates the pipeline with OperationCanceledException. --- .../DisposalRegressionTests.cs | 44 +++++++++++++++++++ ...EnumerableActionAsyncProcessorBuilder_1.cs | 4 +- ...EnumerableActionAsyncProcessorBuilder_2.cs | 4 +- .../AsyncEnumerableBatchProcessor.cs | 19 +++++++- .../AsyncEnumerableOneAtATimeProcessor.cs | 19 +++++++- .../AsyncEnumerableParallelProcessor.cs | 19 +++++++- .../ResultAsyncEnumerableBatchProcessor.cs | 19 +++++++- ...esultAsyncEnumerableOneAtATimeProcessor.cs | 19 +++++++- .../ResultAsyncEnumerableParallelProcessor.cs | 19 +++++++- 9 files changed, 158 insertions(+), 8 deletions(-) diff --git a/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs b/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs index e49995c..429c68f 100644 --- a/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs +++ b/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs @@ -161,6 +161,50 @@ public async Task AsyncEnumerable_Result_Processor_Supports_Await_Using_Without_ } } + [Test, Timeout(30_000)] + public async Task AsyncEnumerable_Processor_Dispose_During_Execution_Cancels_Processing(CancellationToken cancellationToken) + { + var firstItemStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + var processor = InfiniteAsyncEnumerable() + .ForEachAsync(async _ => + { + firstItemStarted.TrySetResult(); + await Task.Yield(); + }) + .ProcessInParallel(maxConcurrency: 1); + + var executeTask = processor.ExecuteAsync(); + await firstItemStarted.Task; + + processor.Dispose(); + + Exception? caught = null; + try + { + await executeTask.WaitAsync(TimeSpan.FromSeconds(10), cancellationToken); + } + catch (Exception exception) + { + caught = exception; + } + + // A TimeoutException here means disposal did not cancel the in-flight run. + await Assert.That(caught is OperationCanceledException).IsTrue(); + } + + private static async IAsyncEnumerable InfiniteAsyncEnumerable( + [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken = default) + { + var i = 0; + while (true) + { + cancellationToken.ThrowIfCancellationRequested(); + yield return i++; + await Task.Yield(); + } + } + private static async IAsyncEnumerable GenerateAsyncEnumerable(int count) { for (var i = 0; i < count; i++) diff --git a/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_1.cs index 83a72a0..c1c8199 100644 --- a/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_1.cs @@ -26,7 +26,9 @@ public AsyncEnumerableActionAsyncProcessorBuilder( { _items = items; _cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - _taskSelector = item => taskSelector(item, _cancellationTokenSource.Token); + // Capture the token now: the source may be disposed by the time a late item runs the selector. + var processorToken = _cancellationTokenSource.Token; + _taskSelector = item => taskSelector(item, processorToken); } /// diff --git a/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_2.cs b/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_2.cs index b40f21b..654f751 100644 --- a/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_2.cs +++ b/EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_2.cs @@ -26,7 +26,9 @@ public AsyncEnumerableActionAsyncProcessorBuilder( { _items = items; _cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - _taskSelector = item => taskSelector(item, _cancellationTokenSource.Token); + // Capture the token now: the source may be disposed by the time a late item runs the selector. + var processorToken = _cancellationTokenSource.Token; + _taskSelector = item => taskSelector(item, processorToken); } /// diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs index bd24f1b..f12c9ff 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs @@ -8,6 +8,7 @@ public sealed class AsyncEnumerableBatchProcessor : IAsyncEnumerableProc private readonly Func _taskSelector; private readonly int _batchSize; private readonly CancellationTokenSource _cancellationTokenSource; + private int _disposed; internal AsyncEnumerableBatchProcessor( IAsyncEnumerable items, @@ -48,7 +49,7 @@ public async Task ExecuteAsync() } finally { - _cancellationTokenSource.Dispose(); + DisposeCancellationSource(cancelFirst: false); } } @@ -60,6 +61,22 @@ private async Task ProcessBatch(List batch) public void Dispose() { + DisposeCancellationSource(cancelFirst: true); + } + + // Explicit disposal cancels in-flight work first; the completion path has nothing left to cancel. + private void DisposeCancellationSource(bool cancelFirst) + { + if (Interlocked.Exchange(ref _disposed, 1) != 0) + { + return; + } + + if (cancelFirst) + { + _cancellationTokenSource.Cancel(); + } + _cancellationTokenSource.Dispose(); } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs index f5fb486..b9a30bf 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs @@ -10,6 +10,7 @@ public sealed class AsyncEnumerableOneAtATimeProcessor : IAsyncEnumerabl private readonly IAsyncEnumerable _items; private readonly Func _taskSelector; private readonly CancellationTokenSource _cancellationTokenSource; + private int _disposed; internal AsyncEnumerableOneAtATimeProcessor( IAsyncEnumerable items, @@ -34,12 +35,28 @@ public async Task ExecuteAsync() } finally { - _cancellationTokenSource.Dispose(); + DisposeCancellationSource(cancelFirst: false); } } public void Dispose() { + DisposeCancellationSource(cancelFirst: true); + } + + // Explicit disposal cancels in-flight work first; the completion path has nothing left to cancel. + private void DisposeCancellationSource(bool cancelFirst) + { + if (Interlocked.Exchange(ref _disposed, 1) != 0) + { + return; + } + + if (cancelFirst) + { + _cancellationTokenSource.Cancel(); + } + _cancellationTokenSource.Dispose(); } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs index 04dd3f1..cc62b54 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs @@ -9,6 +9,7 @@ public sealed class AsyncEnumerableParallelProcessor : IAsyncEnumerableP private readonly int? _maxConcurrency; private readonly bool _scheduleOnThreadPool; private readonly CancellationTokenSource _cancellationTokenSource; + private int _disposed; internal AsyncEnumerableParallelProcessor( IAsyncEnumerable items, @@ -77,12 +78,28 @@ await AsyncEnumerableWorkerPool.ProcessAsync( } finally { - _cancellationTokenSource.Dispose(); + DisposeCancellationSource(cancelFirst: false); } } public void Dispose() { + DisposeCancellationSource(cancelFirst: true); + } + + // Explicit disposal cancels in-flight work first; the completion path has nothing left to cancel. + private void DisposeCancellationSource(bool cancelFirst) + { + if (Interlocked.Exchange(ref _disposed, 1) != 0) + { + return; + } + + if (cancelFirst) + { + _cancellationTokenSource.Cancel(); + } + _cancellationTokenSource.Dispose(); } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs index b5e889e..5a5c9a9 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs @@ -9,6 +9,7 @@ public sealed class ResultAsyncEnumerableBatchProcessor : IAsyn private readonly Func> _taskSelector; private readonly int _batchSize; private readonly CancellationTokenSource _cancellationTokenSource; + private int _disposed; internal ResultAsyncEnumerableBatchProcessor( IAsyncEnumerable items, @@ -55,7 +56,7 @@ public async IAsyncEnumerable ExecuteAsync() } finally { - _cancellationTokenSource.Dispose(); + DisposeCancellationSource(cancelFirst: false); } } @@ -72,6 +73,22 @@ private async IAsyncEnumerable ProcessBatch(List batch, [Enumer public void Dispose() { + DisposeCancellationSource(cancelFirst: true); + } + + // Explicit disposal cancels in-flight work first; the completion path has nothing left to cancel. + private void DisposeCancellationSource(bool cancelFirst) + { + if (Interlocked.Exchange(ref _disposed, 1) != 0) + { + return; + } + + if (cancelFirst) + { + _cancellationTokenSource.Cancel(); + } + _cancellationTokenSource.Dispose(); } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs index 7b49e28..f27536c 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs @@ -10,6 +10,7 @@ public sealed class ResultAsyncEnumerableOneAtATimeProcessor : private readonly IAsyncEnumerable _items; private readonly Func> _taskSelector; private readonly CancellationTokenSource _cancellationTokenSource; + private int _disposed; internal ResultAsyncEnumerableOneAtATimeProcessor( IAsyncEnumerable items, @@ -35,12 +36,28 @@ public async IAsyncEnumerable ExecuteAsync() } finally { - _cancellationTokenSource.Dispose(); + DisposeCancellationSource(cancelFirst: false); } } public void Dispose() { + DisposeCancellationSource(cancelFirst: true); + } + + // Explicit disposal cancels in-flight work first; the completion path has nothing left to cancel. + private void DisposeCancellationSource(bool cancelFirst) + { + if (Interlocked.Exchange(ref _disposed, 1) != 0) + { + return; + } + + if (cancelFirst) + { + _cancellationTokenSource.Cancel(); + } + _cancellationTokenSource.Dispose(); } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs index c98af3c..026894a 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs @@ -10,6 +10,7 @@ public sealed class ResultAsyncEnumerableParallelProcessor : IA private readonly int? _maxConcurrency; private readonly bool _scheduleOnThreadPool; private readonly CancellationTokenSource _cancellationTokenSource; + private int _disposed; internal ResultAsyncEnumerableParallelProcessor( IAsyncEnumerable items, @@ -94,12 +95,28 @@ public async IAsyncEnumerable ExecuteAsync() } finally { - _cancellationTokenSource.Dispose(); + DisposeCancellationSource(cancelFirst: false); } } public void Dispose() { + DisposeCancellationSource(cancelFirst: true); + } + + // Explicit disposal cancels in-flight work first; the completion path has nothing left to cancel. + private void DisposeCancellationSource(bool cancelFirst) + { + if (Interlocked.Exchange(ref _disposed, 1) != 0) + { + return; + } + + if (cancelFirst) + { + _cancellationTokenSource.Cancel(); + } + _cancellationTokenSource.Dispose(); } From 19b679e5848ebd6b1ba7433ee2024142b2bab4f1 Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Tue, 21 Jul 2026 22:40:28 +0100 Subject: [PATCH 4/4] fix: DisposeAsync awaits the in-flight run before returning Async disposal previously cancelled the linked source and returned immediately; a selector already running could still produce side effects after 'await processor.DisposeAsync()'. Mirror ProcessorLifecycle: void processors track their execution task, streaming processors complete a signal when enumeration finishes, and DisposeAsync cancels then waits up to the shared 30-second disposal window before releasing the source. Synchronous Dispose still cancels without blocking. --- .../DisposalRegressionTests.cs | 44 ++++++++++++++++++ .../ProcessorLifecycle.cs | 2 +- .../AsyncEnumerableBatchProcessor.cs | 45 +++++++++++++++++-- .../AsyncEnumerableOneAtATimeProcessor.cs | 45 +++++++++++++++++-- .../AsyncEnumerableParallelProcessor.cs | 45 +++++++++++++++++-- .../ResultAsyncEnumerableBatchProcessor.cs | 39 ++++++++++++++-- ...esultAsyncEnumerableOneAtATimeProcessor.cs | 39 ++++++++++++++-- .../ResultAsyncEnumerableParallelProcessor.cs | 39 ++++++++++++++-- 8 files changed, 276 insertions(+), 22 deletions(-) diff --git a/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs b/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs index 429c68f..736b8f9 100644 --- a/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs +++ b/EnumerableAsyncProcessor.UnitTests/DisposalRegressionTests.cs @@ -193,6 +193,50 @@ public async Task AsyncEnumerable_Processor_Dispose_During_Execution_Cancels_Pro await Assert.That(caught is OperationCanceledException).IsTrue(); } + [Test, Timeout(30_000)] + public async Task AsyncEnumerable_Processor_DisposeAsync_Waits_For_InFlight_Selector(CancellationToken cancellationToken) + { + var firstItemStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var releaseItem = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var itemsCompleted = 0; + + var processor = InfiniteAsyncEnumerable() + .ForEachAsync(async _ => + { + firstItemStarted.TrySetResult(); + await releaseItem.Task; + Interlocked.Increment(ref itemsCompleted); + }) + .ProcessInParallel(maxConcurrency: 1); + + var executeTask = processor.ExecuteAsync(); + await firstItemStarted.Task; + + var disposeTask = processor.DisposeAsync().AsTask(); + + // The selector ignores cancellation and is still blocked, so disposal must still be waiting. + await Assert.That(disposeTask.IsCompleted).IsFalse(); + + releaseItem.TrySetResult(); + await disposeTask.WaitAsync(TimeSpan.FromSeconds(10), cancellationToken); + + // At least the in-flight item finished before disposal returned; the worker may also + // drain one already-buffered item before it observes cancellation. + await Assert.That(itemsCompleted).IsGreaterThanOrEqualTo(1); + + Exception? caught = null; + try + { + await executeTask.WaitAsync(TimeSpan.FromSeconds(10), cancellationToken); + } + catch (Exception exception) + { + caught = exception; + } + + await Assert.That(caught is OperationCanceledException).IsTrue(); + } + private static async IAsyncEnumerable InfiniteAsyncEnumerable( [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken = default) { diff --git a/EnumerableAsyncProcessor/ProcessorLifecycle.cs b/EnumerableAsyncProcessor/ProcessorLifecycle.cs index 08c92bc..df9cfcf 100644 --- a/EnumerableAsyncProcessor/ProcessorLifecycle.cs +++ b/EnumerableAsyncProcessor/ProcessorLifecycle.cs @@ -8,7 +8,7 @@ namespace EnumerableAsyncProcessor; /// internal sealed class ProcessorLifecycle { - private static readonly TimeSpan DisposalTimeout = TimeSpan.FromSeconds(30); + internal static readonly TimeSpan DisposalTimeout = TimeSpan.FromSeconds(30); private readonly CancellationTokenSource _cancellationTokenSource; private readonly Action _trySetCanceledAll; diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs index f12c9ff..82b79f5 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableBatchProcessor.cs @@ -9,6 +9,7 @@ public sealed class AsyncEnumerableBatchProcessor : IAsyncEnumerableProc private readonly int _batchSize; private readonly CancellationTokenSource _cancellationTokenSource; private int _disposed; + private Task? _executionTask; internal AsyncEnumerableBatchProcessor( IAsyncEnumerable items, @@ -22,7 +23,14 @@ internal AsyncEnumerableBatchProcessor( _cancellationTokenSource = cancellationTokenSource; } - public async Task ExecuteAsync() + public Task ExecuteAsync() + { + var executionTask = ExecuteCoreAsync(); + _executionTask = executionTask; + return executionTask; + } + + private async Task ExecuteCoreAsync() { var cancellationToken = _cancellationTokenSource.Token; @@ -80,9 +88,38 @@ private void DisposeCancellationSource(bool cancelFirst) _cancellationTokenSource.Dispose(); } - public ValueTask DisposeAsync() + public async ValueTask DisposeAsync() { - Dispose(); - return ValueTask.CompletedTask; + CancelForDisposal(); + + // Mirror ProcessorLifecycle: give the in-flight run a bounded window to observe cancellation. + if (_executionTask is { IsCompleted: false } executionTask) + { + try + { + await executionTask.WaitAsync(ProcessorLifecycle.DisposalTimeout).ConfigureAwait(false); + } + catch + { + // Cancellation, failure, or timeout of the in-flight run; disposal must not throw. + } + } + + DisposeCancellationSource(cancelFirst: false); + } + + private void CancelForDisposal() + { + try + { + if (Volatile.Read(ref _disposed) == 0) + { + _cancellationTokenSource.Cancel(); + } + } + catch (ObjectDisposedException) + { + // The run completed and disposed the source concurrently - nothing left to cancel. + } } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs index b9a30bf..28c0936 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs @@ -11,6 +11,7 @@ public sealed class AsyncEnumerableOneAtATimeProcessor : IAsyncEnumerabl private readonly Func _taskSelector; private readonly CancellationTokenSource _cancellationTokenSource; private int _disposed; + private Task? _executionTask; internal AsyncEnumerableOneAtATimeProcessor( IAsyncEnumerable items, @@ -22,7 +23,14 @@ internal AsyncEnumerableOneAtATimeProcessor( _cancellationTokenSource = cancellationTokenSource; } - public async Task ExecuteAsync() + public Task ExecuteAsync() + { + var executionTask = ExecuteCoreAsync(); + _executionTask = executionTask; + return executionTask; + } + + private async Task ExecuteCoreAsync() { var cancellationToken = _cancellationTokenSource.Token; @@ -60,9 +68,38 @@ private void DisposeCancellationSource(bool cancelFirst) _cancellationTokenSource.Dispose(); } - public ValueTask DisposeAsync() + public async ValueTask DisposeAsync() { - Dispose(); - return ValueTask.CompletedTask; + CancelForDisposal(); + + // Mirror ProcessorLifecycle: give the in-flight run a bounded window to observe cancellation. + if (_executionTask is { IsCompleted: false } executionTask) + { + try + { + await executionTask.WaitAsync(ProcessorLifecycle.DisposalTimeout).ConfigureAwait(false); + } + catch + { + // Cancellation, failure, or timeout of the in-flight run; disposal must not throw. + } + } + + DisposeCancellationSource(cancelFirst: false); + } + + private void CancelForDisposal() + { + try + { + if (Volatile.Read(ref _disposed) == 0) + { + _cancellationTokenSource.Cancel(); + } + } + catch (ObjectDisposedException) + { + // The run completed and disposed the source concurrently - nothing left to cancel. + } } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs index cc62b54..e824987 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs @@ -10,6 +10,7 @@ public sealed class AsyncEnumerableParallelProcessor : IAsyncEnumerableP private readonly bool _scheduleOnThreadPool; private readonly CancellationTokenSource _cancellationTokenSource; private int _disposed; + private Task? _executionTask; internal AsyncEnumerableParallelProcessor( IAsyncEnumerable items, @@ -25,7 +26,14 @@ internal AsyncEnumerableParallelProcessor( _cancellationTokenSource = cancellationTokenSource; } - public async Task ExecuteAsync() + public Task ExecuteAsync() + { + var executionTask = ExecuteCoreAsync(); + _executionTask = executionTask; + return executionTask; + } + + private async Task ExecuteCoreAsync() { var cancellationToken = _cancellationTokenSource.Token; @@ -103,9 +111,38 @@ private void DisposeCancellationSource(bool cancelFirst) _cancellationTokenSource.Dispose(); } - public ValueTask DisposeAsync() + public async ValueTask DisposeAsync() + { + CancelForDisposal(); + + // Mirror ProcessorLifecycle: give the in-flight run a bounded window to observe cancellation. + if (_executionTask is { IsCompleted: false } executionTask) + { + try + { + await executionTask.WaitAsync(ProcessorLifecycle.DisposalTimeout).ConfigureAwait(false); + } + catch + { + // Cancellation, failure, or timeout of the in-flight run; disposal must not throw. + } + } + + DisposeCancellationSource(cancelFirst: false); + } + + private void CancelForDisposal() { - Dispose(); - return ValueTask.CompletedTask; + try + { + if (Volatile.Read(ref _disposed) == 0) + { + _cancellationTokenSource.Cancel(); + } + } + catch (ObjectDisposedException) + { + // The run completed and disposed the source concurrently - nothing left to cancel. + } } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs index 5a5c9a9..fadeb29 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableBatchProcessor.cs @@ -10,6 +10,7 @@ public sealed class ResultAsyncEnumerableBatchProcessor : IAsyn private readonly int _batchSize; private readonly CancellationTokenSource _cancellationTokenSource; private int _disposed; + private TaskCompletionSource? _executionCompleted; internal ResultAsyncEnumerableBatchProcessor( IAsyncEnumerable items, @@ -25,6 +26,8 @@ internal ResultAsyncEnumerableBatchProcessor( public async IAsyncEnumerable ExecuteAsync() { + var executionCompleted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + _executionCompleted = executionCompleted; var cancellationToken = _cancellationTokenSource.Token; try @@ -57,6 +60,7 @@ public async IAsyncEnumerable ExecuteAsync() finally { DisposeCancellationSource(cancelFirst: false); + executionCompleted.TrySetResult(); } } @@ -92,9 +96,38 @@ private void DisposeCancellationSource(bool cancelFirst) _cancellationTokenSource.Dispose(); } - public ValueTask DisposeAsync() + public async ValueTask DisposeAsync() { - Dispose(); - return ValueTask.CompletedTask; + CancelForDisposal(); + + // Mirror ProcessorLifecycle: give an in-flight enumeration a bounded window to observe cancellation. + if (_executionCompleted is { Task.IsCompleted: false } executionCompleted) + { + try + { + await executionCompleted.Task.WaitAsync(ProcessorLifecycle.DisposalTimeout).ConfigureAwait(false); + } + catch + { + // Timeout of the in-flight enumeration; disposal must not throw. + } + } + + DisposeCancellationSource(cancelFirst: false); + } + + private void CancelForDisposal() + { + try + { + if (Volatile.Read(ref _disposed) == 0) + { + _cancellationTokenSource.Cancel(); + } + } + catch (ObjectDisposedException) + { + // The run completed and disposed the source concurrently - nothing left to cancel. + } } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs index f27536c..a519693 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs @@ -11,6 +11,7 @@ public sealed class ResultAsyncEnumerableOneAtATimeProcessor : private readonly Func> _taskSelector; private readonly CancellationTokenSource _cancellationTokenSource; private int _disposed; + private TaskCompletionSource? _executionCompleted; internal ResultAsyncEnumerableOneAtATimeProcessor( IAsyncEnumerable items, @@ -24,6 +25,8 @@ internal ResultAsyncEnumerableOneAtATimeProcessor( public async IAsyncEnumerable ExecuteAsync() { + var executionCompleted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + _executionCompleted = executionCompleted; var cancellationToken = _cancellationTokenSource.Token; try @@ -37,6 +40,7 @@ public async IAsyncEnumerable ExecuteAsync() finally { DisposeCancellationSource(cancelFirst: false); + executionCompleted.TrySetResult(); } } @@ -61,9 +65,38 @@ private void DisposeCancellationSource(bool cancelFirst) _cancellationTokenSource.Dispose(); } - public ValueTask DisposeAsync() + public async ValueTask DisposeAsync() { - Dispose(); - return ValueTask.CompletedTask; + CancelForDisposal(); + + // Mirror ProcessorLifecycle: give an in-flight enumeration a bounded window to observe cancellation. + if (_executionCompleted is { Task.IsCompleted: false } executionCompleted) + { + try + { + await executionCompleted.Task.WaitAsync(ProcessorLifecycle.DisposalTimeout).ConfigureAwait(false); + } + catch + { + // Timeout of the in-flight enumeration; disposal must not throw. + } + } + + DisposeCancellationSource(cancelFirst: false); + } + + private void CancelForDisposal() + { + try + { + if (Volatile.Read(ref _disposed) == 0) + { + _cancellationTokenSource.Cancel(); + } + } + catch (ObjectDisposedException) + { + // The run completed and disposed the source concurrently - nothing left to cancel. + } } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs index 026894a..aa5af63 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs @@ -11,6 +11,7 @@ public sealed class ResultAsyncEnumerableParallelProcessor : IA private readonly bool _scheduleOnThreadPool; private readonly CancellationTokenSource _cancellationTokenSource; private int _disposed; + private TaskCompletionSource? _executionCompleted; internal ResultAsyncEnumerableParallelProcessor( IAsyncEnumerable items, @@ -28,6 +29,8 @@ internal ResultAsyncEnumerableParallelProcessor( public async IAsyncEnumerable ExecuteAsync() { + var executionCompleted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + _executionCompleted = executionCompleted; var cancellationToken = _cancellationTokenSource.Token; try @@ -96,6 +99,7 @@ public async IAsyncEnumerable ExecuteAsync() finally { DisposeCancellationSource(cancelFirst: false); + executionCompleted.TrySetResult(); } } @@ -120,9 +124,38 @@ private void DisposeCancellationSource(bool cancelFirst) _cancellationTokenSource.Dispose(); } - public ValueTask DisposeAsync() + public async ValueTask DisposeAsync() { - Dispose(); - return ValueTask.CompletedTask; + CancelForDisposal(); + + // Mirror ProcessorLifecycle: give an in-flight enumeration a bounded window to observe cancellation. + if (_executionCompleted is { Task.IsCompleted: false } executionCompleted) + { + try + { + await executionCompleted.Task.WaitAsync(ProcessorLifecycle.DisposalTimeout).ConfigureAwait(false); + } + catch + { + // Timeout of the in-flight enumeration; disposal must not throw. + } + } + + DisposeCancellationSource(cancelFirst: false); + } + + private void CancelForDisposal() + { + try + { + if (Volatile.Read(ref _disposed) == 0) + { + _cancellationTokenSource.Cancel(); + } + } + catch (ObjectDisposedException) + { + // The run completed and disposed the source concurrently - nothing left to cancel. + } } }