Skip to content

Commit e1cb4b7

Browse files
committed
feat: add cancellation-aware selectors
Add token-aware SelectAsync and ForEachAsync overloads for enumerable, async-enumerable, item, and execution-count builders. Adapt each delegate once to the builder's linked processor token so CancelAll, disposal, and external cancellation interrupt in-flight work without changing tokenless paths.
1 parent 691ecf5 commit e1cb4b7

12 files changed

Lines changed: 293 additions & 12 deletions
Lines changed: 133 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,133 @@
1+
using System;
2+
using System.Collections.Generic;
3+
using System.Threading;
4+
using System.Threading.Tasks;
5+
using EnumerableAsyncProcessor.Builders;
6+
using EnumerableAsyncProcessor.Extensions;
7+
8+
namespace EnumerableAsyncProcessor.UnitTests;
9+
10+
public class CancellationAwareSelectorTests
11+
{
12+
[Test]
13+
public async Task CancelAll_Interrupts_InFlight_ForEachAsync_Selector(CancellationToken cancellationToken)
14+
{
15+
var started = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
16+
var interrupted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
17+
18+
await using var processor = new[] { 1 }
19+
.ForEachAsync((_, processorToken) => WaitUntilCanceledAsync(processorToken, started, interrupted))
20+
.ProcessInParallel();
21+
22+
await started.Task.WaitAsync(cancellationToken);
23+
processor.CancelAll();
24+
25+
await interrupted.Task.WaitAsync(cancellationToken);
26+
await Assert.ThrowsAsync<TaskCanceledException>(() => processor.WaitAsync());
27+
}
28+
29+
[Test]
30+
public async Task CancelAll_Interrupts_InFlight_SelectAsync_Selector(CancellationToken cancellationToken)
31+
{
32+
var started = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
33+
var interrupted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
34+
35+
await using var processor = new[] { 1 }
36+
.SelectAsync((item, processorToken) => WaitUntilCanceledAsync(item, processorToken, started, interrupted))
37+
.ProcessInParallel();
38+
39+
await started.Task.WaitAsync(cancellationToken);
40+
processor.CancelAll();
41+
42+
await interrupted.Task.WaitAsync(cancellationToken);
43+
await Assert.ThrowsAsync<TaskCanceledException>(() => processor.GetResultsAsync());
44+
}
45+
46+
[Test]
47+
public async Task CancelAll_Interrupts_ExecutionCount_Selector(CancellationToken cancellationToken)
48+
{
49+
var started = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
50+
var interrupted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
51+
52+
await using var processor = AsyncProcessorBuilder.WithExecutionCount(1)
53+
.ForEachAsync(processorToken => WaitUntilCanceledAsync(processorToken, started, interrupted))
54+
.ProcessInParallel();
55+
56+
await started.Task.WaitAsync(cancellationToken);
57+
processor.CancelAll();
58+
59+
await interrupted.Task.WaitAsync(cancellationToken);
60+
await Assert.ThrowsAsync<TaskCanceledException>(() => processor.WaitAsync());
61+
}
62+
63+
[Test]
64+
public async Task DisposeAsync_Interrupts_InFlight_Selector(CancellationToken cancellationToken)
65+
{
66+
var started = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
67+
var interrupted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
68+
69+
var processor = new[] { 1 }
70+
.ForEachAsync((_, processorToken) => WaitUntilCanceledAsync(processorToken, started, interrupted))
71+
.ProcessInParallel();
72+
73+
await started.Task.WaitAsync(cancellationToken);
74+
var disposeTask = processor.DisposeAsync().AsTask();
75+
76+
await interrupted.Task.WaitAsync(cancellationToken);
77+
await disposeTask.WaitAsync(cancellationToken);
78+
}
79+
80+
[Test]
81+
public async Task ExternalCancellation_Interrupts_AsyncEnumerable_Selector(CancellationToken cancellationToken)
82+
{
83+
using var cancellationTokenSource = new CancellationTokenSource();
84+
var started = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
85+
var interrupted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
86+
87+
var processor = GetItemsAsync()
88+
.ForEachAsync(
89+
(_, processorToken) => WaitUntilCanceledAsync(processorToken, started, interrupted),
90+
cancellationTokenSource.Token)
91+
.ProcessInParallel();
92+
93+
var executeTask = processor.ExecuteAsync();
94+
await started.Task.WaitAsync(cancellationToken);
95+
await cancellationTokenSource.CancelAsync();
96+
97+
await interrupted.Task.WaitAsync(cancellationToken);
98+
await Assert.ThrowsAsync<OperationCanceledException>(() => executeTask);
99+
}
100+
101+
private static async Task WaitUntilCanceledAsync(
102+
CancellationToken cancellationToken,
103+
TaskCompletionSource started,
104+
TaskCompletionSource interrupted)
105+
{
106+
started.TrySetResult();
107+
try
108+
{
109+
await Task.Delay(Timeout.InfiniteTimeSpan, cancellationToken);
110+
}
111+
catch (OperationCanceledException)
112+
{
113+
interrupted.TrySetResult();
114+
throw;
115+
}
116+
}
117+
118+
private static async Task<T> WaitUntilCanceledAsync<T>(
119+
T result,
120+
CancellationToken cancellationToken,
121+
TaskCompletionSource started,
122+
TaskCompletionSource interrupted)
123+
{
124+
await WaitUntilCanceledAsync(cancellationToken, started, interrupted);
125+
return result;
126+
}
127+
128+
private static async IAsyncEnumerable<int> GetItemsAsync()
129+
{
130+
yield return 1;
131+
await Task.CompletedTask;
132+
}
133+
}

EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder.cs

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,13 @@ public ActionAsyncProcessorBuilder(int count, Func<Task> taskSelector, Cancellat
1717
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
1818
}
1919

20+
public ActionAsyncProcessorBuilder(int count, Func<CancellationToken, Task> taskSelector, CancellationToken cancellationToken)
21+
{
22+
_count = count;
23+
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
24+
_taskSelector = () => taskSelector(_cancellationTokenSource.Token);
25+
}
26+
2027
public IAsyncProcessor ProcessInBatches(int batchSize)
2128
{
2229
return new BatchAsyncProcessor(batchSize, _count, _taskSelector, _cancellationTokenSource).StartProcessing();
@@ -77,4 +84,4 @@ public IAsyncProcessor ProcessOneAtATime()
7784
return new OneAtATimeAsyncProcessor(_count, _taskSelector, _cancellationTokenSource).StartProcessing();
7885
}
7986

80-
}
87+
}

EnumerableAsyncProcessor/Builders/ActionAsyncProcessorBuilder_1.cs

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,13 @@ internal ActionAsyncProcessorBuilder(int count, Func<Task<TOutput>> taskSelector
1717
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
1818
}
1919

20+
internal ActionAsyncProcessorBuilder(int count, Func<CancellationToken, Task<TOutput>> taskSelector, CancellationToken cancellationToken)
21+
{
22+
_count = count;
23+
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
24+
_taskSelector = () => taskSelector(_cancellationTokenSource.Token);
25+
}
26+
2027
public IAsyncProcessor<TOutput> ProcessInBatches(int batchSize)
2128
{
2229
return new ResultBatchAsyncProcessor<TOutput>(batchSize, _count, _taskSelector, _cancellationTokenSource).StartProcessing();
@@ -77,4 +84,4 @@ public IAsyncProcessor<TOutput> ProcessOneAtATime()
7784
return new ResultOneAtATimeAsyncProcessor<TOutput>(_count, _taskSelector, _cancellationTokenSource).StartProcessing();
7885
}
7986

80-
}
87+
}

EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_1.cs

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,16 @@ public AsyncEnumerableActionAsyncProcessorBuilder(
1919
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
2020
}
2121

22+
public AsyncEnumerableActionAsyncProcessorBuilder(
23+
IAsyncEnumerable<TInput> items,
24+
Func<TInput, CancellationToken, Task> taskSelector,
25+
CancellationToken cancellationToken)
26+
{
27+
_items = items;
28+
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
29+
_taskSelector = item => taskSelector(item, _cancellationTokenSource.Token);
30+
}
31+
2232
/// <summary>
2333
/// Process items in parallel without concurrency limits.
2434
/// </summary>
@@ -70,7 +80,6 @@ public IAsyncEnumerableProcessor ProcessInParallel(int? maxConcurrency, bool sch
7080
_items, _taskSelector, maxConcurrency, scheduleOnThreadPool, _cancellationTokenSource);
7181
}
7282

73-
7483
/// <summary>
7584
/// Process items one at a time (sequential processing).
7685
/// </summary>
@@ -91,4 +100,4 @@ public IAsyncEnumerableProcessor ProcessInBatches(int batchSize)
91100
_items, _taskSelector, batchSize, _cancellationTokenSource);
92101
}
93102

94-
}
103+
}

EnumerableAsyncProcessor/Builders/AsyncEnumerableActionAsyncProcessorBuilder_2.cs

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,16 @@ public AsyncEnumerableActionAsyncProcessorBuilder(
1919
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
2020
}
2121

