Skip to content

Commit bfe0bb2

Browse files
authored
+semver:minor - Merge pull request #306 from thomhurst/fix/deadlock-configureawait-issues
fix: Fix deadlock issues by adding ConfigureAwait(false)
2 parents 79e25cb + 58f705a commit bfe0bb2

33 files changed

Lines changed: 148 additions & 126 deletions

EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@ public static ItemActionAsyncProcessorBuilder<T> ForEachAsync<T>(this IEnumerabl
2727
internal static async IAsyncEnumerable<T> ToIAsyncEnumerable<T>(this IEnumerable<Task<T>> tasks)
2828
{
2929
#if NET9_0_OR_GREATER
30-
await foreach (var task in Task.WhenEach(tasks))
30+
await foreach (var task in Task.WhenEach(tasks).ConfigureAwait(false))
3131
{
3232
yield return task.Result;
3333
}
@@ -36,9 +36,9 @@ internal static async IAsyncEnumerable<T> ToIAsyncEnumerable<T>(this IEnumerable
3636

3737
while (managedTasksList.Count != 0)
3838
{
39-
var finishedTask = await Task.WhenAny(managedTasksList);
39+
var finishedTask = await Task.WhenAny(managedTasksList).ConfigureAwait(false);
4040
managedTasksList.Remove(finishedTask);
41-
yield return await finishedTask;
41+
yield return await finishedTask.ConfigureAwait(false);
4242
}
4343
#endif
4444
}

EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -148,7 +148,16 @@ protected virtual ValueTask DisposeAsyncCore()
148148

149149
public void Dispose()
150150
{
151-
// Synchronous disposal calls async disposal and blocks
152-
DisposeAsync().GetAwaiter().GetResult();
151+
// Use async disposal with ConfigureAwait(false) to avoid deadlocks
152+
// and add a timeout to prevent indefinite blocking
153+
try
154+
{
155+
var disposeTask = DisposeAsync().ConfigureAwait(false);
156+
disposeTask.GetAwaiter().GetResult();
157+
}
158+
catch
159+
{
160+
// Suppress exceptions during disposal as per IDisposable pattern
161+
}
153162
}
154163
}

EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableChannelBasedProcessor.cs

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -54,17 +54,17 @@ public async Task ExecuteAsync()
5454
.ToArray();
5555

5656
// Wait for all tasks
57-
await producerTask;
58-
await Task.WhenAll(consumerTasks);
57+
await producerTask.ConfigureAwait(false);
58+
await Task.WhenAll(consumerTasks).ConfigureAwait(false);
5959
}
6060

6161
private async Task ProduceAsync(ChannelWriter<TInput> writer, CancellationToken cancellationToken)
6262
{
6363
try
6464
{
65-
await foreach (var item in _items.WithCancellation(cancellationToken))
65+
await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false))
6666
{
67-
await writer.WriteAsync(item, cancellationToken);
67+
await writer.WriteAsync(item, cancellationToken).ConfigureAwait(false);
6868
}
6969
}
7070
catch (OperationCanceledException)
@@ -82,21 +82,21 @@ private async Task ConsumeAsync(ChannelReader<TInput> reader, CancellationToken
8282
if (_options.IsIOBound)
8383
{
8484
// For I/O-bound tasks, process directly without Task.Run
85-
await foreach (var item in reader.ReadAllAsync(cancellationToken))
85+
await foreach (var item in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
8686
{
87-
await _taskSelector(item);
87+
await _taskSelector(item).ConfigureAwait(false);
8888
}
8989
}
9090
else
9191
{
9292
// For CPU-bound tasks, use Task.Run to avoid blocking
9393
await Task.Run(async () =>
9494
{
95-
await foreach (var item in reader.ReadAllAsync(cancellationToken))
95+
await foreach (var item in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
9696
{
97-
await _taskSelector(item);
97+
await _taskSelector(item).ConfigureAwait(false);
9898
}
99-
}, cancellationToken);
99+
}, cancellationToken).ConfigureAwait(false);
100100
}
101101
}
102102
}

EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableIOBoundParallelProcessor.cs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -35,9 +35,9 @@ public async Task ExecuteAsync()
3535

3636
try
3737
{
38-
await foreach (var item in _items.WithCancellation(cancellationToken))
38+
await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false))
3939
{
40-
await semaphore.WaitAsync(cancellationToken);
40+
await semaphore.WaitAsync(cancellationToken).ConfigureAwait(false);
4141

4242
// Start task without Task.Run for I/O-bound operations
4343
var task = ProcessItemAsync(item, semaphore, cancellationToken);
@@ -50,7 +50,7 @@ public async Task ExecuteAsync()
5050
}
5151
}
5252

