diff --git a/CLAUDE.md b/CLAUDE.md index e8024c7..d285057 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -29,8 +29,6 @@ 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 diff --git a/EnumerableAsyncProcessor.UnitTests/V3BinaryCompatibilityTests.cs b/EnumerableAsyncProcessor.UnitTests/V3BinaryCompatibilityTests.cs deleted file mode 100644 index d6759e4..0000000 --- a/EnumerableAsyncProcessor.UnitTests/V3BinaryCompatibilityTests.cs +++ /dev/null @@ -1,47 +0,0 @@ -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 47fe659..11a3965 100644 --- a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs +++ b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs @@ -54,16 +54,6 @@ 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)). - /// - [System.ComponentModel.EditorBrowsable(System.ComponentModel.EditorBrowsableState.Never)] - public IAsyncProcessor ProcessInParallel(int maxConcurrency) - { - return ProcessInParallel((int?)maxConcurrency); - } - public IAsyncProcessor ProcessOneAtATime() { return new OneAtATimeAsyncProcessor(_count, _taskSelector, _cancellationTokenSource).StartProcessing(); diff --git a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs index 0737a92..1c9d36b 100644 --- a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs @@ -54,16 +54,6 @@ 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)). - /// - [System.ComponentModel.EditorBrowsable(System.ComponentModel.EditorBrowsableState.Never)] - public IAsyncProcessor ProcessInParallel(int maxConcurrency) - { - return ProcessInParallel((int?)maxConcurrency); - } - public IAsyncProcessor ProcessOneAtATime() { return new ResultOneAtATimeAsyncProcessor(_count, _taskSelector, _cancellationTokenSource).StartProcessing(); diff --git a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs index bdb6ebb..766e4d9 100644 --- a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs @@ -57,16 +57,6 @@ public IAsyncProcessor ProcessInParallel(int? maxConcurrency = null, bool schedu .StartProcessing(); } - /// - /// 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); - } - public IAsyncProcessor ProcessOneAtATime() { return new OneAtATimeAsyncProcessor(_items, _taskSelector, _cancellationTokenSource) diff --git a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs index d687b80..9d45ca4 100644 --- a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs +++ b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs @@ -83,16 +83,6 @@ 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)). - /// - [System.ComponentModel.EditorBrowsable(System.ComponentModel.EditorBrowsableState.Never)] - public IAsyncProcessor ProcessInParallel(int maxConcurrency) - { - return ProcessInParallel((int?)maxConcurrency); - } - /// /// Process items one at a time sequentially. /// diff --git a/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs b/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs index 32b61ac..01f961c 100644 --- a/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs +++ b/EnumerableAsyncProcessor/Extensions/AsyncEnumerableExtensions.cs @@ -150,26 +150,6 @@ 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) - { - var results = new List(); - - await foreach (var item in items.WithCancellation(cancellationToken).ConfigureAwait(false)) - { - results.Add(item); - } - - return results; - } - /// /// Process items in parallel with transformation and return all results as IEnumerable when awaited. /// diff --git a/EnumerableAsyncProcessor/PublicAPI.Shipped.txt b/EnumerableAsyncProcessor/PublicAPI.Shipped.txt index 3464632..ab34a09 100644 --- a/EnumerableAsyncProcessor/PublicAPI.Shipped.txt +++ b/EnumerableAsyncProcessor/PublicAPI.Shipped.txt @@ -3,14 +3,12 @@ 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! @@ -48,7 +46,6 @@ 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! @@ -57,7 +54,6 @@ 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! @@ -155,7 +151,6 @@ static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.ForEachAsyn 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.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! static EnumerableAsyncProcessor.Extensions.AsyncEnumerableExtensions.SelectMany(this System.Collections.Generic.IAsyncEnumerable! items, System.Func!>! selector, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) -> System.Collections.Generic.IAsyncEnumerable! diff --git a/README.md b/README.md index 586daf4..6047321 100644 --- a/README.md +++ b/README.md @@ -26,7 +26,7 @@ Version 4 is a major release. Review these source and behaviour changes before u - **.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` 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. +- **The no-selector `IAsyncEnumerable.ProcessInParallel(...)` 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, or enumerate the stream directly. - **`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.