Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ public static ItemActionAsyncProcessorBuilder<T> ForEachAsync<T>(this IEnumerabl
internal static async IAsyncEnumerable<T> ToIAsyncEnumerable<T>(this IEnumerable<Task<T>> 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;
}
Expand All @@ -36,9 +36,9 @@ internal static async IAsyncEnumerable<T> ToIAsyncEnumerable<T>(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
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<TInput> 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)
Expand All @@ -82,21 +82,21 @@ private async Task ConsumeAsync(ChannelReader<TInput> 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
{
// 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);
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -50,7 +50,7 @@ public async Task ExecuteAsync()
}
}

await Task.WhenAll(tasks);
await Task.WhenAll(tasks).ConfigureAwait(false);
}
finally
{
Expand All @@ -62,7 +62,7 @@ private async Task ProcessItemAsync(TInput item, SemaphoreSlim semaphore, Cancel
{
try
{
await _taskSelector(item);
await _taskSelector(item).ConfigureAwait(false);
}
finally
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -32,15 +32,15 @@ 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);
tasks.Add(task);
}

// Wait for all tasks to complete
await Task.WhenAll(tasks);
await Task.WhenAll(tasks).ConfigureAwait(false);
}
}
#endif
Original file line number Diff line number Diff line change
Expand Up @@ -32,14 +32,14 @@ public async IAsyncEnumerable<TOutput> 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;
}
Expand All @@ -62,18 +62,18 @@ private async IAsyncEnumerable<TOutput> 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<TOutput> ExecuteWithOrderPreservationAsync(CancellationToken cancellationToken)
Expand All @@ -90,30 +90,30 @@ private async IAsyncEnumerable<TOutput> 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);

// Start ordered consumer
var consumerTasks = new List<Task>();
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
var yieldingTask = Task.Run(async () =>
{
while (!producerCompleted || nextYieldIndex < totalProduced || orderingDictionary.Count > 0)
{
await orderingLock.WaitAsync(cancellationToken);
await orderingLock.WaitAsync(cancellationToken).ConfigureAwait(false);
TaskCompletionSource<TOutput>? tcs = null;
var found = false;

Expand All @@ -127,29 +127,29 @@ private async IAsyncEnumerable<TOutput> 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)
{
break;
}
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(
Expand All @@ -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<TOutput>();
Expand Down Expand Up @@ -218,9 +218,9 @@ private async Task ProduceAsync(ChannelWriter<TInput> 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
Expand All @@ -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
Expand All @@ -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)
Expand Down
Loading