NetTaskPipeline 0.14.0

There is a newer version of this package available.
See the version list below for details.
dotnet add package NetTaskPipeline --version 0.14.0
                    
NuGet\Install-Package NetTaskPipeline -Version 0.14.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="0.14.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="NetTaskPipeline" Version="0.14.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 0.14.0
                    
#r "nuget: NetTaskPipeline, 0.14.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@0.14.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=0.14.0
                    
Install as a Cake Addin
#tool nuget:?package=NetTaskPipeline&version=0.14.0
                    
Install as a Cake Tool

NetTaskPipeline: Async Task Pipeline for .NET

NuGet examples build and tests coverage

A lightweight async task pipeline for .NET with sequential and parallel execution support.

👉 Supports Dependency Injection through fluent service provider registration

Features

  • Sequential task execution
  • Parallel task groups
  • Generic task registration with AddTask<TTask>()
  • Inline task registration with delegate-based AddTask(...) overloads
  • RabbitMQ RPC task execution with delegate-based AddTaskRpc<TRequest, TResponse>(...)
  • HTTP task execution with delegate-based AddTaskHttp<TRequest, TResponse>(...)
  • Fluent context-based branching
  • Shared execution context
  • Cancellation support
  • Retry support
  • CI-enforced minimum 80% line coverage
  • Timeout support
  • Error handling modes
  • Execution result reporting
  • Maximum degree of parallelism for parallel groups
  • Runnable simple, advanced, branching, HTTP, 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, register the service provider with WithServiceProvider(...) once and use the generic task methods normally. Do not pass the service provider to AddTask.

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();

Tasks are resolved internally using:

ActivatorUtilities.GetServiceOrCreateInstance(serviceProvider, taskType)

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, RabbitMQ RPC tasks, and HTTP 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 RabbitMQ RPC tasks

Use AddTaskRpc<TRequest, TResponse>(requestFactory, configure) when the pipeline needs to send a typed RabbitMQ RPC request.

using NetTaskPipeline;

var context = new 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);

The RabbitMQ RPC response is deserialized as TResponse and stored automatically using the configured ResponseKey.

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

Context usage with HTTP tasks

Use AddTaskHttp<TRequest, TResponse>(requestFactory, configure) when the pipeline needs to send a typed HTTP request and store the typed response in the shared context.

using System.Net.Http;
using NetTaskPipeline;

var context = new TaskContext();
context.Set("CustomerId", 123);

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";
            options.Headers["x-api-key"] = "secret";
        })
    .ExecuteAsync(context);

The HTTP response is deserialized as TResponse and stored automatically using the configured ResponseKey.

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

When the endpoint returns plain text, use string as the response type.

var result = await new TaskPipeline()
    .AddTaskHttp<object, string>(
        _ => new object(),
        options =>
        {
            options.RequestUri = "https://api.example.com/health";
            options.Method = HttpMethod.Get;
            options.ResponseKey = "HealthStatus";
        })
    .ExecuteAsync();

var healthStatus = result.Context.Get<string>("HealthStatus");

See examples/HttpExample for a runnable example that calls two public APIs.

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:

  • StopOnFirstError
  • ContinueOnError

Retry

await new TaskPipeline()
    .WithRetry(3)
    .AddTask<CallExternalApiTask>()
    .ExecuteAsync();

Configure an optional delay between retries. Exponential backoff keeps the first delay unchanged and doubles it for each subsequent retry:

await new TaskPipeline()
    .WithRetry(3)
    .WithRetryDelay(TimeSpan.FromMilliseconds(250), exponentialBackoff: true)
    .AddTask<CallExternalApiTask>()
    .ExecuteAsync();

Without WithRetryDelay, retries remain immediate for backward compatibility.

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. GitHub Actions runs all examples on pushes and pull requests to main, including the RabbitMQ RPC Docker example.

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/HttpExample/HttpExample.csproj

The RabbitMQ RPC Docker example can be started with Docker Compose:

cd examples/RpcDockerExample
docker compose up --build

Development

Pull requests must keep line coverage at or above 80% for unit-testable code. RabbitMQ transport/RPC files are integration-bound and excluded from the unit coverage gate; their integration validation remains separate. Code changes should update affected tests, examples, and this README when behavior or public APIs change. Project documentation is intentionally kept in this main README.

Dependency Injection example

Run the dependency injection example from the repository root:

dotnet run --project examples/DependencyInjectionExample/DependencyInjectionExample.csproj

Register task dependencies with Microsoft.Extensions.DependencyInjection, build the service provider, and configure the pipeline once with WithServiceProvider(serviceProvider). Typed tasks are then resolved through ActivatorUtilities.GetServiceOrCreateInstance, including tasks inside parallel groups and branches.

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 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. 
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