From 5d1547b184702297d1f1e144fec873f54b02adb3 Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Sat, 9 Aug 2025 01:13:45 +0100 Subject: [PATCH] fix: Fix critical deadlock issues in async processors MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Fix TaskCompletionSource not being set on exceptions in ProcessOrderedItemAsync Previously, if _taskSelector threw an exception, no TaskCompletionSource was added to the ordering dictionary, causing infinite waits. Now always add TCS first, then set result/exception appropriately. - Replace GetAwaiter().GetResult() with Task.Run + timeout in Dispose methods Prevents synchronization context deadlocks in UI/ASP.NET contexts by running async disposal on thread pool with 30-second timeout. - Add proper exception handling to channel.Writer.Complete() calls Ensures channels are properly completed even when exceptions occur during producer/consumer task execution. Uses finally blocks to guarantee completion regardless of where exceptions are thrown. 🤖 Generated with [Claude Code](https://claude.ai/code) Co-Authored-By: Claude --- .claude/settings.local.json | 11 ++- .../Abstract/AbstractAsyncProcessorBase.cs | 12 ++- .../AsyncEnumerableParallelProcessor.cs | 11 ++- ...ultAsyncEnumerableChannelBasedProcessor.cs | 93 +++++++++++++++---- .../ResultAbstractAsyncProcessorBase.cs | 12 ++- 5 files changed, 112 insertions(+), 27 deletions(-) diff --git a/.claude/settings.local.json b/.claude/settings.local.json index 1a96a5b..6775c53 100644 --- a/.claude/settings.local.json +++ b/.claude/settings.local.json @@ -12,7 +12,16 @@ "Bash(timeout 30 dotnet run:*)", "Bash(.EnumerableAsyncProcessor.Example.exe)", "Bash(EnumerableAsyncProcessor.Example.exe)", - "Bash(del \"EnumerableAsyncProcessor.UnitTests\\ValidationTests.cs\")" + "Bash(del \"EnumerableAsyncProcessor.UnitTests\\ValidationTests.cs\")", + "Bash(mkdir:*)", + "Bash(git add:*)", + "Bash(dotnet build)", + "Bash(dotnet build:*)", + "Bash(dotnet test)", + "Bash(dotnet test:*)", + "Bash(rg:*)", + "Bash(find:*)", + "Bash(timeout 30 dotnet test --no-build --verbosity minimal)" ], "deny": [] } diff --git a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs index 827e4ee..f02468d 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessorBase.cs @@ -148,12 +148,16 @@ protected virtual ValueTask DisposeAsyncCore() public void Dispose() { - // Use async disposal with ConfigureAwait(false) to avoid deadlocks - // and add a timeout to prevent indefinite blocking + // Use Task.Run to avoid deadlocks by running async disposal on thread pool + // Add timeout to prevent indefinite blocking try { - var disposeTask = DisposeAsync().ConfigureAwait(false); - disposeTask.GetAwaiter().GetResult(); + var disposeTask = Task.Run(async () => await DisposeAsync().ConfigureAwait(false)); + if (!disposeTask.Wait(TimeSpan.FromSeconds(30))) + { + // Log warning if disposal times out, but don't throw + // as per IDisposable pattern + } } catch { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs index 4e50e9b..972dcc7 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/AsyncEnumerableParallelProcessor.cs @@ -37,10 +37,17 @@ public async Task ExecuteAsync() { await channel.Writer.WriteAsync(item, cancellationToken).ConfigureAwait(false); } + channel.Writer.Complete(); } - finally + catch (OperationCanceledException) { - channel.Writer.Complete(); + channel.Writer.TryComplete(); + throw; + } + catch (Exception ex) + { + channel.Writer.TryComplete(ex); + throw; } }, cancellationToken); diff --git a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableChannelBasedProcessor.cs b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableChannelBasedProcessor.cs index e0e3389..9a0e56f 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableChannelBasedProcessor.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/AsyncEnumerable/ResultProcessors/ResultAsyncEnumerableChannelBasedProcessor.cs @@ -62,9 +62,32 @@ private async IAsyncEnumerable ExecuteWithoutOrderPreservationAsync(Can // Complete output when all processing is done var completionTask = Task.Run(async () => { - await producerTask.ConfigureAwait(false); - await Task.WhenAll(consumerTasks).ConfigureAwait(false); - outputChannel.Writer.Complete(); + Exception? exception = null; + try + { + await producerTask.ConfigureAwait(false); + await Task.WhenAll(consumerTasks).ConfigureAwait(false); + } + catch (Exception ex) + { + exception = ex; + } + finally + { + if (exception != null) + { + outputChannel.Writer.TryComplete(exception); + } + else + { + outputChannel.Writer.TryComplete(); + } + } + + if (exception != null) + { + throw exception; + } }, cancellationToken); // Yield results as they complete @@ -138,7 +161,36 @@ private async IAsyncEnumerable ExecuteWithOrderPreservationAsync(Cancel await Task.Delay(10, cancellationToken).ConfigureAwait(false); } } - outputChannel.Writer.Complete(); + }, cancellationToken); + + // Ensure channel completion when all tasks finish + var channelCompletionTask = Task.Run(async () => + { + Exception? exception = null; + try + { + await yieldingTask.ConfigureAwait(false); + } + catch (Exception ex) + { + exception = ex; + } + finally + { + if (exception != null) + { + outputChannel.Writer.TryComplete(exception); + } + else + { + outputChannel.Writer.TryComplete(); + } + } + + if (exception != null) + { + throw exception; + } }, cancellationToken); // Yield from output channel @@ -160,21 +212,30 @@ private async Task ProcessOrderedItemAsync( SemaphoreSlim semaphore, CancellationToken cancellationToken) { + var tcs = new TaskCompletionSource(); + + await orderingLock.WaitAsync(cancellationToken).ConfigureAwait(false); + try + { + orderingDictionary[index] = tcs; + } + finally + { + orderingLock.Release(); + } + try { var result = await _taskSelector(item).ConfigureAwait(false); - - await orderingLock.WaitAsync(cancellationToken).ConfigureAwait(false); - try - { - var tcs = new TaskCompletionSource(); - tcs.SetResult(result); - orderingDictionary[index] = tcs; - } - finally - { - orderingLock.Release(); - } + tcs.TrySetResult(result); + } + catch (OperationCanceledException) + { + tcs.TrySetCanceled(); + } + catch (Exception ex) + { + tcs.TrySetException(ex); } finally { diff --git a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs index fdd756d..2fc608b 100644 --- a/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs +++ b/EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessorBase.cs @@ -155,12 +155,16 @@ protected virtual ValueTask DisposeAsyncCore() public void Dispose() { - // Use async disposal with ConfigureAwait(false) to avoid deadlocks - // and add a timeout to prevent indefinite blocking + // Use Task.Run to avoid deadlocks by running async disposal on thread pool + // Add timeout to prevent indefinite blocking try { - var disposeTask = DisposeAsync().ConfigureAwait(false); - disposeTask.GetAwaiter().GetResult(); + var disposeTask = Task.Run(async () => await DisposeAsync().ConfigureAwait(false)); + if (!disposeTask.Wait(TimeSpan.FromSeconds(30))) + { + // Log warning if disposal times out, but don't throw + // as per IDisposable pattern + } } catch {