22+
public AsyncEnumerableActionAsyncProcessorBuilder(
23+
IAsyncEnumerable<TInput> items,
24+
Func<TInput, CancellationToken, Task<TOutput>> taskSelector,
25+
CancellationToken cancellationToken)
26+
{
27+
_items = items;
28+
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
29+
_taskSelector = item => taskSelector(item, _cancellationTokenSource.Token);
30+
}
31+
2232
/// <summary>
2333
/// Process items in parallel without concurrency limits and return results.
2434
/// </summary>
@@ -70,7 +80,6 @@ public IAsyncEnumerableProcessor<TOutput> ProcessInParallel(int? maxConcurrency,
7080
_items, _taskSelector, maxConcurrency, scheduleOnThreadPool, _cancellationTokenSource);
7181
}
7282

73-
7483
/// <summary>
7584
/// Process items one at a time and return results in order.
7685
/// </summary>
@@ -91,4 +100,4 @@ public IAsyncEnumerableProcessor<TOutput> ProcessInBatches(int batchSize)
91100
_items, _taskSelector, batchSize, _cancellationTokenSource);
92101
}
93102

94-
}
103+
}

EnumerableAsyncProcessor/Builders/AsyncEnumerableAsyncProcessorBuilder.cs

Lines changed: 27 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,19 @@ public AsyncEnumerableActionAsyncProcessorBuilder<TInput, TOutput> SelectAsync<T
2222
return new AsyncEnumerableActionAsyncProcessorBuilder<TInput, TOutput>(_items, taskSelector, cancellationToken);
2323
}
2424

25+
public AsyncEnumerableActionAsyncProcessorBuilder<TInput, TOutput> SelectAsync<TOutput>(
26+
Func<TInput, CancellationToken, Task<TOutput>> taskSelector)
27+
{
28+
return SelectAsync(taskSelector, CancellationToken.None);
29+
}
30+
31+
public AsyncEnumerableActionAsyncProcessorBuilder<TInput, TOutput> SelectAsync<TOutput>(
32+
Func<TInput, CancellationToken, Task<TOutput>> taskSelector,
33+
CancellationToken cancellationToken)
34+
{
35+
return new AsyncEnumerableActionAsyncProcessorBuilder<TInput, TOutput>(_items, taskSelector, cancellationToken);
36+
}
37+
2538
public AsyncEnumerableActionAsyncProcessorBuilder<TInput> ForEachAsync(
2639
Func<TInput, Task> taskSelector)
2740
{
@@ -34,4 +47,17 @@ public AsyncEnumerableActionAsyncProcessorBuilder<TInput> ForEachAsync(
3447
{
3548
return new AsyncEnumerableActionAsyncProcessorBuilder<TInput>(_items, taskSelector, cancellationToken);
3649
}
37-
}
50+
51+
public AsyncEnumerableActionAsyncProcessorBuilder<TInput> ForEachAsync(
52+
Func<TInput, CancellationToken, Task> taskSelector)
53+
{
54+
return ForEachAsync(taskSelector, CancellationToken.None);
55+
}
56+
57+
public AsyncEnumerableActionAsyncProcessorBuilder<TInput> ForEachAsync(
58+
Func<TInput, CancellationToken, Task> taskSelector,
59+
CancellationToken cancellationToken)
60+
{
61+
return new AsyncEnumerableActionAsyncProcessorBuilder<TInput>(_items, taskSelector, cancellationToken);
62+
}
63+
}

EnumerableAsyncProcessor/Builders/ExecutionCountAsyncProcessorBuilder.cs

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,16 @@ public ActionAsyncProcessorBuilder<TOutput> SelectAsync<TOutput>(Func<Task<TOutp
1919
return new ActionAsyncProcessorBuilder<TOutput>(_count, taskSelector, cancellationToken);
2020
}
2121

22+
public ActionAsyncProcessorBuilder<TOutput> SelectAsync<TOutput>(Func<CancellationToken, Task<TOutput>> taskSelector)
23+
{
24+
return SelectAsync(taskSelector, CancellationToken.None);
25+
}
26+
27+
public ActionAsyncProcessorBuilder<TOutput> SelectAsync<TOutput>(Func<CancellationToken, Task<TOutput>> taskSelector, CancellationToken cancellationToken)
28+
{
29+
return new ActionAsyncProcessorBuilder<TOutput>(_count, taskSelector, cancellationToken);
30+
}
31+
2232
public ActionAsyncProcessorBuilder ForEachAsync(Func<Task> taskSelector)
2333
{
2434
return ForEachAsync(taskSelector, CancellationToken.None);
@@ -28,4 +38,14 @@ public ActionAsyncProcessorBuilder ForEachAsync(Func<Task> taskSelector, Cancell
2838
{
2939
return new ActionAsyncProcessorBuilder(_count, taskSelector, cancellationToken);
3040
}
31-
}
41+
42+
public ActionAsyncProcessorBuilder ForEachAsync(Func<CancellationToken, Task> taskSelector)
43+
{
44+
return ForEachAsync(taskSelector, CancellationToken.None);
45+
}
46+
47+
public ActionAsyncProcessorBuilder ForEachAsync(Func<CancellationToken, Task> taskSelector, CancellationToken cancellationToken)
48+
{
49+
return new ActionAsyncProcessorBuilder(_count, taskSelector, cancellationToken);
50+
}
51+
}

EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_1.cs

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,13 @@ public ItemActionAsyncProcessorBuilder(IEnumerable<TInput> items, Func<TInput,Ta
1717
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
1818
}
1919

20+
public ItemActionAsyncProcessorBuilder(IEnumerable<TInput> items, Func<TInput, CancellationToken, Task> taskSelector, CancellationToken cancellationToken)
21+
{
22+
_items = items;
23+
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
24+
_taskSelector = item => taskSelector(item, _cancellationTokenSource.Token);
25+
}
26+
2027
public IAsyncProcessor ProcessInBatches(int batchSize)
2128
{
2229
return new BatchAsyncProcessor<TInput>(batchSize, _items, _taskSelector, _cancellationTokenSource)
@@ -82,4 +89,4 @@ public IAsyncProcessor ProcessOneAtATime()
8289
.StartProcessing();
8390
}
8491

85-
}
92+
}

EnumerableAsyncProcessor/Builders/ItemActionAsyncProcessorBuilder_2.cs

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,13 @@ internal ItemActionAsyncProcessorBuilder(IEnumerable<TInput> items, Func<TInput,
1717
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
1818
}
1919

20+
internal ItemActionAsyncProcessorBuilder(IEnumerable<TInput> items, Func<TInput, CancellationToken, Task<TOutput>> taskSelector, CancellationToken cancellationToken)
21+
{
22+
_items = items;
23+
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
24+
_taskSelector = item => taskSelector(item, _cancellationTokenSource.Token);
25+
}
26+
2027
/// <summary>
2128
/// Processes items in batches of the specified size.
2229
/// </summary>
@@ -137,4 +144,4 @@ public IAsyncProcessor<TOutput> ProcessOneAtATime()
137144
return new ResultOneAtATimeAsyncProcessor<TInput, TOutput>(_items, _taskSelector, _cancellationTokenSource).StartProcessing();
138145
}
139146

140-
}
147+
}

EnumerableAsyncProcessor/Builders/ItemAsyncProcessorBuilder.cs

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,16 @@ public ItemActionAsyncProcessorBuilder<TInput, TOutput> SelectAsync<TOutput>(Fun
1919
return new ItemActionAsyncProcessorBuilder<TInput, TOutput>(_items, taskSelector, cancellationToken);
2020
}
2121

22+
public ItemActionAsyncProcessorBuilder<TInput, TOutput> SelectAsync<TOutput>(Func<TInput, CancellationToken, Task<TOutput>> taskSelector)
23+
{
24+
return SelectAsync(taskSelector, CancellationToken.None);
25+
}
26+
27+
public ItemActionAsyncProcessorBuilder<TInput, TOutput> SelectAsync<TOutput>(Func<TInput, CancellationToken, Task<TOutput>> taskSelector, CancellationToken cancellationToken)
28+
{
29+
return new ItemActionAsyncProcessorBuilder<TInput, TOutput>(_items, taskSelector, cancellationToken);
30+
}
31+
2232
public ItemActionAsyncProcessorBuilder<TInput> ForEachAsync(Func<TInput, Task> taskSelector)
2333
{
2434
return ForEachAsync(taskSelector, CancellationToken.None);
@@ -28,4 +38,14 @@ public ItemActionAsyncProcessorBuilder<TInput> ForEachAsync(Func<TInput, Task> t
2838
{
2939
return new ItemActionAsyncProcessorBuilder<TInput>(_items, taskSelector, cancellationToken);
3040
}
31-
}
41+
42+
public ItemActionAsyncProcessorBuilder<TInput> ForEachAsync(Func<TInput, CancellationToken, Task> taskSelector)
43+
{
44+
return ForEachAsync(taskSelector, CancellationToken.None);
45+
}
46+
47+
public ItemActionAsyncProcessorBuilder<TInput> ForEachAsync(Func<TInput, CancellationToken, Task> taskSelector, CancellationToken cancellationToken)
48+
{
49+
return new ItemActionAsyncProcessorBuilder<TInput>(_items, taskSelector, cancellationToken);
50+
}
51+
}

0 commit comments

Comments
 (0)