NetTaskPipeline 1.0.0

dotnet add package NetTaskPipeline --version 1.0.0
                    
NuGet\Install-Package NetTaskPipeline -Version 1.0.0
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="NetTaskPipeline" Version="1.0.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="NetTaskPipeline" Version="1.0.0" />
                    
Directory.Packages.props
<PackageReference Include="NetTaskPipeline" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add NetTaskPipeline --version 1.0.0
                    
#r "nuget: NetTaskPipeline, 1.0.0"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package NetTaskPipeline@1.0.0
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=NetTaskPipeline&version=1.0.0
                    
Install as a Cake Addin
#tool nuget:?package=NetTaskPipeline&version=1.0.0
                    
Install as a Cake Tool

NetTaskPipeline

NuGet examples build and tests coverage

A lightweight async task pipeline targeting .NET 10 LTS and .NET Standard 2.0, with sequential and parallel execution, branching, retries, timeouts, Dependency Injection, HTTP tasks, and RabbitMQ RPC.

Installation

dotnet add package NetTaskPipeline

Quick start

using NetTaskPipeline;

var result = await new TaskPipeline()
    .AddTask("Load customer", context =>
    {
        context.Set("CustomerId", 123);
        return Task.CompletedTask;
    })
    .AddTask("Process customer", context =>
    {
        Console.WriteLine(context.Get<int>("CustomerId"));
        return Task.CompletedTask;
    })
    .ExecuteAsync();

Console.WriteLine($"Pipeline success: {result.Success}");

Capabilities

Capability API
Sequential tasks AddTask(...) / AddTask<TTask>()
Parallel tasks AddParallel<...>()
Branching AddBranch(...)
Shared state TaskContext
Retry WithRetry(...) / per-task retryCount
Retry delay/backoff WithRetryDelay(...)
Retry filtering WithRetryPolicy(...)
Timeout WithTimeout(...) / per-task timeout
Error handling OnError(...)
Dependency Injection WithServiceProvider(...)
HTTP AddTaskHttp<...>()
RabbitMQ RPC AddTaskRpc<...>()

Core execution

Each AddTask creates a sequential execution group. AddParallel runs tasks in the same group concurrently before the next group starts.

await new TaskPipeline()
    .AddTask<ValidateCustomerTask>()
    .AddParallel<GeneratePdfTask, SendEmailTask, SaveLogTask>()
    .AddTask<SaveOrderTask>()
    .ExecuteAsync();
ValidateCustomerTask
        ↓
GeneratePdfTask + SendEmailTask + SaveLogTask
        ↓
SaveOrderTask

A task implements ITask:

public sealed class ValidateCustomerTask : ITask
{
    public async Task ExecuteAsync(TaskContext context, CancellationToken cancellationToken = default)
    {
        await Task.Delay(500, cancellationToken);
        context.Set("CustomerId", 123);
    }
}

Generic registration requires a public parameterless constructor unless Dependency Injection is configured. Small operations can instead use the delegate-based AddTask(...) overloads shown in the quick start.

TaskContext

One TaskContext is shared by the complete pipeline. Use Set to add or replace data, Get<T> for required values, and TryGet<T> for optional values.

var context = new TaskContext();
context.Set("CorrelationId", Guid.NewGuid().ToString("N"));
context.Set("CustomerId", 123);

var result = await new TaskPipeline()
    .AddTask<LoadCustomerTask>()
    .AddTask("Audit", ctx =>
    {
        if (ctx.TryGet<string>("CorrelationId", out var correlationId))
            Console.WriteLine(correlationId);

        return Task.CompletedTask;
    })
    .ExecuteAsync(context);

var customerId = result.Context.Get<int>("CustomerId");

The final context is available from TaskPipelineResult.Context. GetOrAdd uses ConcurrentDictionary semantics: insertion is atomic, but its value factory may execute more than once under contention, so factories should not contain side effects.

Branching

Branches select a pipeline flow from the shared context.

await new TaskPipeline()
    .AddBranch(
        selector: ctx => ctx.Get<string>("CustomerType"),
        configure: branch => branch
            .When("premium", flow => flow.AddTask<ApplyPremiumDiscountTask>())
            .When("standard", flow => flow.AddTask<ApplyStandardDiscountTask>())
            .Default(flow => flow.AddTask<ReviewCustomerManuallyTask>()),
        name: "Customer type decision")
    .AddTask<SaveOrderTask>()
    .ExecuteAsync(context);

The selector also has an asynchronous overload that receives a CancellationToken. If no case matches and no Default flow is configured, the branch returns a failed execution result instead of silently succeeding.

Reliability

Configure retry, optional retry delay/backoff, timeout, and error handling globally:

await new TaskPipeline()
    .OnError(ErrorMode.StopOnFirstError)
    .WithRetry(3)
    .WithRetryDelay(TimeSpan.FromMilliseconds(250), exponentialBackoff: true)
    .WithTimeout(TimeSpan.FromSeconds(30))
    .AddTask<CallExternalApiTask>()
    .ExecuteAsync();

