NetTaskPipeline 0.8.1
See the version list below for details.
dotnet add package NetTaskPipeline --version 0.8.1
NuGet\Install-Package NetTaskPipeline -Version 0.8.1
<PackageReference Include="NetTaskPipeline" Version="0.8.1" />
<PackageVersion Include="NetTaskPipeline" Version="0.8.1" />
<PackageReference Include="NetTaskPipeline" />
paket add NetTaskPipeline --version 0.8.1
#r "nuget: NetTaskPipeline, 0.8.1"
#:package NetTaskPipeline@0.8.1
#addin nuget:?package=NetTaskPipeline&version=0.8.1
#tool nuget:?package=NetTaskPipeline&version=0.8.1
NetTaskPipeline: Async Task Pipeline for .NET
A lightweight async task pipeline for .NET with sequential and parallel execution support.
👉 Supports Dependency Injection (constructor injection built-in)
Features
- Sequential task execution
- Parallel task groups
- Generic task registration with
AddTask<TTask>() - Inline task registration with delegate-based
AddTask(...)overloads - RPC task execution with
AddTaskRpc<TRequest, TResponse>(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, advanced, branching, and RabbitMQ RPC Docker 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 when NOT using dependency injection.
Inline task registration
Use the delegate-based AddTask(...) overloads when you need a small task without creating a dedicated ITask class.
var result = await new TaskPipeline()
.AddTask("Set customer", context =>
{
context.Set("CustomerId", 123);
return Task.CompletedTask;
})
.AddTask("Send notification", async (context, cancellationToken) =>
{
var customerId = context.Get<int>("CustomerId");
await Task.Delay(500, cancellationToken);
Console.WriteLine($"Notification sent for customer {customerId}.");
})
.ExecuteAsync();
Inline tasks support the same retry and timeout settings as regular tasks.
await new TaskPipeline()
.AddTask(
"Call external API",
async (_, cancellationToken) =>
{
await Task.Delay(500, cancellationToken);
},
retryCount: 3,
timeout: TimeSpan.FromSeconds(5))
.ExecuteAsync();
When no task name is provided, the default name is InlineTask.
await new TaskPipeline()
.AddTask(context =>
{
context.Set("StartedAt", DateTimeOffset.UtcNow);
return Task.CompletedTask;
})
.ExecuteAsync();
Dependency Injection
If your tasks require constructor dependencies, use the DI-enabled overloads:
.AddTask<MyTask>(serviceProvider)
.AddParallel<TaskA, TaskB>(serviceProvider)
Example:
var services = new ServiceCollection();
services.AddSingleton<ICustomerRepository, CustomerRepository>();
services.AddTransient<LoadCustomerTask>();
var serviceProvider = services.BuildServiceProvider();
await new TaskPipeline()
.AddTask<LoadCustomerTask>(serviceProvider)
.ExecuteAsync();
Tasks are resolved using:
ActivatorUtilities.GetServiceOrCreateInstance<TTask>(serviceProvider)
Shared context
Every pipeline execution uses a TaskContext. The same context instance is passed to each task, so values written by one task can be read by later tasks, branch selectors, and RPC tasks.
When ExecuteAsync() is called without arguments, the pipeline creates a new empty context. When the caller needs to provide initial data, create a TaskContext and pass it to ExecuteAsync(context).
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);
The final context is available through the pipeline result:
var correlationId = result.Context.Get<string>("CorrelationId");
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");
}
}
A later task can read the values written by LoadCustomerTask:
await new TaskPipeline()
.AddTask<LoadCustomerTask>()
.AddTask<SendCustomerEmailTask>()
.ExecuteAsync();
Reading values from the shared context
Use context.Get<T>(key) when a value is required. It returns the value using the expected type and throws if the key is missing or if the stored value is not compatible with T.
using NetTaskPipeline;
public sealed class SendCustomerEmailTask : ITask
{
public async Task ExecuteAsync(TaskContext context, CancellationToken cancellationToken = default)
{
var customerId = context.Get<int>("CustomerId");
var customerName = context.Get<string>("CustomerName");
await Task.Delay(1000, cancellationToken);
Console.WriteLine($"Email sent for customer {customerId} - {customerName}.");
}
}
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;
}
}
Context usage with branching
AddBranch can choose the next flow from a value stored in TaskContext. Every When option receives an Action<TaskPipeline> so each branch can configure one or more tasks.
var context = new TaskContext();
context.Set("CustomerType", "premium");
var result = await new TaskPipeline()
.AddBranch(
selector: ctx => ctx.Get<string>("CustomerType"),
configure: branch => branch
.When("premium", flow => flow
.AddTask<ApplyPremiumDiscountTask>()
.AddTask<SendPremiumEmailTask>())
.When("standard", flow => flow
.AddTask<ApplyStandardDiscountTask>()
.AddTask<SendStandardEmailTask>())
.When("blocked", flow => flow
.AddTask<BlockOrderTask>())
.Default(flow => flow
.AddTask<ReviewCustomerManuallyTask>()),
name: "Customer type decision")
.AddTask<SaveOrderTask>()
.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("high-value", flow => flow.AddTask<RequireManagerApprovalTask>())
.When("low-value", flow => flow.AddTask<AutoApproveTask>()),
name: "Approval decision")
.ExecuteAsync(context);
Context usage with RPC tasks
Use AddTaskRpc<TRequest, TResponse>(key) when the pipeline needs to send a typed RPC request. The key is used to read the request object from TaskContext and is also used as the RPC endpoint name.
var context = new TaskContext();
context.Set("CustomerRequest", new GetCustomerRequest
{
CustomerId = 123
});
var result = await new TaskPipeline()
.AddTaskRpc<GetCustomerRequest, GetCustomerResponse>("CustomerRequest")
.ExecuteAsync(context);
The response is deserialized as TResponse and stored automatically using the same key plus Response.
var response = result.Context.Get<GetCustomerResponse>("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; }
}
public sealed class GetCustomerResponse
{
public int CustomerId { get; set; }
public string Name { get; set; } = string.Empty;
}
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 runnable examples.
dotnet run --project examples/SimpleExample/SimpleExample.csproj
dotnet run --project examples/AdvancedExample/AdvancedExample.csproj
dotnet run --project examples/BranchingExample/BranchingExample.csproj
The RabbitMQ RPC Docker example can be started with Docker Compose:
cd examples/RpcDockerExample
docker compose up --build
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
- Microsoft.Extensions.DependencyInjection (>= 10.0.1)
- RabbitMQ.Client (>= 7.2.1)
- System.Text.Json (>= 10.0.6)
-
net8.0
- Microsoft.Extensions.DependencyInjection (>= 10.0.1)
- 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 | 87 | 9/22/2026 |
| 0.14.0 | 85 | 9/22/2026 |
| 0.13.0 | 112 | 4/29/2026 |
| 0.12.0 | 114 | 4/28/2026 |
| 0.11.0 | 107 | 4/28/2026 |
| 0.10.0 | 102 | 4/28/2026 |
| 0.8.1 | 113 | 4/27/2026 |
| 0.8.0 | 108 | 4/27/2026 |
| 0.7.0 | 106 | 4/27/2026 |
| 0.5.0 | 119 | 4/27/2026 |
| 0.4.0 | 111 | 4/26/2026 |
| 0.3.0 | 106 | 4/25/2026 |
| 0.2.0 | 113 | 4/25/2026 |
| 0.1.0 | 110 | 4/25/2026 |