NetTaskPipeline 0.4.0
See the version list below for details.
dotnet add package NetTaskPipeline --version 0.4.0
NuGet\Install-Package NetTaskPipeline -Version 0.4.0
<PackageReference Include="NetTaskPipeline" Version="0.4.0" />
<PackageVersion Include="NetTaskPipeline" Version="0.4.0" />
<PackageReference Include="NetTaskPipeline" />
paket add NetTaskPipeline --version 0.4.0
#r "nuget: NetTaskPipeline, 0.4.0"
#:package NetTaskPipeline@0.4.0
#addin nuget:?package=NetTaskPipeline&version=0.4.0
#tool nuget:?package=NetTaskPipeline&version=0.4.0
NetTaskPipeline: Async Task Pipeline for .NET
A lightweight async task pipeline for .NET with sequential and parallel execution support.
Features
- Sequential task execution
- Parallel task groups
- Generic task registration with
AddTask<TTask>() - RPC task execution with
AddTaskRpc(key) - Fluent context-based branching
- Shared execution context
- Cancellation support
- Retry support
- Timeout support
- Error handling modes
- Execution result reporting
- Maximum degree of parallelism for parallel groups
- Runnable simple and advanced examples
Basic usage
using NetTaskPipeline;
var result = await new TaskPipeline()
.OnError(ErrorMode.StopOnFirstError)
.WithRetry(2)
.WithTimeout(TimeSpan.FromSeconds(10))
.WithMaxDegreeOfParallelism(3)
.AddTask<ValidateCustomerTask>()
.AddParallel<GeneratePdfTask, SendEmailTask, SaveLogTask>()
.ExecuteAsync();
Console.WriteLine($"Pipeline success: {result.Success}");
Creating a task
using NetTaskPipeline;
public sealed class ValidateCustomerTask : ITask
{
public async Task ExecuteAsync(TaskContext context, CancellationToken cancellationToken = default)
{
await Task.Delay(500, cancellationToken);
context.Set("CustomerId", 123);
}
}
Generic task registration
Use AddTask<TTask>() to add a sequential task without creating the instance manually.
await new TaskPipeline()
.AddTask<ValidateCustomerTask>()
.AddTask<LoadCustomerTask>()
.ExecuteAsync();
Use AddParallel<TTask1, TTask2>() or AddParallel<TTask1, TTask2, TTask3>() to add a parallel task group by type.
await new TaskPipeline()
.AddTask<ValidateCustomerTask>()
.AddParallel<GeneratePdfTask, SendEmailTask, SaveLogTask>()
.AddTask<SaveOrderTask>()
.ExecuteAsync();
Generic task registration requires a public parameterless constructor because the pipeline creates the task instance internally.
public sealed class SaveOrderTask : ITask
{
public Task ExecuteAsync(TaskContext context, CancellationToken cancellationToken = default)
{
return Task.CompletedTask;
}
}
Tasks that require constructor dependencies are not supported by the generic registration API yet. Keep task dependencies in the shared TaskContext, or add a factory/DI integration layer on top of the pipeline.
Adding values to the shared context
Use context.Set(key, value) inside any task to add or replace a value in the shared pipeline context.
using NetTaskPipeline;
public sealed class LoadCustomerTask : ITask
{
public async Task ExecuteAsync(TaskContext context, CancellationToken cancellationToken = default)
{
await Task.Delay(500, cancellationToken);
context.Set("CustomerId", 123);
context.Set("CustomerName", "John Smith");
context.Set("CustomerRequest", new GetCustomerRequest { CustomerId = 123 });
context.Set("CustomerType", "premium");
}
}
Every task receives the same TaskContext instance during a pipeline execution, so values added by one task can be read by later tasks.
await new TaskPipeline()
.AddTask<LoadCustomerTask>()
.AddTask<SendCustomerEmailTask>()
.ExecuteAsync();
You can also create the context before executing the pipeline and pass initial values to it.
var context = new TaskContext();
context.Set("CorrelationId", Guid.NewGuid().ToString("N"));
context.Set("RequestedBy", "system");
var result = await new TaskPipeline()
.AddTask<LoadCustomerTask>()
.AddTask<SendCustomerEmailTask>()
.ExecuteAsync(context);
Reading from the shared context
Use context.Get<T>(key) when the value is required. It throws an exception if the key does not exist or if the value has a different type.
using NetTaskPipeline;
public sealed class SendEmailTask : ITask
{
public async Task ExecuteAsync(TaskContext context, CancellationToken cancellationToken = default)
{
var customerId = context.Get<int>("CustomerId");
await Task.Delay(1000, cancellationToken);
Console.WriteLine($"Email sent for customer {customerId}.");
}
}
Use context.TryGet<T>(key, out var value) when the value is optional.
public sealed class AuditTask : ITask
{
public Task ExecuteAsync(TaskContext context, CancellationToken cancellationToken = default)
{
if (context.TryGet<string>("CorrelationId", out var correlationId))
{
Console.WriteLine($"Correlation ID: {correlationId}");
}
return Task.CompletedTask;
}
}
RPC task
Use AddTaskRpc(key) to send an RPC request from the pipeline. The only parameter is the key used to get the outgoing request object from the shared TaskContext.
var context = new TaskContext();
context.Set("CustomerRequest", new GetCustomerRequest
{
CustomerId = 123
});
var result = await new TaskPipeline()
.AddTaskRpc("CustomerRequest")
.ExecuteAsync(context);
The RPC endpoint name is the same as the key. In the example above, the endpoint name is CustomerRequest.
The response is stored automatically using the same key plus Response.
var response = result.Context.Get<object>("CustomerRequestResponse");
By default, the RPC connection is read from the NET_TASK_PIPELINE_RPC_URI environment variable. If the variable is not set, the local development connection is used.
NET_TASK_PIPELINE_RPC_URI=amqp://guest:guest@localhost:5672/
public sealed class GetCustomerRequest
{
public int CustomerId { get; set; }
}
RPC consumer
The consumer must listen to a queue with the same name used in AddTaskRpc(key).
.AddTaskRpc("CustomerRequest")
For this call, the consumer must listen to:
CustomerRequest
The consumer receives a JSON request, processes it, and publishes a JSON response back to the reply queue sent by the RPC caller.
using System.Text;
using System.Text.Json;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
var factory = new ConnectionFactory
{
Uri = new Uri(
Environment.GetEnvironmentVariable("NET_TASK_PIPELINE_RPC_URI")
?? "amqp://guest:guest@localhost:5672/")
};
await using var connection = await factory.CreateConnectionAsync();
await using var channel = await connection.CreateChannelAsync();
const string queueName = "CustomerRequest";
await channel.QueueDeclareAsync(
queue: queueName,
durable: false,
exclusive: false,
autoDelete: false,
arguments: null);
await channel.BasicQosAsync(
prefetchSize: 0,
prefetchCount: 1,
global: false);
var consumer = new AsyncEventingBasicConsumer(channel);
consumer.ReceivedAsync += async (_, eventArgs) =>
{
string responseJson;
try
{
var requestJson = Encoding.UTF8.GetString(eventArgs.Body.ToArray());
var request = JsonSerializer.Deserialize<GetCustomerRequest>(requestJson)
?? throw new InvalidOperationException("Invalid request payload.");
var response = await ProcessCustomerAsync(request);
responseJson = JsonSerializer.Serialize(response);
}
catch (Exception ex)
{
responseJson = JsonSerializer.Serialize(new
{
error = true,
message = ex.Message
});
}
var responseBytes = Encoding.UTF8.GetBytes(responseJson);
var replyProperties = new BasicProperties
{
CorrelationId = eventArgs.BasicProperties.CorrelationId,
ContentType = "application/json"
};
await channel.BasicPublishAsync(
exchange: string.Empty,
routingKey: eventArgs.BasicProperties.ReplyTo!,
mandatory: false,
basicProperties: replyProperties,
body: responseBytes);
await channel.BasicAckAsync(
deliveryTag: eventArgs.DeliveryTag,
multiple: false);
};
await channel.BasicConsumeAsync(
queue: queueName,
autoAck: false,
consumer: consumer);
Console.WriteLine($"RPC consumer listening on queue '{queueName}'.");
Console.WriteLine("Press ENTER to stop.");
Console.ReadLine();
static Task<GetCustomerResponse> ProcessCustomerAsync(GetCustomerRequest request)
{
return Task.FromResult(new GetCustomerResponse
{
CustomerId = request.CustomerId,
Name = $"Customer {request.CustomerId}"
});
}
public sealed class GetCustomerRequest
{
public int CustomerId { get; set; }
}
public sealed class GetCustomerResponse
{
public int CustomerId { get; set; }
public string Name { get; set; } = string.Empty;
}
The response must preserve the original correlation id.
CorrelationId = eventArgs.BasicProperties.CorrelationId
The response must be published to the reply destination received in the request.
routingKey: eventArgs.BasicProperties.ReplyTo!
Fluent context-based branching
Use AddBranch when the pipeline needs to choose between multiple flows based on a string value from the shared TaskContext.
var context = new TaskContext();
context.Set("CustomerType", "premium");
var result = await new TaskPipeline()
.AddBranch(
ctx => ctx.Get<string>("CustomerType"),
branch => branch
.When<ApplyPremiumDiscountTask, SendPremiumEmailTask>("premium")
.When<ApplyStandardDiscountTask, SendStandardEmailTask>("standard")
.When<BlockOrderTask>("blocked")
.Default<ReviewCustomerManuallyTask>(),
name: "Customer type decision")
.AddTask<SaveOrderTask>()
.ExecuteAsync(context);
For a branch with a full sub-pipeline, use the When overload that receives a pipeline.
await new TaskPipeline()
.AddBranch(
ctx => ctx.Get<string>("CustomerType"),
branch => branch
.When("premium", premium => premium
.WithRetry(2)
.AddTask<ApplyPremiumDiscountTask>()
.AddParallel<SendPremiumEmailTask, NotifySalesTeamTask>())
.When("standard", standard => standard
.AddTask<ApplyStandardDiscountTask>())
.Default(fallback => fallback
.AddTask<ReviewCustomerManuallyTask>()))
.ExecuteAsync(context);
The branch selector can also be asynchronous.
await new TaskPipeline()
.AddBranch(
async (ctx, cancellationToken) =>
{
await Task.Delay(100, cancellationToken);
return ctx.Get<decimal>("Total") >= 1000m
? "high-value"
: "low-value";
},
branch => branch
.When<RequireManagerApprovalTask>("high-value")
.When<AutoApproveTask>("low-value"),
name: "Approval decision")
.ExecuteAsync(context);
The lower-level AddNamedBranch API is still available when you prefer to pass a dictionary explicitly.
Execution model
Each AddTask call creates one execution group.
await new TaskPipeline()
.AddTask<FirstTask>()
.AddParallel<SecondTask, ThirdTask>()
.AddTask<FourthTask>()
.ExecuteAsync();
Execution order:
FirstTask
↓
SecondTask + ThirdTask in parallel
↓
FourthTask
Error handling
await new TaskPipeline()
.OnError(ErrorMode.ContinueOnError)
.AddTask<FirstTask>()
.AddTask<SecondTask>()
.ExecuteAsync();
Available modes:
StopOnFirstErrorContinueOnError
Retry
await new TaskPipeline()
.WithRetry(3)
.AddTask<CallExternalApiTask>()
.ExecuteAsync();
Per-task retry:
await new TaskPipeline()
.AddTask<CallExternalApiTask>(retryCount: 3)
.ExecuteAsync();
Timeout
await new TaskPipeline()
.WithTimeout(TimeSpan.FromSeconds(30))
.AddTask<LongRunningTask>()
.ExecuteAsync();
Per-task timeout:
await new TaskPipeline()
.AddTask<LongRunningTask>(timeout: TimeSpan.FromSeconds(5))
.ExecuteAsync();
Results
TaskPipelineResult result = await pipeline.ExecuteAsync();
foreach (var taskResult in result.TaskResults)
{
Console.WriteLine($"{taskResult.TaskName}: {taskResult.Status} in {taskResult.Duration}");
}
Runnable examples
The repository includes two executable examples.
Simple example
Covers sequential and parallel execution with a final result summary.
dotnet run --project examples/SimpleExample/SimpleExample.csproj
Advanced example
Covers most pipeline features in a single runnable flow:
- Shared
TaskContext - Sequential execution
- Parallel task groups
- Fluent context-based branching
- Generic task registration
ContinueOnError- Global retry
- Per-task retry
- Global timeout
- Per-task timeout
- Maximum degree of parallelism
- Execution result report
dotnet run --project examples/AdvancedExample/AdvancedExample.csproj
License
MIT
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net5.0 was computed. net5.0-windows was computed. net6.0 was computed. net6.0-android was computed. net6.0-ios was computed. net6.0-maccatalyst was computed. net6.0-macos was computed. net6.0-tvos was computed. net6.0-windows was computed. net7.0 was computed. net7.0-android was computed. net7.0-ios was computed. net7.0-maccatalyst was computed. net7.0-macos was computed. net7.0-tvos was computed. net7.0-windows was computed. net8.0 is compatible. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 was computed. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 was computed. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
| .NET Core | netcoreapp2.0 was computed. netcoreapp2.1 was computed. netcoreapp2.2 was computed. netcoreapp3.0 was computed. netcoreapp3.1 was computed. |
| .NET Standard | netstandard2.0 is compatible. netstandard2.1 was computed. |
| .NET Framework | net461 was computed. net462 was computed. net463 was computed. net47 was computed. net471 was computed. net472 was computed. net48 was computed. net481 was computed. |
| MonoAndroid | monoandroid was computed. |
| MonoMac | monomac was computed. |
| MonoTouch | monotouch was computed. |
| Tizen | tizen40 was computed. tizen60 was computed. |
| Xamarin.iOS | xamarinios was computed. |
| Xamarin.Mac | xamarinmac was computed. |
| Xamarin.TVOS | xamarintvos was computed. |
| Xamarin.WatchOS | xamarinwatchos was computed. |
-
.NETStandard 2.0
- RabbitMQ.Client (>= 7.2.1)
- System.Text.Json (>= 10.0.6)
-
net8.0
- RabbitMQ.Client (>= 7.2.1)
- System.Text.Json (>= 10.0.6)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 1.0.0 | 88 | 9/22/2026 |
| 0.14.0 | 86 | 9/22/2026 |
| 0.13.0 | 113 | 4/29/2026 |
| 0.12.0 | 115 | 4/28/2026 |
| 0.11.0 | 108 | 4/28/2026 |
| 0.10.0 | 103 | 4/28/2026 |
| 0.8.1 | 114 | 4/27/2026 |
| 0.8.0 | 109 | 4/27/2026 |
| 0.7.0 | 107 | 4/27/2026 |
| 0.5.0 | 120 | 4/27/2026 |
| 0.4.0 | 112 | 4/26/2026 |
| 0.3.0 | 107 | 4/25/2026 |
| 0.2.0 | 114 | 4/25/2026 |
| 0.1.0 | 111 | 4/25/2026 |