Skip to content

Commit d9edb29

Browse files
committed
feat: add token-bucket rate limiting
Timed processors limit operation starts per window. New overload separates permit count from maximum in-flight concurrency. Closes #335
1 parent 8c0ac6f commit d9edb29

16 files changed

Lines changed: 201 additions & 55 deletions

EnumerableAsyncProcessor.UnitTests/ValidationRegressionTests.cs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,9 @@ await AssertThrows<ArgumentOutOfRangeException>(() =>
5050
// The result variant previously skipped this validation entirely
5151
await AssertThrows<ArgumentOutOfRangeException>(() =>
5252
new[] { 1 }.SelectAsync(i => Task.FromResult(i)).ProcessInParallel(0, TimeSpan.FromSeconds(1)));
53+
54+
await AssertThrows<ArgumentOutOfRangeException>(() =>
55+
new[] { 1 }.SelectAsync(i => Task.FromResult(i)).ProcessInParallel(1, TimeSpan.FromSeconds(1), 0));
5356
}
5457

5558
[Test]

EnumerableAsyncProcessor.UnitTests/WorkerPoolBehaviourTests.cs

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
using System;
22
using System.Collections.Concurrent;
3+
using System.Diagnostics;
34
using System.Linq;
45
using System.Threading;
56
using System.Threading.Tasks;
@@ -133,6 +134,71 @@ public async Task Timed_RateLimited_Processor_Cancels_Unprocessed_Items_Promptly
133134
await processor.DisposeAsync();
134135
}
135136