Without WithRetryDelay, retries are immediate. Exponential backoff is overflow-safe, and optional jitter can spread concurrent retries (WithRetryDelay(delay, exponentialBackoff: true, jitter: true)). Use WithRetryPolicy(exception => ...) to retry only selected failures. External cancellation is never retried; task timeouts remain failures and can be filtered by the retry policy. Timeouts are cooperative: the pipeline cancels the token at the configured deadline, so tasks should observe the supplied CancellationToken. If a task ignores cancellation, the pipeline waits for it to return and still records a timeout failure once it completes. Per-task settings override the pipeline defaults where supported:

await new TaskPipeline()
    .AddTask<CallExternalApiTask>(
        retryCount: 3,
        timeout: TimeSpan.FromSeconds(5))
    .ExecuteAsync();

Error modes are StopOnFirstError and ContinueOnError.

Integrations

HTTP

AddTaskHttp<TRequest, TResponse> sends a typed HTTP request and stores the typed response in TaskContext.

var result = await new TaskPipeline()
    .AddTaskHttp<GetCustomerRequest, GetCustomerResponse>(
        ctx => new GetCustomerRequest { CustomerId = ctx.Get<int>("CustomerId") },
        options =>
        {
            options.RequestUri = "https://api.example.com/customers";
            options.Method = HttpMethod.Post;
            options.ResponseKey = "CustomerResponse";
        })
    .ExecuteAsync(context);

var response = result.Context.Get<GetCustomerResponse>("CustomerResponse");

Use string as TResponse for plain-text responses.

RabbitMQ RPC

AddTaskRpc<TRequest, TResponse> publishes a typed request, waits for the correlated response, deserializes it, and stores it in TaskContext.

var result = await new TaskPipeline()
    .AddTaskRpc<GetCustomerRequest, GetCustomerResponse>(
        ctx => new GetCustomerRequest { CustomerId = ctx.Get<int>("CustomerId") },
        options =>
        {
            options.ConnectionUri = "amqp://guest:guest@localhost:5672/";
            options.RoutingKey = "CustomerRequest";
            options.ResponseKey = "CustomerResponse";
        })
    .ExecuteAsync(context);

var response = result.Context.Get<GetCustomerResponse>("CustomerResponse");

Dependency Injection

Configure the service provider once. Typed tasks are then resolved from it, including tasks inside parallel groups and branches.

using Microsoft.Extensions.DependencyInjection;
using NetTaskPipeline;

var services = new ServiceCollection();
services.AddSingleton<ICustomerRepository, CustomerRepository>();
services.AddTransient<LoadCustomerTask>();
services.AddTransient<SendCustomerNotificationTask>();

using var serviceProvider = services.BuildServiceProvider();

await new TaskPipeline()
    .WithServiceProvider(serviceProvider)
    .AddTask<LoadCustomerTask>()
    .AddTask<SendCustomerNotificationTask>()
    .ExecuteAsync();

Internally, task resolution uses ActivatorUtilities.GetServiceOrCreateInstance.

Results

ExecuteAsync returns a TaskPipelineResult containing the final context and task execution results. Execution result properties are read-only to consumers so recorded status, attempts, timing, and exceptions cannot be mutated after execution.

var result = await pipeline.ExecuteAsync();

foreach (var taskResult in result.TaskResults)
    Console.WriteLine($"{taskResult.TaskName}: {taskResult.Status} in {taskResult.Duration}");

Examples

Example Purpose
SimpleExample Basic sequential pipeline
AdvancedExample Advanced pipeline behavior
BranchingExample Context-based branching
DependencyInjectionExample Service-provider task resolution
HttpExample Typed HTTP tasks
RpcDockerExample RabbitMQ RPC with Docker Compose

Run the .NET examples from the repository root:

dotnet run --project examples/SimpleExample/SimpleExample.csproj
dotnet run --project examples/AdvancedExample/AdvancedExample.csproj
dotnet run --project examples/BranchingExample/BranchingExample.csproj
dotnet run --project examples/DependencyInjectionExample/DependencyInjectionExample.csproj
dotnet run --project examples/HttpExample/HttpExample.csproj

Run the RabbitMQ RPC example with Docker:

cd examples/RpcDockerExample
docker compose up --build

GitHub Actions executes all examples on pushes and pull requests to main and validates their expected results.

Development

Check Behavior
Build and unit tests Runs for pushes and pull requests to main
Unit coverage Requires at least 80% line coverage for unit-testable code
RabbitMQ integration Runs against a real RabbitMQ service and tests RPC round-trip and timeout behavior
Integration coverage Collected separately as Cobertura and uploaded as a workflow artifact
Examples All runnable examples must produce their expected results
Concurrency Superseded CI runs for the same branch are cancelled
NuGet release Published only from a GitHub Release tag such as v1.0.0 after the full CI and a packed-package smoke test succeed
Package debugging Source Link and .snupkg symbols are published with the package

RabbitMQ transport/RPC files are excluded from the unit coverage gate because they are integration-bound; they are covered separately by the RabbitMQ integration job. The ≥80% badge therefore represents the unit-testable-code gate, not aggregate coverage across every source file.

The project uses .NET 10 LTS for development, tests, CI, examples, and release validation while retaining netstandard2.0 as a library target for broad consumer compatibility.\n\nCode changes should keep tests and examples current. Project documentation is intentionally maintained only in this root README.md.

License

MIT

Product 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 was computed.  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 is compatible.  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. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

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 110 4/26/2026
0.3.0 105 4/25/2026
0.2.0 112 4/25/2026
0.1.0 109 4/25/2026

See the GitHub release for changes in this version.