53-
await Task.WhenAll(tasks);
53+
await Task.WhenAll(tasks).ConfigureAwait(false);
5454
}
5555
finally
5656
{
@@ -62,7 +62,7 @@ private async Task ProcessItemAsync(TInput item, SemaphoreSlim semaphore, Cancel
6262
{
6363
try
6464
{
65-
await _taskSelector(item);
65+
await _taskSelector(item).ConfigureAwait(false);
6666
}
6767
finally
6868
{

EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,9 +26,9 @@ public async Task ExecuteAsync()
2626
{
2727
var cancellationToken = _cancellationTokenSource.Token;
2828

29-
await foreach (var item in _items.WithCancellation(cancellationToken))
29+
await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false))
3030
{
31-
await _taskSelector(item);
31+
await _taskSelector(item).ConfigureAwait(false);
3232
}
3333
}
3434
}

EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -33,9 +33,9 @@ public async Task ExecuteAsync()
3333
{
3434
try
3535
{
36-
await foreach (var item in _items.WithCancellation(cancellationToken))
36+
await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false))
3737
{
38-
await channel.Writer.WriteAsync(item, cancellationToken);
38+
await channel.Writer.WriteAsync(item, cancellationToken).ConfigureAwait(false);
3939
}
4040
}
4141
finally
@@ -48,16 +48,16 @@ public async Task ExecuteAsync()
4848
var consumerTasks = Enumerable.Range(0, _maxConcurrency)
4949
.Select(_ => Task.Run(async () =>
5050
{
51-
await foreach (var item in channel.Reader.ReadAllAsync(cancellationToken))
51+
await foreach (var item in channel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
5252
{
53-
await _taskSelector(item);
53+
await _taskSelector(item).ConfigureAwait(false);
5454
}
5555
}, cancellationToken))
5656
.ToArray();
5757

5858
// Wait for producer and all consumers to complete
59-
await producerTask;
60-
await Task.WhenAll(consumerTasks);
59+
await producerTask.ConfigureAwait(false);
60+
await Task.WhenAll(consumerTasks).ConfigureAwait(false);
6161
}
6262
}
6363
#endif

EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableUnboundedParallelProcessor.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,15 +32,15 @@ public async Task ExecuteAsync()
3232

3333
// Start a task for each item immediately as it arrives
3434
// No throttling or concurrency control
35-
await foreach (var item in _items.WithCancellation(cancellationToken))
35+
await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false))
3636
{
3737
// Start task immediately without waiting
3838
var task = _taskSelector(item);
3939
tasks.Add(task);
4040
}
4141

4242
// Wait for all tasks to complete
43-
await Task.WhenAll(tasks);
43+
await Task.WhenAll(tasks).ConfigureAwait(false);
4444
}
4545
}
4646
#endif

EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableChannelBasedProcessor.cs

Lines changed: 30 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -32,14 +32,14 @@ public async IAsyncEnumerable<TOutput> ExecuteAsync()
3232

3333
if (_options.PreserveOrder)
3434
{
35-
await foreach (var result in ExecuteWithOrderPreservationAsync(cancellationToken))
35+
await foreach (var result in ExecuteWithOrderPreservationAsync(cancellationToken).ConfigureAwait(false))
3636
{
3737
yield return result;
3838
}
3939
}
4040
else
4141
{
42-
await foreach (var result in ExecuteWithoutOrderPreservationAsync(cancellationToken))
42+
await foreach (var result in ExecuteWithoutOrderPreservationAsync(cancellationToken).ConfigureAwait(false))
4343
{
4444
yield return result;
4545
}
@@ -62,18 +62,18 @@ private async IAsyncEnumerable<TOutput> ExecuteWithoutOrderPreservationAsync(Can
6262
// Complete output when all processing is done
6363
var completionTask = Task.Run(async () =>
6464
{
65-
await producerTask;
66-
await Task.WhenAll(consumerTasks);
65+
await producerTask.ConfigureAwait(false);
66+
await Task.WhenAll(consumerTasks).ConfigureAwait(false);
6767
outputChannel.Writer.Complete();
6868
}, cancellationToken);
6969

7070
// Yield results as they complete
71-
await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken))
71+
await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
7272
{
7373
yield return result;
7474
}
7575

76-
await completionTask;
76+
await completionTask.ConfigureAwait(false);
7777
}
7878

