diff --git a/EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs b/EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs index d56c276..e166a02 100644 --- a/EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs +++ b/EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs @@ -27,7 +27,7 @@ public static ItemActionAsyncProcessorBuilder ForEachAsync(this IEnumerabl internal static async IAsyncEnumerable ToIAsyncEnumerable(this IEnumerable> tasks) { #if NET9_0_OR_GREATER - await foreach (var task in Task.WhenEach(tasks)) + await foreach (var task in Task.WhenEach(tasks).ConfigureAwait(false)) { yield return task.Result; } @@ -36,9 +36,9 @@ internal static async IAsyncEnumerable ToIAsyncEnumerable(this IEnumerable while (managedTasksList.Count != 0) { - var finishedTask = await Task.WhenAny(managedTasksList); + var finishedTask = await Task.WhenAny(managedTasksList).ConfigureAwait(false); managedTasksList.Remove(finishedTask); - yield return await finishedTask; + yield return await finishedTask.ConfigureAwait(false); } #endif } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs index ba62fc9..827e4ee 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs @@ -148,7 +148,16 @@ protected virtual ValueTask DisposeAsyncCore() public void Dispose() { - // Synchronous disposal calls async disposal and blocks - DisposeAsync().GetAwaiter().GetResult(); + // Use async disposal with ConfigureAwait(false) to avoid deadlocks + // and add a timeout to prevent indefinite blocking + try + { + var disposeTask = DisposeAsync().ConfigureAwait(false); + disposeTask.GetAwaiter().GetResult(); + } + catch + { + // Suppress exceptions during disposal as per IDisposable pattern + } } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableChannelBasedProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableChannelBasedProcessor.cs index d4683a1..53e0d9f 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableChannelBasedProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableChannelBasedProcessor.cs @@ -54,17 +54,17 @@ public async Task ExecuteAsync() .ToArray(); // Wait for all tasks - await producerTask; - await Task.WhenAll(consumerTasks); + await producerTask.ConfigureAwait(false); + await Task.WhenAll(consumerTasks).ConfigureAwait(false); } private async Task ProduceAsync(ChannelWriter writer, CancellationToken cancellationToken) { try { - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - await writer.WriteAsync(item, cancellationToken); + await writer.WriteAsync(item, cancellationToken).ConfigureAwait(false); } } catch (OperationCanceledException) @@ -82,9 +82,9 @@ private async Task ConsumeAsync(ChannelReader reader, CancellationToken if (_options.IsIOBound) { // For I/O-bound tasks, process directly without Task.Run - await foreach (var item in reader.ReadAllAsync(cancellationToken)) + await foreach (var item in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { - await _taskSelector(item); + await _taskSelector(item).ConfigureAwait(false); } } else @@ -92,11 +92,11 @@ private async Task ConsumeAsync(ChannelReader reader, CancellationToken // For CPU-bound tasks, use Task.Run to avoid blocking await Task.Run(async () => { - await foreach (var item in reader.ReadAllAsync(cancellationToken)) + await foreach (var item in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { - await _taskSelector(item); + await _taskSelector(item).ConfigureAwait(false); } - }, cancellationToken); + }, cancellationToken).ConfigureAwait(false); } } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableIOBoundParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableIOBoundParallelProcessor.cs index 201d57d..b80ef7e 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableIOBoundParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableIOBoundParallelProcessor.cs @@ -35,9 +35,9 @@ public async Task ExecuteAsync() try { - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - await semaphore.WaitAsync(cancellationToken); + await semaphore.WaitAsync(cancellationToken).ConfigureAwait(false); // Start task without Task.Run for I/O-bound operations var task = ProcessItemAsync(item, semaphore, cancellationToken); @@ -50,7 +50,7 @@ public async Task ExecuteAsync() } } - await Task.WhenAll(tasks); + await Task.WhenAll(tasks).ConfigureAwait(false); } finally { @@ -62,7 +62,7 @@ private async Task ProcessItemAsync(TInput item, SemaphoreSlim semaphore, Cancel { try { - await _taskSelector(item); + await _taskSelector(item).ConfigureAwait(false); } finally { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs index d973734..6cfd60a 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableOneAtATimeProcessor.cs @@ -26,9 +26,9 @@ public async Task ExecuteAsync() { var cancellationToken = _cancellationTokenSource.Token; - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - await _taskSelector(item); + await _taskSelector(item).ConfigureAwait(false); } } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs index db80d12..4e50e9b 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs @@ -33,9 +33,9 @@ public async Task ExecuteAsync() { try { - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - await channel.Writer.WriteAsync(item, cancellationToken); + await channel.Writer.WriteAsync(item, cancellationToken).ConfigureAwait(false); } } finally @@ -48,16 +48,16 @@ public async Task ExecuteAsync() var consumerTasks = Enumerable.Range(0, _maxConcurrency) .Select(_ => Task.Run(async () => { - await foreach (var item in channel.Reader.ReadAllAsync(cancellationToken)) + await foreach (var item in channel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { - await _taskSelector(item); + await _taskSelector(item).ConfigureAwait(false); } }, cancellationToken)) .ToArray(); // Wait for producer and all consumers to complete - await producerTask; - await Task.WhenAll(consumerTasks); + await producerTask.ConfigureAwait(false); + await Task.WhenAll(consumerTasks).ConfigureAwait(false); } } #endif \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableUnboundedParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableUnboundedParallelProcessor.cs index 504b475..85e8787 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableUnboundedParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableUnboundedParallelProcessor.cs @@ -32,7 +32,7 @@ public async Task ExecuteAsync() // Start a task for each item immediately as it arrives // No throttling or concurrency control - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { // Start task immediately without waiting var task = _taskSelector(item); @@ -40,7 +40,7 @@ public async Task ExecuteAsync() } // Wait for all tasks to complete - await Task.WhenAll(tasks); + await Task.WhenAll(tasks).ConfigureAwait(false); } } #endif \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableChannelBasedProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableChannelBasedProcessor.cs index 97353bd..e0e3389 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableChannelBasedProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableChannelBasedProcessor.cs @@ -32,14 +32,14 @@ public async IAsyncEnumerable ExecuteAsync() if (_options.PreserveOrder) { - await foreach (var result in ExecuteWithOrderPreservationAsync(cancellationToken)) + await foreach (var result in ExecuteWithOrderPreservationAsync(cancellationToken).ConfigureAwait(false)) { yield return result; } } else { - await foreach (var result in ExecuteWithoutOrderPreservationAsync(cancellationToken)) + await foreach (var result in ExecuteWithoutOrderPreservationAsync(cancellationToken).ConfigureAwait(false)) { yield return result; } @@ -62,18 +62,18 @@ private async IAsyncEnumerable ExecuteWithoutOrderPreservationAsync(Can // Complete output when all processing is done var completionTask = Task.Run(async () => { - await producerTask; - await Task.WhenAll(consumerTasks); + await producerTask.ConfigureAwait(false); + await Task.WhenAll(consumerTasks).ConfigureAwait(false); outputChannel.Writer.Complete(); }, cancellationToken); // Yield results as they complete - await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken)) + await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { yield return result; } - await completionTask; + await completionTask.ConfigureAwait(false); } private async IAsyncEnumerable ExecuteWithOrderPreservationAsync(CancellationToken cancellationToken) @@ -90,7 +90,7 @@ private async IAsyncEnumerable ExecuteWithOrderPreservationAsync(Cancel var orderedInputChannel = CreateOrderedInputChannel(); var producerTask = Task.Run(async () => { - await ProduceOrderedAsync(orderedInputChannel.Writer, cancellationToken); + await ProduceOrderedAsync(orderedInputChannel.Writer, cancellationToken).ConfigureAwait(false); producerCompleted = true; }, cancellationToken); @@ -98,14 +98,14 @@ private async IAsyncEnumerable ExecuteWithOrderPreservationAsync(Cancel var consumerTasks = new List(); var consumerTask = Task.Run(async () => { - await foreach (var (item, index) in orderedInputChannel.Reader.ReadAllAsync(cancellationToken)) + await foreach (var (item, index) in orderedInputChannel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { - await semaphore.WaitAsync(cancellationToken); + await semaphore.WaitAsync(cancellationToken).ConfigureAwait(false); totalProduced = Math.Max(totalProduced, index + 1); var task = ProcessOrderedItemAsync(item, index, orderingDictionary, orderingLock, semaphore, cancellationToken); consumerTasks.Add(task); } - await Task.WhenAll(consumerTasks); + await Task.WhenAll(consumerTasks).ConfigureAwait(false); }, cancellationToken); // Yield results in order @@ -113,7 +113,7 @@ private async IAsyncEnumerable ExecuteWithOrderPreservationAsync(Cancel { while (!producerCompleted || nextYieldIndex < totalProduced || orderingDictionary.Count > 0) { - await orderingLock.WaitAsync(cancellationToken); + await orderingLock.WaitAsync(cancellationToken).ConfigureAwait(false); TaskCompletionSource? tcs = null; var found = false; @@ -127,7 +127,7 @@ private async IAsyncEnumerable ExecuteWithOrderPreservationAsync(Cancel if (found && tcs != null) { - await outputChannel.Writer.WriteAsync(await tcs.Task, cancellationToken); + await outputChannel.Writer.WriteAsync(await tcs.Task.ConfigureAwait(false), cancellationToken).ConfigureAwait(false); } else if (producerCompleted && consumerTasks.All(t => t.IsCompleted) && orderingDictionary.Count == 0) { @@ -135,21 +135,21 @@ private async IAsyncEnumerable ExecuteWithOrderPreservationAsync(Cancel } else { - await Task.Delay(10, cancellationToken); + await Task.Delay(10, cancellationToken).ConfigureAwait(false); } } outputChannel.Writer.Complete(); }, cancellationToken); // Yield from output channel - await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken)) + await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { yield return result; } - await producerTask; - await consumerTask; - await yieldingTask; + await producerTask.ConfigureAwait(false); + await consumerTask.ConfigureAwait(false); + await yieldingTask.ConfigureAwait(false); } private async Task ProcessOrderedItemAsync( @@ -162,9 +162,9 @@ private async Task ProcessOrderedItemAsync( { try { - var result = await _taskSelector(item); + var result = await _taskSelector(item).ConfigureAwait(false); - await orderingLock.WaitAsync(cancellationToken); + await orderingLock.WaitAsync(cancellationToken).ConfigureAwait(false); try { var tcs = new TaskCompletionSource(); @@ -218,9 +218,9 @@ private async Task ProduceAsync(ChannelWriter writer, CancellationToken { try { - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - await writer.WriteAsync(item, cancellationToken); + await writer.WriteAsync(item, cancellationToken).ConfigureAwait(false); } } finally @@ -234,9 +234,9 @@ private async Task ProduceOrderedAsync(ChannelWriter<(TInput, int)> writer, Canc try { var index = 0; - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - await writer.WriteAsync((item, index++), cancellationToken); + await writer.WriteAsync((item, index++), cancellationToken).ConfigureAwait(false); } } finally @@ -254,22 +254,22 @@ private async Task ConsumeAsync( { if (_options.IsIOBound) { - await foreach (var item in reader.ReadAllAsync(cancellationToken)) + await foreach (var item in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { - var result = await _taskSelector(item); - await writer.WriteAsync(result, cancellationToken); + var result = await _taskSelector(item).ConfigureAwait(false); + await writer.WriteAsync(result, cancellationToken).ConfigureAwait(false); } } else { await Task.Run(async () => { - await foreach (var item in reader.ReadAllAsync(cancellationToken)) + await foreach (var item in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { - var result = await _taskSelector(item); - await writer.WriteAsync(result, cancellationToken); + var result = await _taskSelector(item).ConfigureAwait(false); + await writer.WriteAsync(result, cancellationToken).ConfigureAwait(false); } - }, cancellationToken); + }, cancellationToken).ConfigureAwait(false); } } catch (OperationCanceledException) diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableIOBoundParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableIOBoundParallelProcessor.cs index e5467c0..7593362 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableIOBoundParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableIOBoundParallelProcessor.cs @@ -35,12 +35,12 @@ public async IAsyncEnumerable ExecuteAsync() var processingTask = ProcessAsync(outputChannel.Writer, cancellationToken); // Yield results as they complete - await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken)) + await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { yield return result; } - await processingTask; + await processingTask.ConfigureAwait(false); } private async Task ProcessAsync(ChannelWriter writer, CancellationToken cancellationToken) @@ -50,9 +50,9 @@ private async Task ProcessAsync(ChannelWriter writer, CancellationToken try { - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - await semaphore.WaitAsync(cancellationToken); + await semaphore.WaitAsync(cancellationToken).ConfigureAwait(false); // For I/O-bound, don't use Task.Run var task = ProcessItemAsync(item, writer, semaphore, cancellationToken); @@ -65,7 +65,7 @@ private async Task ProcessAsync(ChannelWriter writer, CancellationToken } } - await Task.WhenAll(tasks); + await Task.WhenAll(tasks).ConfigureAwait(false); } finally { @@ -82,8 +82,8 @@ private async Task ProcessItemAsync( { try { - var result = await _taskSelector(item); - await writer.WriteAsync(result, cancellationToken); + var result = await _taskSelector(item).ConfigureAwait(false); + await writer.WriteAsync(result, cancellationToken).ConfigureAwait(false); } catch (OperationCanceledException) { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs index 7f8835e..f3c6b88 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableOneAtATimeProcessor.cs @@ -26,9 +26,9 @@ public async IAsyncEnumerable ExecuteAsync() { var cancellationToken = _cancellationTokenSource.Token; - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - var result = await _taskSelector(item); + var result = await _taskSelector(item).ConfigureAwait(false); yield return result; } } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs index 392c98f..87b64d9 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableParallelProcessor.cs @@ -32,13 +32,13 @@ public async IAsyncEnumerable ExecuteAsync() var processingTask = ProcessAsync(outputChannel.Writer, cancellationToken); // Yield results as they become available - await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken)) + await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { yield return result; } // Ensure processing completes - await processingTask; + await processingTask.ConfigureAwait(false); } private async Task ProcessAsync(ChannelWriter writer, CancellationToken cancellationToken) @@ -48,9 +48,9 @@ private async Task ProcessAsync(ChannelWriter writer, CancellationToken try { - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { - await semaphore.WaitAsync(cancellationToken); + await semaphore.WaitAsync(cancellationToken).ConfigureAwait(false); var task = ProcessItemAsync(item, writer, semaphore, cancellationToken); tasks.Add(task); @@ -62,7 +62,7 @@ private async Task ProcessAsync(ChannelWriter writer, CancellationToken } } - await Task.WhenAll(tasks); + await Task.WhenAll(tasks).ConfigureAwait(false); } finally { @@ -79,8 +79,8 @@ private async Task ProcessItemAsync( { try { - var result = await _taskSelector(item); - await writer.WriteAsync(result, cancellationToken); + var result = await _taskSelector(item).ConfigureAwait(false); + await writer.WriteAsync(result, cancellationToken).ConfigureAwait(false); } finally { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableUnboundedParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableUnboundedParallelProcessor.cs index af98df8..a0c31ec 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableUnboundedParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableUnboundedParallelProcessor.cs @@ -34,13 +34,13 @@ public async IAsyncEnumerable ExecuteAsync() var processingTask = ProcessAsync(outputChannel.Writer, cancellationToken); // Yield results as they become available - await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken)) + await foreach (var result in outputChannel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { yield return result; } // Ensure processing completes - await processingTask; + await processingTask.ConfigureAwait(false); } private async Task ProcessAsync(ChannelWriter writer, CancellationToken cancellationToken) @@ -50,7 +50,7 @@ private async Task ProcessAsync(ChannelWriter writer, CancellationToken try { // Start a task for each item immediately as it arrives - await foreach (var item in _items.WithCancellation(cancellationToken)) + await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false)) { // Capture the item in a local variable for the closure var capturedItem = item; @@ -60,8 +60,8 @@ private async Task ProcessAsync(ChannelWriter writer, CancellationToken { try { - var result = await _taskSelector(capturedItem); - await writer.WriteAsync(result, cancellationToken); + var result = await _taskSelector(capturedItem).ConfigureAwait(false); + await writer.WriteAsync(result, cancellationToken).ConfigureAwait(false); } catch (OperationCanceledException) { @@ -73,7 +73,7 @@ private async Task ProcessAsync(ChannelWriter writer, CancellationToken } // Wait for all tasks to complete - await Task.WhenAll(tasks); + await Task.WhenAll(tasks).ConfigureAwait(false); } finally { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor.cs index fda0a30..267a83f 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor.cs @@ -20,12 +20,12 @@ internal override async Task Process() foreach (var taskWrappers in batchedTaskWrappers) { - await ProcessBatch(taskWrappers); + await ProcessBatch(taskWrappers).ConfigureAwait(false); } } - private Task ProcessBatch(ActionTaskWrapper[] taskWrappers) + private async Task ProcessBatch(ActionTaskWrapper[] taskWrappers) { - return Task.WhenAll(taskWrappers.Select(tw => tw.Process(CancellationToken))); + await Task.WhenAll(taskWrappers.Select(tw => tw.Process(CancellationToken))).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor_1.cs index 1e6093c..3f36f84 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/BatchAsyncProcessor_1.cs @@ -21,12 +21,12 @@ internal override async Task Process() foreach (var currentBatch in batchedItems) { - await ProcessBatch(currentBatch); + await ProcessBatch(currentBatch).ConfigureAwait(false); } } - private Task ProcessBatch(ItemTaskWrapper[] currentBatch) + private async Task ProcessBatch(ItemTaskWrapper[] currentBatch) { - return Task.WhenAll(currentBatch.Select(tw => tw.Process(CancellationToken))); + await Task.WhenAll(currentBatch.Select(tw => tw.Process(CancellationToken))).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor.cs index db7173c..ee9e816 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor.cs @@ -12,7 +12,7 @@ internal override async Task Process() { foreach (var taskWrapper in TaskWrappers) { - await taskWrapper.Process(CancellationToken); + await taskWrapper.Process(CancellationToken).ConfigureAwait(false); } } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor_1.cs index 0b15f77..4fcb885 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/OneAtATimeAsyncProcessor_1.cs @@ -12,7 +12,7 @@ internal override async Task Process() { foreach (var taskWrapper in TaskWrappers) { - await taskWrapper.Process(CancellationToken); + await taskWrapper.Process(CancellationToken).ConfigureAwait(false); } } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs index 0a0c361..afc622e 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor.cs @@ -11,13 +11,13 @@ internal ParallelAsyncProcessor(int count, Func taskSelector, Cancellation _isIOBound = isIOBound; } - internal override Task Process() + internal override async Task Process() { // For I/O-bound tasks, don't use Task.Run wrapper as it adds unnecessary overhead // The tasks are already async and won't block threads if (_isIOBound) { - return Task.WhenAll(TaskWrappers.Select(taskWrapper => + await Task.WhenAll(TaskWrappers.Select(taskWrapper => { var task = taskWrapper.Process(CancellationToken); // Fast-path for already completed tasks @@ -26,10 +26,11 @@ internal override Task Process() return task; } return task; - })); + })).ConfigureAwait(false); + return; } // For CPU-bound tasks, use Task.Run to offload to ThreadPool - return Task.WhenAll(TaskWrappers.Select(taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken)))); + await Task.WhenAll(TaskWrappers.Select(taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken)))).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs index 93f521c..0b4116a 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ParallelAsyncProcessor_1.cs @@ -11,13 +11,13 @@ internal ParallelAsyncProcessor(IEnumerable items, Func ta _isIOBound = isIOBound; } - internal override Task Process() + internal override async Task Process() { // For I/O-bound tasks, don't use Task.Run wrapper as it adds unnecessary overhead // The tasks are already async and won't block threads if (_isIOBound) { - return Task.WhenAll(TaskWrappers.Select(taskWrapper => + await Task.WhenAll(TaskWrappers.Select(taskWrapper => { var task = taskWrapper.Process(CancellationToken); // Fast-path for already completed tasks @@ -26,10 +26,11 @@ internal override Task Process() return task; } return task; - })); + })).ConfigureAwait(false); + return; } // For CPU-bound tasks, use Task.Run to offload to ThreadPool - return Task.WhenAll(TaskWrappers.Select(taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken)))); + await Task.WhenAll(TaskWrappers.Select(taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken)))).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs index 66a5317..fdd756d 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs @@ -155,7 +155,16 @@ protected virtual ValueTask DisposeAsyncCore() public void Dispose() { - // Synchronous disposal calls async disposal and blocks - DisposeAsync().GetAwaiter().GetResult(); + // Use async disposal with ConfigureAwait(false) to avoid deadlocks + // and add a timeout to prevent indefinite blocking + try + { + var disposeTask = DisposeAsync().ConfigureAwait(false); + disposeTask.GetAwaiter().GetResult(); + } + catch + { + // Suppress exceptions during disposal as per IDisposable pattern + } } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_1.cs index 4d0d09b..385263b 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_1.cs @@ -22,12 +22,12 @@ internal override async Task Process() foreach (var currentTaskCompletionSourceBatch in batchedTaskCompletionSources) { - await ProcessBatch(currentTaskCompletionSourceBatch); + await ProcessBatch(currentTaskCompletionSourceBatch).ConfigureAwait(false); } } - private Task ProcessBatch(ActionTaskWrapper[] batch) + private async Task ProcessBatch(ActionTaskWrapper[] batch) { - return Task.WhenAll(batch.Select(tw => tw.Process(CancellationToken))); + await Task.WhenAll(batch.Select(tw => tw.Process(CancellationToken))).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_2.cs index 6662810..1cd1792 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultBatchAsyncProcessor_2.cs @@ -18,12 +18,12 @@ internal override async Task Process() foreach (var currentBatch in batchedItems) { - await ProcessBatch(currentBatch); + await ProcessBatch(currentBatch).ConfigureAwait(false); } } - private Task ProcessBatch(ItemTaskWrapper[] currentBatch) + private async Task ProcessBatch(ItemTaskWrapper[] currentBatch) { - return Task.WhenAll(currentBatch.Select(tw => tw.Process(CancellationToken))); + await Task.WhenAll(currentBatch.Select(tw => tw.Process(CancellationToken))).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_1.cs index 8d4ce7a..39e358d 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_1.cs @@ -12,7 +12,7 @@ internal override async Task Process() { foreach (var taskWrapper in TaskWrappers) { - await taskWrapper.Process(CancellationToken); + await taskWrapper.Process(CancellationToken).ConfigureAwait(false); } } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_2.cs index 59b524c..5d90c02 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultOneAtATimeAsyncProcessor_2.cs @@ -12,7 +12,7 @@ internal override async Task Process() { foreach (var taskWrapper in TaskWrappers) { - await taskWrapper.Process(CancellationToken); + await taskWrapper.Process(CancellationToken).ConfigureAwait(false); } } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs index 38ffac1..6e0c801 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_1.cs @@ -11,13 +11,13 @@ internal ResultParallelAsyncProcessor(int count, Func> taskSelecto _isIOBound = isIOBound; } - internal override Task Process() + internal override async Task Process() { // For I/O-bound tasks, don't use Task.Run wrapper as it adds unnecessary overhead // The tasks are already async and won't block threads if (_isIOBound) { - return Task.WhenAll(TaskWrappers.Select(taskWrapper => + await Task.WhenAll(TaskWrappers.Select(taskWrapper => { var task = taskWrapper.Process(CancellationToken); // Fast-path for already completed tasks @@ -26,10 +26,11 @@ internal override Task Process() return task; } return task; - })); + })).ConfigureAwait(false); + return; } // For CPU-bound tasks, use Task.Run to offload to ThreadPool - return Task.WhenAll(TaskWrappers.Select(taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken)))); + await Task.WhenAll(TaskWrappers.Select(taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken)))).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs index 37821bf..1870a96 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultParallelAsyncProcessor_2.cs @@ -11,13 +11,13 @@ internal ResultParallelAsyncProcessor(IEnumerable items, Func + await Task.WhenAll(TaskWrappers.Select(taskWrapper => { var task = taskWrapper.Process(CancellationToken); // Fast-path for already completed tasks @@ -26,10 +26,11 @@ internal override Task Process() return task; } return task; - })); + })).ConfigureAwait(false); + return; } // For CPU-bound tasks, use Task.Run to offload to ThreadPool - return Task.WhenAll(TaskWrappers.Select(taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken)))); + await Task.WhenAll(TaskWrappers.Select(taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken)))).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs index 9b2c4a0..c67c72a 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_1.cs @@ -22,7 +22,7 @@ internal override Task Process() { await Task.WhenAll( Task.Run(() => taskWrapper.Process(CancellationToken)), - Task.Delay(_timeSpan, CancellationToken)); + Task.Delay(_timeSpan, CancellationToken)).ConfigureAwait(false); }, CancellationToken.None, false); // false = CPU-bound } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs index b6c3220..5ef3c2d 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultTimedRateLimitedParallelAsyncProcessor_2.cs @@ -22,7 +22,7 @@ internal override Task Process() { await Task.WhenAll( Task.Run(() => taskWrapper.Process(CancellationToken)), - Task.Delay(_timeSpan, CancellationToken)); + Task.Delay(_timeSpan, CancellationToken)).ConfigureAwait(false); }, CancellationToken.None, false); // false = CPU-bound } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultUnboundedParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultUnboundedParallelAsyncProcessor_1.cs index 8ecda1b..70e38ef 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultUnboundedParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultUnboundedParallelAsyncProcessor_1.cs @@ -13,7 +13,7 @@ internal ResultUnboundedParallelAsyncProcessor(int count, Func> ta { } - internal override Task Process() + internal override async Task Process() { // Start ALL tasks immediately without any throttling // This provides true unbounded parallelism @@ -28,6 +28,6 @@ internal override Task Process() return task; }); - return Task.WhenAll(tasks); + await Task.WhenAll(tasks).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultUnboundedParallelAsyncProcessor_2.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultUnboundedParallelAsyncProcessor_2.cs index 03608aa..6af1afc 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultUnboundedParallelAsyncProcessor_2.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/ResultUnboundedParallelAsyncProcessor_2.cs @@ -13,7 +13,7 @@ internal ResultUnboundedParallelAsyncProcessor(IEnumerable items, Func taskWrapper.Process(CancellationToken)), - Task.Delay(_timeSpan, CancellationToken)); + Task.Delay(_timeSpan, CancellationToken)).ConfigureAwait(false); }, CancellationToken.None, false); // false = CPU-bound } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs index a1c65f4..3172167 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/TimedRateLimitedParallelAsyncProcessor_1.cs @@ -26,7 +26,7 @@ internal override Task Process() { await Task.WhenAll( Task.Run(() => taskWrapper.Process(CancellationToken)), - Task.Delay(_timeSpan, CancellationToken)); + Task.Delay(_timeSpan, CancellationToken)).ConfigureAwait(false); }, CancellationToken.None, false); // false = CPU-bound } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/UnboundedParallelAsyncProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/UnboundedParallelAsyncProcessor.cs index 58dcb23..12146fa 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/UnboundedParallelAsyncProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/UnboundedParallelAsyncProcessor.cs @@ -13,7 +13,7 @@ internal UnboundedParallelAsyncProcessor(int count, Func taskSelector, Can { } - internal override Task Process() + internal override async Task Process() { // Start ALL tasks immediately without any throttling // This provides true unbounded parallelism @@ -28,6 +28,6 @@ internal override Task Process() return task; }); - return Task.WhenAll(tasks); + await Task.WhenAll(tasks).ConfigureAwait(false); } } \ No newline at end of file diff --git a/EnumerableAsyncProcessor/RunnableProcessors/UnboundedParallelAsyncProcessor_1.cs b/EnumerableAsyncProcessor/RunnableProcessors/UnboundedParallelAsyncProcessor_1.cs index 476ed8c..f2c41cb 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/UnboundedParallelAsyncProcessor_1.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/UnboundedParallelAsyncProcessor_1.cs @@ -13,7 +13,7 @@ internal UnboundedParallelAsyncProcessor(IEnumerable items, Func