diff --git a/EnumerableAsyncProcessor.UnitTests/ValidationRegressionTests.cs b/EnumerableAsyncProcessor.UnitTests/ValidationRegressionTests.cs index 9fbe11a..209b55d 100644 --- a/EnumerableAsyncProcessor.UnitTests/ValidationRegressionTests.cs +++ b/EnumerableAsyncProcessor.UnitTests/ValidationRegressionTests.cs @@ -50,6 +50,9 @@ await AssertThrows(() => // The result variant previously skipped this validation entirely await AssertThrows(() => new[] { 1 }.SelectAsync(i => Task.FromResult(i)).ProcessInParallel(0, TimeSpan.FromSeconds(1))); + + await AssertThrows(() => + new[] { 1 }.SelectAsync(i => Task.FromResult(i)).ProcessInParallel(1, TimeSpan.FromSeconds(1), 0)); } [Test] diff --git a/EnumerableAsyncProcessor.UnitTests/WorkerPoolBehaviourTests.cs b/EnumerableAsyncProcessor.UnitTests/WorkerPoolBehaviourTests.cs index 5aa469d..7caa47b 100644 --- a/EnumerableAsyncProcessor.UnitTests/WorkerPoolBehaviourTests.cs +++ b/EnumerableAsyncProcessor.UnitTests/WorkerPoolBehaviourTests.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Concurrent; +using System.Diagnostics; using System.Linq; using System.Threading; using System.Threading.Tasks; @@ -133,6 +134,71 @@ public async Task Timed_RateLimited_Processor_Cancels_Unprocessed_Items_Promptly await processor.DisposeAsync(); } + [Test, Timeout(10_000)] + public async Task Timed_Rate_Limit_Allows_Concurrency_Independent_Of_Permit_Count(CancellationToken cancellationToken) + { + const int itemCount = 6; + + var startedCount = 0; + var allStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + await using var processor = Enumerable.Range(0, itemCount).ToList() + .ForEachAsync(async _ => + { + if (Interlocked.Increment(ref startedCount) == itemCount) + { + allStarted.TrySetResult(); + } + + await release.Task; + }, cancellationToken) + .ProcessInParallel( + permitsPerWindow: 2, + window: TimeSpan.FromMilliseconds(100), + maxConcurrency: itemCount); + + try + { + await allStarted.Task.WaitAsync(TimeSpan.FromSeconds(3), cancellationToken); + } + finally + { + release.TrySetResult(); + } + + await processor.WaitAsync(); + + await Assert.That(startedCount).IsEqualTo(itemCount); + } + + [Test, Retry(3), Timeout(10_000)] + public async Task Timed_Rate_Limit_Does_Not_Exceed_Permits_At_Replenishment(CancellationToken cancellationToken) + { + const int permitsPerWindow = 3; + var window = TimeSpan.FromMilliseconds(200); + var startedAt = new ConcurrentBag(); + var stopwatch = Stopwatch.StartNew(); + + await using var processor = Enumerable.Range(0, 9).ToList() + .ForEachAsync(_ => + { + startedAt.Add(stopwatch.Elapsed); + return Task.CompletedTask; + }, cancellationToken) + .ProcessInParallel(permitsPerWindow, window, maxConcurrency: 9); + + await processor.WaitAsync(); + + var orderedStarts = startedAt.OrderBy(x => x).ToArray(); + + for (var i = permitsPerWindow; i < orderedStarts.Length; i++) + { + await Assert.That(orderedStarts[i] - orderedStarts[i - permitsPerWindow]) + .IsGreaterThan(window / 2); + } + } + [Test] public async Task Result_Order_Is_Preserved_Regardless_Of_Completion_Order() { diff --git a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs index 9c32159..9c2c1b2 100644 --- a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs +++ b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs @@ -31,7 +31,16 @@ public IAsyncProcessor ProcessInBatches(int batchSize) public IAsyncProcessor ProcessInParallel(int maxConcurrency, TimeSpan timeSpan) { - return new TimedRateLimitedParallelAsyncProcessor(_count, _taskSelector, maxConcurrency, timeSpan, _cancellationTokenSource).StartProcessing(); + return ProcessInParallel(maxConcurrency, timeSpan, maxConcurrency); + } + + /// Processes tasks with independent start-rate and concurrency limits. + /// Maximum operations that may start in each window. + /// Rate-limit replenishment window. + /// Maximum operations that may remain in flight. + public IAsyncProcessor ProcessInParallel(int permitsPerWindow, TimeSpan window, int maxConcurrency) + { + return new TimedRateLimitedParallelAsyncProcessor(_count, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource).StartProcessing(); } /// diff --git a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs index 9eb3a91..65b4ec9 100644 --- a/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs @@ -31,7 +31,16 @@ public IAsyncProcessor ProcessInBatches(int batchSize) public IAsyncProcessor ProcessInParallel(int maxConcurrency, TimeSpan timeSpan) { - return new ResultTimedRateLimitedParallelAsyncProcessor(_count, _taskSelector, maxConcurrency, timeSpan, _cancellationTokenSource).StartProcessing(); + return ProcessInParallel(maxConcurrency, timeSpan, maxConcurrency); + } + + /// Processes tasks with independent start-rate and concurrency limits. + /// Maximum operations that may start in each window. + /// Rate-limit replenishment window. + /// Maximum operations that may remain in flight. + public IAsyncProcessor ProcessInParallel(int permitsPerWindow, TimeSpan window, int maxConcurrency) + { + return new ResultTimedRateLimitedParallelAsyncProcessor(_count, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource).StartProcessing(); } /// diff --git a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs index 3ac2223..402d3fe 100644 --- a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs +++ b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs @@ -32,7 +32,16 @@ public IAsyncProcessor ProcessInBatches(int batchSize) public IAsyncProcessor ProcessInParallel(int maxConcurrency, TimeSpan timeSpan) { - return new TimedRateLimitedParallelAsyncProcessor(_items, _taskSelector, maxConcurrency, timeSpan, _cancellationTokenSource) + return ProcessInParallel(maxConcurrency, timeSpan, maxConcurrency); + } + + /// Processes items with independent start-rate and concurrency limits. + /// Maximum operations that may start in each window. + /// Rate-limit replenishment window. + /// Maximum operations that may remain in flight. + public IAsyncProcessor ProcessInParallel(int permitsPerWindow, TimeSpan window, int maxConcurrency) + { + return new TimedRateLimitedParallelAsyncProcessor(_items, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource) .StartProcessing(); } diff --git a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs index fbad93f..ca902a6 100644 --- a/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs +++ b/EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs @@ -52,7 +52,19 @@ public IAsyncProcessor ProcessInBatches(int batchSize) /// public IAsyncProcessor ProcessInParallel(int maxConcurrency, TimeSpan timeSpan) { - return new ResultTimedRateLimitedParallelAsyncProcessor(_items, _taskSelector, maxConcurrency, timeSpan, _cancellationTokenSource).StartProcessing(); + return ProcessInParallel(maxConcurrency, timeSpan, maxConcurrency); + } + + /// + /// Processes items with independent start-rate and concurrency limits. + /// + /// Maximum operations that may start in each window. + /// Rate-limit replenishment window. + /// Maximum operations that may remain in flight. + /// An async processor that implements IDisposable and IAsyncDisposable. + public IAsyncProcessor ProcessInParallel(int permitsPerWindow, TimeSpan window, int maxConcurrency) + { + return new ResultTimedRateLimitedParallelAsyncProcessor(_items, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource).StartProcessing(); } /// diff --git a/EnumerableAsyncProcessor/EnumerableAsyncProcessor.csproj b/EnumerableAsyncProcessor/EnumerableAsyncProcessor.csproj index e917f3e..7fd772d 100644 --- a/EnumerableAsyncProcessor/EnumerableAsyncProcessor.csproj +++ b/EnumerableAsyncProcessor/EnumerableAsyncProcessor.csproj @@ -29,6 +29,7 @@ + diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs index e512d84..8e7f6a2 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs @@ -46,6 +46,6 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper => // Throttled processing runs on a fixed worker pool: P worker tasks instead of // one queued task, closure and semaphore wait per item - await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false); + await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, rateLimiter: null, CancellationToken).ConfigureAwait(false); } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs index 9dcdc42..fd7a364 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs @@ -46,6 +46,6 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper => // Throttled processing runs on a fixed worker pool: P worker tasks instead of // one queued task, closure and semaphore wait per item - await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false); + await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, rateLimiter: null, CancellationToken).ConfigureAwait(false); } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs index f85d1c4..a25bb9d 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs @@ -46,6 +46,6 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper => // Throttled processing runs on a fixed worker pool: P worker tasks instead of // one queued task, closure and semaphore wait per item - await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false); + await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, rateLimiter: null, CancellationToken).ConfigureAwait(false); } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs index 2348103..aa7f7e0 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs @@ -46,6 +46,6 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper => // Throttled processing runs on a fixed worker pool: P worker tasks instead of // one queued task, closure and semaphore wait per item - await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false); + await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, rateLimiter: null, CancellationToken).ConfigureAwait(false); } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs index bc85f94..cbb9752 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs @@ -5,20 +5,23 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; public class ResultTimedRateLimitedParallelAsyncProcessor : ResultAbstractAsyncProcessor { - private readonly int _levelsOfParallelism; - private readonly TimeSpan _timeSpan; + private readonly int _permitsPerWindow; + private readonly TimeSpan _window; + private readonly int _maxConcurrency; - internal ResultTimedRateLimitedParallelAsyncProcessor(int count, Func> taskSelector, int levelsOfParallelism, TimeSpan timeSpan, CancellationTokenSource cancellationTokenSource) : base(count, taskSelector, cancellationTokenSource) + internal ResultTimedRateLimitedParallelAsyncProcessor(int count, Func> taskSelector, int permitsPerWindow, TimeSpan window, int maxConcurrency, CancellationTokenSource cancellationTokenSource) : base(count, taskSelector, cancellationTokenSource) { - ValidationHelper.ThrowIfNegativeOrZero(levelsOfParallelism); - ValidationHelper.ThrowIfNegative(timeSpan); + ValidationHelper.ThrowIfNegativeOrZero(permitsPerWindow); + ValidationHelper.ThrowIfNegative(window); + ValidationHelper.ThrowIfNegativeOrZero(maxConcurrency); - _levelsOfParallelism = levelsOfParallelism; - _timeSpan = timeSpan; + _permitsPerWindow = permitsPerWindow; + _window = window; + _maxConcurrency = maxConcurrency; } internal override Task Process() { - return WorkerPool.ProcessAsync(TaskWrappers, _levelsOfParallelism, minimumIterationTime: _timeSpan, CancellationToken); + return WorkerPool.ProcessRateLimitedAsync(TaskWrappers, _maxConcurrency, _permitsPerWindow, _window, CancellationToken); } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs index e1d1160..f680ed9 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs @@ -5,20 +5,23 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors; public class ResultTimedRateLimitedParallelAsyncProcessor : ResultAbstractAsyncProcessor { - private readonly int _levelsOfParallelism; - private readonly TimeSpan _timeSpan; + private readonly int _permitsPerWindow; + private readonly TimeSpan _window; + private readonly int _maxConcurrency; - internal ResultTimedRateLimitedParallelAsyncProcessor(IEnumerable items, Func> taskSelector, int levelsOfParallelism, TimeSpan timeSpan, CancellationTokenSource cancellationTokenSource) : base(items, taskSelector, cancellationTokenSource) + internal ResultTimedRateLimitedParallelAsyncProcessor(IEnumerable items, Func> taskSelector, int permitsPerWindow, TimeSpan window, int maxConcurrency, CancellationTokenSource cancellationTokenSource) : base(items, taskSelector, cancellationTokenSource) { - ValidationHelper.ThrowIfNegativeOrZero(levelsOfParallelism); - ValidationHelper.ThrowIfNegative(timeSpan); + ValidationHelper.ThrowIfNegativeOrZero(permitsPerWindow); + ValidationHelper.ThrowIfNegative(window); + ValidationHelper.ThrowIfNegativeOrZero(maxConcurrency); - _levelsOfParallelism = levelsOfParallelism; - _timeSpan = timeSpan; + _permitsPerWindow = permitsPerWindow; + _window = window; + _maxConcurrency = maxConcurrency; } internal override Task Process() { - return WorkerPool.ProcessAsync(TaskWrappers, _levelsOfParallelism, minimumIterationTime: _timeSpan, CancellationToken); + return WorkerPool.ProcessRateLimitedAsync(TaskWrappers, _maxConcurrency, _permitsPerWindow, _window, CancellationToken); } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor.cs index 430cdf5..253c47f 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor.cs @@ -5,20 +5,23 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; public class TimedRateLimitedParallelAsyncProcessor : AbstractAsyncProcessor { - private readonly int _levelsOfParallelism; - private readonly TimeSpan _timeSpan; + private readonly int _permitsPerWindow; + private readonly TimeSpan _window; + private readonly int _maxConcurrency; - internal TimedRateLimitedParallelAsyncProcessor(int count, Func taskSelector, int levelsOfParallelism, TimeSpan timeSpan, CancellationTokenSource cancellationTokenSource) : base(count, taskSelector, cancellationTokenSource) + internal TimedRateLimitedParallelAsyncProcessor(int count, Func taskSelector, int permitsPerWindow, TimeSpan window, int maxConcurrency, CancellationTokenSource cancellationTokenSource) : base(count, taskSelector, cancellationTokenSource) { - ValidationHelper.ThrowIfNegativeOrZero(levelsOfParallelism); - ValidationHelper.ThrowIfNegative(timeSpan); + ValidationHelper.ThrowIfNegativeOrZero(permitsPerWindow); + ValidationHelper.ThrowIfNegative(window); + ValidationHelper.ThrowIfNegativeOrZero(maxConcurrency); - _levelsOfParallelism = levelsOfParallelism; - _timeSpan = timeSpan; + _permitsPerWindow = permitsPerWindow; + _window = window; + _maxConcurrency = maxConcurrency; } internal override Task Process() { - return WorkerPool.ProcessAsync(TaskWrappers, _levelsOfParallelism, minimumIterationTime: _timeSpan, CancellationToken); + return WorkerPool.ProcessRateLimitedAsync(TaskWrappers, _maxConcurrency, _permitsPerWindow, _window, CancellationToken); } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs index ba8c2c0..f9ea993 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs @@ -5,20 +5,23 @@ namespace EnumerableAsyncProcessor.RunnableProcessors; public class TimedRateLimitedParallelAsyncProcessor : AbstractAsyncProcessor { - private readonly int _levelsOfParallelism; - private readonly TimeSpan _timeSpan; + private readonly int _permitsPerWindow; + private readonly TimeSpan _window; + private readonly int _maxConcurrency; - internal TimedRateLimitedParallelAsyncProcessor(IEnumerable items, Func taskSelector, int levelsOfParallelism, TimeSpan timeSpan, CancellationTokenSource cancellationTokenSource) : base(items, taskSelector, cancellationTokenSource) + internal TimedRateLimitedParallelAsyncProcessor(IEnumerable items, Func taskSelector, int permitsPerWindow, TimeSpan window, int maxConcurrency, CancellationTokenSource cancellationTokenSource) : base(items, taskSelector, cancellationTokenSource) { - ValidationHelper.ThrowIfNegativeOrZero(levelsOfParallelism); - ValidationHelper.ThrowIfNegative(timeSpan); + ValidationHelper.ThrowIfNegativeOrZero(permitsPerWindow); + ValidationHelper.ThrowIfNegative(window); + ValidationHelper.ThrowIfNegativeOrZero(maxConcurrency); - _levelsOfParallelism = levelsOfParallelism; - _timeSpan = timeSpan; + _permitsPerWindow = permitsPerWindow; + _window = window; + _maxConcurrency = maxConcurrency; } internal override Task Process() { - return WorkerPool.ProcessAsync(TaskWrappers, _levelsOfParallelism, minimumIterationTime: _timeSpan, CancellationToken); + return WorkerPool.ProcessRateLimitedAsync(TaskWrappers, _maxConcurrency, _permitsPerWindow, _window, CancellationToken); } } diff --git a/EnumerableAsyncProcessor/WorkerPool.cs b/EnumerableAsyncProcessor/WorkerPool.cs index 04691bd..c65ce21 100644 --- a/EnumerableAsyncProcessor/WorkerPool.cs +++ b/EnumerableAsyncProcessor/WorkerPool.cs @@ -1,3 +1,5 @@ +using System.Threading.RateLimiting; + namespace EnumerableAsyncProcessor; /// @@ -7,14 +9,10 @@ namespace EnumerableAsyncProcessor; /// internal static class WorkerPool { - /// - /// When set, each worker holds its slot for at least this long per item, which caps throughput - /// at (workerCount / minimumIterationTime) operations for the timed rate-limited processors. - /// internal static Task ProcessAsync( TWrapper[] taskWrappers, int workerCount, - TimeSpan? minimumIterationTime, + RateLimiter? rateLimiter, CancellationToken cancellationToken) where TWrapper : ITaskWrapper { workerCount = Math.Min(workerCount, taskWrappers.Length); @@ -42,22 +40,49 @@ internal static Task ProcessAsync( return; } - // Process never throws; it completes the item's TaskCompletionSource instead, - // so one failed item cannot stop the worker from draining the rest. - var processTask = taskWrappers[index].Process(cancellationToken); - - if (minimumIterationTime is { } minimumTime) - { - await Task.WhenAll(processTask, Task.Delay(minimumTime, cancellationToken)).ConfigureAwait(false); - } - else + if (rateLimiter is not null) { - await processTask.ConfigureAwait(false); + using var lease = await rateLimiter.AcquireAsync(1, cancellationToken).ConfigureAwait(false); + + if (!lease.IsAcquired) + { + throw new InvalidOperationException("The rate limiter could not acquire a permit."); + } } + + // Process never throws; it completes the item's TaskCompletionSource instead, + // so one failed item cannot stop the worker from draining the rest. + await taskWrappers[index].Process(cancellationToken).ConfigureAwait(false); } }, cancellationToken); } return Task.WhenAll(workers); } + + internal static async Task ProcessRateLimitedAsync( + TWrapper[] taskWrappers, + int workerCount, + int permitsPerWindow, + TimeSpan window, + CancellationToken cancellationToken) where TWrapper : ITaskWrapper + { + if (window == TimeSpan.Zero) + { + await ProcessAsync(taskWrappers, workerCount, rateLimiter: null, cancellationToken).ConfigureAwait(false); + return; + } + + using var rateLimiter = new TokenBucketRateLimiter(new TokenBucketRateLimiterOptions + { + TokenLimit = permitsPerWindow, + TokensPerPeriod = permitsPerWindow, + ReplenishmentPeriod = window, + AutoReplenishment = true, + QueueProcessingOrder = QueueProcessingOrder.OldestFirst, + QueueLimit = workerCount + }); + + await ProcessAsync(taskWrappers, workerCount, rateLimiter, cancellationToken).ConfigureAwait(false); + } }