7979
private async IAsyncEnumerable<TOutput> ExecuteWithOrderPreservationAsync(CancellationToken cancellationToken)
@@ -90,30 +90,30 @@ private async IAsyncEnumerable<TOutput> ExecuteWithOrderPreservationAsync(Cancel
9090
var orderedInputChannel = CreateOrderedInputChannel();
9191
var producerTask = Task.Run(async () =>
9292
{
93-
await ProduceOrderedAsync(orderedInputChannel.Writer, cancellationToken);
93+
await ProduceOrderedAsync(orderedInputChannel.Writer, cancellationToken).ConfigureAwait(false);
9494
producerCompleted = true;
9595
}, cancellationToken);
9696

9797
// Start ordered consumer
9898
var consumerTasks = new List<Task>();
9999
var consumerTask = Task.Run(async () =>
100100
{
101-
await foreach (var (item, index) in orderedInputChannel.Reader.ReadAllAsync(cancellationToken))
101+
await foreach (var (item, index) in orderedInputChannel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
102102
{
103-
await semaphore.WaitAsync(cancellationToken);
103+
await semaphore.WaitAsync(cancellationToken).ConfigureAwait(false);
104104
totalProduced = Math.Max(totalProduced, index + 1);
105105
var task = ProcessOrderedItemAsync(item, index, orderingDictionary, orderingLock, semaphore, cancellationToken);
106106
consumerTasks.Add(task);
107107
}
108-
await Task.WhenAll(consumerTasks);
108+
await Task.WhenAll(consumerTasks).ConfigureAwait(false);
109109
}, cancellationToken);
110110

111111
// Yield results in order
112112
var yieldingTask = Task.Run(async () =>
113113
{
114114
while (!producerCompleted || nextYieldIndex < totalProduced || orderingDictionary.Count > 0)
115115
{
116-
await orderingLock.WaitAsync(cancellationToken);
116+
await orderingLock.WaitAsync(cancellationToken).ConfigureAwait(false);
117117
TaskCompletionSource<TOutput>? tcs = null;
118118
var found = false;
119119

@@ -127,29 +127,29 @@ private async IAsyncEnumerable<TOutput> ExecuteWithOrderPreservationAsync(Cancel
127127

128128
if (found && tcs != null)
129129
{
130-
await outputChannel.Writer.WriteAsync(await tcs.Task, cancellationToken);
130+
await outputChannel.Writer.WriteAsync(await tcs.Task.ConfigureAwait(false), cancellationToken).ConfigureAwait(false);
131131
}
132132
else if (producerCompleted && consumerTasks.All(t => t.IsCompleted) && orderingDictionary.Count == 0)
133133
{
134134
break;
135135
}
136136
else
137137
{
138-
await Task.Delay(10, cancellationToken);
138+
await Task.Delay(10, cancellationToken).ConfigureAwait(false);
139139
}
140140
}
141141
outputChannel.Writer.Complete();
142142
}, cancellationToken);
143143

144144
// Yield from output channel
145-
await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken))
145+
await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
146146
{
147147
yield return result;
148148
}
149149

150-
await producerTask;
151-
await consumerTask;
152-
await yieldingTask;
150+
await producerTask.ConfigureAwait(false);
151+
await consumerTask.ConfigureAwait(false);
152+
await yieldingTask.ConfigureAwait(false);
153153
}
154154

155155
private async Task ProcessOrderedItemAsync(
@@ -162,9 +162,9 @@ private async Task ProcessOrderedItemAsync(
162162
{
163163
try
164164
{
165-
var result = await _taskSelector(item);
165+
var result = await _taskSelector(item).ConfigureAwait(false);
166166

167-
await orderingLock.WaitAsync(cancellationToken);
167+
await orderingLock.WaitAsync(cancellationToken).ConfigureAwait(false);
168168
try
169169
{
170170
var tcs = new TaskCompletionSource<TOutput>();
@@ -218,9 +218,9 @@ private async Task ProduceAsync(ChannelWriter<TInput> writer, CancellationToken
218218
{
219219
try
220220
{
221-
await foreach (var item in _items.WithCancellation(cancellationToken))
221+
await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false))
222222
{
223-
await writer.WriteAsync(item, cancellationToken);
223+
await writer.WriteAsync(item, cancellationToken).ConfigureAwait(false);
224224
}
225225
}
226226
finally
@@ -234,9 +234,9 @@ private async Task ProduceOrderedAsync(ChannelWriter<(TInput, int)> writer, Canc
234234
try
235235
{
236236
var index = 0;
237-
await foreach (var item in _items.WithCancellation(cancellationToken))
237+
await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false))
238238
{
239-
await writer.WriteAsync((item, index++), cancellationToken);
239+
await writer.WriteAsync((item, index++), cancellationToken).ConfigureAwait(false);
240240
}
241241
}
242242
finally
@@ -254,22 +254,22 @@ private async Task ConsumeAsync(
254254
{
255255
if (_options.IsIOBound)
256256
{
257-
await foreach (var item in reader.ReadAllAsync(cancellationToken))
257+
await foreach (var item in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
258258
{
259-
var result = await _taskSelector(item);
260-
await writer.WriteAsync(result, cancellationToken);
259+
var result = await _taskSelector(item).ConfigureAwait(false);
260+
await writer.WriteAsync(result, cancellationToken).ConfigureAwait(false);
261261
}
262262
}
263263
else
264264
{
265265
await Task.Run(async () =>
266266
{
267-
await foreach (var item in reader.ReadAllAsync(cancellationToken))
267+
await foreach (var item in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
268268
{
269-
var result = await _taskSelector(item);
270-
await writer.WriteAsync(result, cancellationToken);
269+
var result = await _taskSelector(item).ConfigureAwait(false);
270+
await writer.WriteAsync(result, cancellationToken).ConfigureAwait(false);
271271
}
272-
}, cancellationToken);
272+
}, cancellationToken).ConfigureAwait(false);
273273
}
274274
}
275275
catch (OperationCanceledException)

0 commit comments

Comments
 (0)