Skip to content

Commit bc5a93f

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 1f41055 commit bc5a93f

20 files changed

Lines changed: 209 additions & 63 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;
@@ -118,6 +119,71 @@ public async Task Timed_RateLimited_Processor_Cancels_Unprocessed_Items_Promptly
118119
await processor.DisposeAsync();
119120
}
120121

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

EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,16 @@ public IAsyncProcessor ProcessInParallel(int levelOfParallelism)
2929

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

3544
/// <summary>
@@ -77,4 +86,4 @@ public IAsyncProcessor ProcessOneAtATime()
7786
return new OneAtATimeAsyncProcessor(_count, _taskSelector, _cancellationTokenSource).StartProcessing();
7887
}
7988

80-
}
89+
}

EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,16 @@ public IAsyncProcessor<TOutput> ProcessInParallel(int levelOfParallelism)
2929

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

3544
/// <summary>
@@ -77,4 +86,4 @@ public IAsyncProcessor<TOutput> ProcessOneAtATime()
7786
return new ResultOneAtATimeAsyncProcessor<TOutput>(_count, _taskSelector, _cancellationTokenSource).StartProcessing();
7887
}
7988

80-
}
89+
}

EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,16 @@ public IAsyncProcessor ProcessInParallel(int levelOfParallelism)
3131

3232
public IAsyncProcessor ProcessInParallel(int levelOfParallelism, TimeSpan timeSpan)
3333
{
34-
return new TimedRateLimitedParallelAsyncProcessor<TInput>(_items, _taskSelector, levelOfParallelism, timeSpan, _cancellationTokenSource)
34+
return ProcessInParallel(levelOfParallelism, timeSpan, levelOfParallelism);
35+
}
36+
37+
/// <summary>Processes items 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<TInput>(_items, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource)
3544
.StartProcessing();
3645
}
3746

@@ -82,4 +91,4 @@ public IAsyncProcessor ProcessOneAtATime()
8291
.StartProcessing();
8392
}
8493

85-
}
94+
}

EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,19 @@ public IAsyncProcessor<TOutput> ProcessInParallel(int levelOfParallelism)
6060
/// </remarks>
6161
public IAsyncProcessor<TOutput> ProcessInParallel(int levelOfParallelism, TimeSpan timeSpan)
6262
{
63-
return new ResultTimedRateLimitedParallelAsyncProcessor<TInput, TOutput>(_items, _taskSelector, levelOfParallelism, timeSpan, _cancellationTokenSource).StartProcessing();
63+
return ProcessInParallel(levelOfParallelism, timeSpan, levelOfParallelism);
64+
}
65+
66+
/// <summary>
67+
/// Processes items with independent start-rate and concurrency limits.
68+
/// </summary>
69+
/// <param name="permitsPerWindow">Maximum operations that may start in each window.</param>
70+
/// <param name="window">Rate-limit replenishment window.</param>
71+
/// <param name="maxConcurrency">Maximum operations that may remain in flight.</param>
72+
/// <returns>An async processor that implements IDisposable and IAsyncDisposable.</returns>
73+
public IAsyncProcessor<TOutput> ProcessInParallel(int permitsPerWindow, TimeSpan window, int maxConcurrency)
74+
{
75+
return new ResultTimedRateLimitedParallelAsyncProcessor<TInput, TOutput>(_items, _taskSelector, permitsPerWindow, window, maxConcurrency, _cancellationTokenSource).StartProcessing();
6476
}
6577

6678
/// <summary>
@@ -137,4 +149,4 @@ public IAsyncProcessor<TOutput> ProcessOneAtATime()
137149
return new ResultOneAtATimeAsyncProcessor<TInput, TOutput>(_items, _taskSelector, _cancellationTokenSource).StartProcessing();
138150
}
139151

140-
}
152+
}

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/RateLimitedParallelAsyncProcessor.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,6 @@ internal RateLimitedParallelAsyncProcessor(int count, Func<Task> taskSelector, i
1616

1717
internal override Task Process()
1818
{
19-
return WorkerPool.ProcessAsync(TaskWrappers, _levelsOfParallelism, minimumIterationTime: null, CancellationToken);
19+
return WorkerPool.ProcessAsync(TaskWrappers, _levelsOfParallelism, rateLimiter: null, CancellationToken);
2020
}
2121
}

0 commit comments

Comments
 (0)