137+
[Test, Timeout(10_000)]
138+
public async Task Timed_Rate_Limit_Allows_Concurrency_Independent_Of_Permit_Count(CancellationToken cancellationToken)
139+
{
140+
const int itemCount = 6;
141+
142+
var startedCount = 0;
143+
var allStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
144+
var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
145+
146+
await using var processor = Enumerable.Range(0, itemCount).ToList()
147+
.ForEachAsync(async _ =>
148+
{
149+
if (Interlocked.Increment(ref startedCount) == itemCount)
150+
{
151+
allStarted.TrySetResult();
152+
}
153+
154+
await release.Task;
155+
}, cancellationToken)
156+
.ProcessInParallel(
157+
permitsPerWindow: 2,
158+
window: TimeSpan.FromMilliseconds(100),
159+
maxConcurrency: itemCount);
160+
161+
try
162+
{
163+
await allStarted.Task.WaitAsync(TimeSpan.FromSeconds(3), cancellationToken);
164+
}
165+
finally
166+
{
167+
release.TrySetResult();
168+
}
169+
170+
await processor.WaitAsync();
171+
172+
await Assert.That(startedCount).IsEqualTo(itemCount);
173+
}
174+
175+
[Test, Retry(3), Timeout(10_000)]
176+
public async Task Timed_Rate_Limit_Does_Not_Exceed_Permits_At_Replenishment(CancellationToken cancellationToken)
177+
{
178+
const int permitsPerWindow = 3;
179+
var window = TimeSpan.FromMilliseconds(200);
180+
var startedAt = new ConcurrentBag<TimeSpan>();
181+
var stopwatch = Stopwatch.StartNew();
182+
183+
await using var processor = Enumerable.Range(0, 9).ToList()
184+
.ForEachAsync(_ =>
185+
{
186+
startedAt.Add(stopwatch.Elapsed);
187+
return Task.CompletedTask;
188+
}, cancellationToken)
189+
.ProcessInParallel(permitsPerWindow, window, maxConcurrency: 9);
190+
191+
await processor.WaitAsync();
192+
193+
var orderedStarts = startedAt.OrderBy(x => x).ToArray();
194+
195+
for (var i = permitsPerWindow; i < orderedStarts.Length; i++)
196+
{
197+
await Assert.That(orderedStarts[i] - orderedStarts[i - permitsPerWindow])
198+
.IsGreaterThan(window / 2);
199+
}
200+
}
201+
136202
[Test]
137203
public async Task Result_Order_Is_Preserved_Regardless_Of_Completion_Order()
138204
{

EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,16 @@ public IAsyncProcessor ProcessInBatches(int batchSize)
3131

3232
public IAsyncProcessor ProcessInParallel(int maxConcurrency, TimeSpan timeSpan)
3333
{
34-
return new TimedRateLimitedParallelAsyncProcessor(_count, _taskSelector, maxConcurrency, timeSpan, _cancellationTokenSource).StartProcessing();
34+
return ProcessInParallel(maxConcurrency, timeSpan, maxConcurrency);
35+
}
36+
37+
/// <summary>Processes tasks with independent start-rate and concurrency limits.</summary>
38+
/// <param name="permitsPerWindow">Maximum operations that may start in each window.</param>
39+
/// <param name="window">Rate-limit replenishment window.</param>
40+
/// <param name="maxConcurrency">Maximum operations that may remain in flight.</param>
41+
public IAsyncProcessor ProcessInParallel(int permitsPerWindow, TimeSpan window, int maxConcurrency)
42+
{
43+
return new TimedRateLimitedParallelAsyncProcessor(_count, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource).StartProcessing();
3544
}
3645

3746
/// <summary>

EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,16 @@ public IAsyncProcessor<TOutput> ProcessInBatches(int batchSize)
3131

3232
public IAsyncProcessor<TOutput> ProcessInParallel(int maxConcurrency, TimeSpan timeSpan)
3333
{
34-
return new ResultTimedRateLimitedParallelAsyncProcessor<TOutput>(_count, _taskSelector, maxConcurrency, timeSpan, _cancellationTokenSource).StartProcessing();
34+
return ProcessInParallel(maxConcurrency, timeSpan, maxConcurrency);
35+
}
36+
37+
/// <summary>Processes tasks with independent start-rate and concurrency limits.</summary>
38+
/// <param name="permitsPerWindow">Maximum operations that may start in each window.</param>
39+
/// <param name="window">Rate-limit replenishment window.</param>
40+
/// <param name="maxConcurrency">Maximum operations that may remain in flight.</param>
41+
public IAsyncProcessor<TOutput> ProcessInParallel(int permitsPerWindow, TimeSpan window, int maxConcurrency)
42+
{
43+
return new ResultTimedRateLimitedParallelAsyncProcessor<TOutput>(_count, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource).StartProcessing();
3544
}
3645

3746
/// <summary>

EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,16 @@ public IAsyncProcessor ProcessInBatches(int batchSize)
3232

3333
public IAsyncProcessor ProcessInParallel(int maxConcurrency, TimeSpan timeSpan)
3434
{
35-
return new TimedRateLimitedParallelAsyncProcessor<TInput>(_items, _taskSelector, maxConcurrency, timeSpan, _cancellationTokenSource)
35+
return ProcessInParallel(maxConcurrency, timeSpan, maxConcurrency);
36+
}
37+
38+
/// <summary>Processes items with independent start-rate and concurrency limits.</summary>
39+
/// <param name="permitsPerWindow">Maximum operations that may start in each window.</param>
40+
/// <param name="window">Rate-limit replenishment window.</param>
41+
/// <param name="maxConcurrency">Maximum operations that may remain in flight.</param>
42+
public IAsyncProcessor ProcessInParallel(int permitsPerWindow, TimeSpan window, int maxConcurrency)
43+
{
44+
return new TimedRateLimitedParallelAsyncProcessor<TInput>(_items, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource)
3645
.StartProcessing();
3746
}
3847

EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,19 @@ public IAsyncProcessor<TOutput> ProcessInBatches(int batchSize)
5252
/// </remarks>
5353
public IAsyncProcessor<TOutput> ProcessInParallel(int maxConcurrency, TimeSpan timeSpan)
5454
{
55-
return new ResultTimedRateLimitedParallelAsyncProcessor<TInput, TOutput>(_items, _taskSelector, maxConcurrency, timeSpan, _cancellationTokenSource).StartProcessing();
55+
return ProcessInParallel(maxConcurrency, timeSpan, maxConcurrency);
56+
}
57+
58+
/// <summary>
59+
/// Processes items with independent start-rate and concurrency limits.
60+
/// </summary>
61+
/// <param name="permitsPerWindow">Maximum operations that may start in each window.</param>
62+
/// <param name="window">Rate-limit replenishment window.</param>
63+
/// <param name="maxConcurrency">Maximum operations that may remain in flight.</param>
64+
/// <returns>An async processor that implements IDisposable and IAsyncDisposable.</returns>
65+
public IAsyncProcessor<TOutput> ProcessInParallel(int permitsPerWindow, TimeSpan window, int maxConcurrency)
66+
{
67+
return new ResultTimedRateLimitedParallelAsyncProcessor<TInput, TOutput>(_items, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource).StartProcessing();
5668
}
5769

5870
/// <summary>

EnumerableAsyncProcessor/EnumerableAsyncProcessor.csproj

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
<None Include="$(MSBuildThisFileDirectory)..\README.md" Pack="true" PackagePath="\" />
3030

3131
<PackageReference Include="Microsoft.SourceLink.GitHub" Version="10.0.301" PrivateAssets="All"/>
32+
<PackageReference Include="System.Threading.RateLimiting" Version="10.0.10" />
3233

3334
</ItemGroup>
3435

EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,6 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper =>
4646

4747
// Throttled processing runs on a fixed worker pool: P worker tasks instead of
4848
// one queued task, closure and semaphore wait per item
49-
await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false);
49+
await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, rateLimiter: null, CancellationToken).ConfigureAwait(false);
5050
}
5151
}

EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,6 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper =>
4646

4747
// Throttled processing runs on a fixed worker pool: P worker tasks instead of
4848
// one queued task, closure and semaphore wait per item
49-
await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false);
49+
await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, rateLimiter: null, CancellationToken).ConfigureAwait(false);
5050
}
5151
}

EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,6 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper =>
4646

4747
// Throttled processing runs on a fixed worker pool: P worker tasks instead of
4848
// one queued task, closure and semaphore wait per item
49-
await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false);
49+
await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, rateLimiter: null, CancellationToken).ConfigureAwait(false);
5050
}
5151
}

0 commit comments

Comments
 (0)