Pipeliner.Net 2.5.1

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

Pipeliner.Net <img src="Pipeliner.Net/Icon.png" alt="Pipeliner.Net Icon" width="40" />

Continous Integration Release

A strongly typed, high-performance .NET pipeline framework with a fluent API for orchestrating complex workflows with branching, resilience, and streaming support.

Why use Pipeliner.Net?

Pipeliner.Net is useful when you need to:

  • compose multiple business steps into one reusable workflow,
  • keep request flow strongly typed from input to output,
  • run synchronous and asynchronous operations in the same pipeline,
  • add resilience policies (like retries) around execution,
  • process data in batches,
  • branch, fork, and merge workflow paths,
  • halt workflows at explicit gates for manual review,
  • observe run and step events for diagnostics and profiling.

Typical use cases:

  • API request processing pipelines,
  • ETL and transformation workflows,
  • command and validation orchestration,
  • integration workflows that call multiple external systems,
  • background job processing with explicit, testable steps.

Installation

dotnet add package Pipeliner.Net

Quick start

var pipeline = Pipeline
    .For<string>()
    .Then<int>(value => int.Parse(value))
    .Branch(
        value => value >= 0,
        value => value,
        _ => 0)
    .Then<int>(value => value * 2)
    .ThenAsync<int>(async (value, cancellationToken) =>
    {
        await Task.Delay(25, cancellationToken);
        return value + 10;
    })
    .Build("Quick start workflow");

var result = await pipeline.RunAsync("50");
Console.WriteLine(result);
// 110

Core concepts

OperationPipeline<TParam, TResult>

OperationPipeline<TParam, TResult> is the runtime pipeline type produced by the builder.

Use it for execution:

  • Run(...) and RunAsync(...)
  • RunBatchAsync(...)

Pipeline.For<TInput>() builder

Pipeline.For<TInput>() is the entry point for pipeline composition and produces an OperationPipeline via Build().

Use it for all pipeline definitions in application code.

Fluent builder examples

1) Typed synchronous + asynchronous chain

var pipeline = Pipeline
    .For<string>()
    .Then<int>(Convert.ToInt32)
    .ThenAsync<int>(async (value, cancellationToken) =>
    {
        await Task.Delay(25, cancellationToken);
        return value + 5;
    })
    .Build("Parse and increment");

var result = await pipeline.RunAsync("10");
// 15

2) Factory-based step registration

public sealed class AddTaxStep : IPipelineStep<decimal, decimal>
{
    public ValueTask<decimal> ExecuteAsync(decimal input, CancellationToken cancellationToken = default) =>
        ValueTask.FromResult(input * 1.2m);
}

var pipeline = Pipeline
    .For<decimal>()
    .Then<AddTaxStep, decimal>(() => new AddTaxStep())
    .Build();

var total = pipeline.Run(100m);
// 120.0

3) Step-level retry policy

var pipeline = Pipeline
    .For<int>()
    .ThenAsync(
        async (value, cancellationToken) =>
        {
            await Task.Delay(5, cancellationToken);
            return value + 1;
        },
        StepExecutionOptions.WithPolicy(new RetryExecutionPolicy(3)))
    .Build();

4) Step-level concurrency and rate limiting

using System.Threading.RateLimiting;

var limiter = new TokenBucketRateLimiter(
    new TokenBucketRateLimiterOptions
    {
        TokenLimit = 100,
        TokensPerPeriod = 100,
        ReplenishmentPeriod = TimeSpan.FromMinutes(1),
        QueueLimit = 250,
        AutoReplenishment = true
    });

var pipeline = Pipeline
    .For<Order>()
    .ThenAsync(
        "Send to ERP",
        SendToErpAsync,
        StepExecutionOptions.Create(
            maxConcurrency: 4,
            rateLimiter: limiter))
    .Build();

A rejected rate-limit lease throws PipelineRateLimitRejectedException.

5) Saga compensation

var pipeline = Pipeline
    .For<CreateOrderCommand>()
    .ThenSaga(
        "Reserve inventory",
        ReserveInventoryAsync,
        (reservation, cancellationToken) => ReleaseInventoryAsync(reservation.Id, cancellationToken))
    .ThenSaga(
        "Capture payment",
        CapturePaymentAsync,
        (payment, cancellationToken) => RefundPaymentAsync(payment.Id, cancellationToken))
    .Build("Create order");

If a later step fails, completed saga compensations run in reverse order. If compensation itself fails, PipelineSagaCompensationException exposes the original pipeline exception and all compensation failures.

6) Per-run state

public sealed class OrderState
{
    public int Attempts { get; set; }
    public DateTimeOffset? ValidatedAt { get; set; }
}

var pipeline = Pipeline
    .For<Order>()
    .WithState(() => new OrderState())
    .ThenAsync("Validate", async (order, state, cancellationToken) =>
    {
        state.Attempts++;
        state.ValidatedAt = DateTimeOffset.UtcNow;
        return await ValidateAsync(order, cancellationToken);
    })
    .Then("Apply state", (order, state) => order with { Attempts = state.Attempts })
    .Build("Stateful order workflow");

State is created once per pipeline run, so concurrent executions do not share mutable state.

7) Context-aware steps

Use the PipelineExecutionContext overload when a step needs run metadata, pipeline metadata, step metadata, or checkpoint helpers.

var pipeline = Pipeline
    .For<Order>()
    .ThenAsync("Enrich", async (order, context, cancellationToken) =>
    {
        Console.WriteLine($"Run {context.RunId} executing {context.StepName}");
        return await EnrichOrderAsync(order, context.PipelineVersion, cancellationToken);
    })
    .Build("Order enrichment", "2.4.0");

The context exposes the run ID, pipeline ID, pipeline name, pipeline version, current step ID, current step name, current step kind, current attempt number, and recorded attempts for the run.

8) Halt gates

Use HaltWhen(...) when a workflow should stop at a controlled point instead of continuing downstream. This is useful for manual review, compliance holds, approval workflows, or external repair before retrying later.

var pipeline = Pipeline
    .For<Order>()
    .Then("Validate", ValidateOrder)
    .HaltWhen("Manual review", order => order.RiskScore > 80)
    .ThenAsync("Submit", SubmitOrderAsync)
    .Build("Order submission", "2.4.0");
PipelineRunOutcome<OrderSubmission> outcome =
    await pipeline.RunWithStatusAsync(order, cancellationToken);

if (outcome.IsHalted)
{
    Console.WriteLine($"Halted at {outcome.Halt!.HaltName} for run {outcome.Halt.RunId}");
    return;
}

OrderSubmission submission = outcome.Value!;

Use RunWithStatusAsync(...) when halt is an expected workflow outcome. It returns PipelineRunStatus.Completed, PipelineRunStatus.Halted, PipelineRunStatus.Failed, or PipelineRunStatus.Cancelled without requiring callers to catch an exception for controlled halts.

RunAsync(...) still throws PipelineHaltedException for callers that prefer the existing success/failure execution model. Observer hooks receive run and step halt events in both modes.

9) Pipeline-level execution policy

var pipeline = Pipeline
    .For<int>()
    .ThenAsync<int>((value, _) => ValueTask.FromResult(value + 1))
    .WithPolicy(new RetryExecutionPolicy(2))
    .Build();

10) Dynamic routing

var pipeline = Pipeline
    .For<Payment>()
    .RouteBy<PaymentMethod, PaymentResult>(
        "Payment method route",
        payment => payment.Method,
        routes => routes
            .When(PaymentMethod.Card, ChargeCard)
            .WhenAsync(PaymentMethod.BankTransfer, StartBankTransferAsync)
            .Default(payment => PaymentResult.Rejected(payment.Id)))
    .Build();

If no route matches and no default route is configured, PipelineRouteNotFoundException is thrown.

11) Branch and branch async

var pipeline = Pipeline
    .For<int>()
    .Branch(
        value => value >= 0,
        value => value,
        _ => 0)
    .BranchAsync(
        value => value > 100,
        (value, _) => ValueTask.FromResult($"large:{value}"),
        (value, _) => ValueTask.FromResult($"small:{value}"))
    .Build();

var label = await pipeline.RunAsync(150);
// large:150

12) How fork works

Fork(...) fans out the current pipeline value to multiple branch delegates and executes those branches concurrently. Every branch receives the same input value and returns the same branch output type.

var forkOnly = Pipeline
    .For<Order>()
    .Fork<OrderCheckResult>(
        (order, cancellationToken) => CheckInventoryAsync(order, cancellationToken),
        (order, cancellationToken) => CheckPaymentRiskAsync(order, cancellationToken),
        (order, cancellationToken) => CheckShippingAsync(order, cancellationToken))
    .Build("Order checks");

ForkExecutionResult<OrderCheckResult> result = await forkOnly.RunAsync(order);

The output of a fork step is ForkExecutionResult<TBranch>. It contains one ForkResult<TBranch> per branch, in the same order the branches were registered.

Each ForkResult<TBranch> contains:

  • Index: the branch index from the original registration order,
  • IsSuccess: whether that branch completed successfully,
  • Value: the branch output when successful,
  • Error: the branch exception when failed.

Fork execution has a deliberate failure model: one failed branch does not automatically fail the fork step. The fork captures every branch outcome first. The next Merge(...) step decides whether branch failures should fail the pipeline, be ignored, or allow the first successful branch to win.

This is useful when independent checks or data fetches can run at the same time:

  • call several external services using the same request,
  • enrich a record from multiple independent sources,
  • run independent validations and aggregate the results,
  • calculate competing projections and choose one,
  • fan out to optional providers where partial success is acceptable.

Fork(...) is different from ThenParallel(...):

  • Fork(...) runs several different branch delegates against one current value.
  • ThenParallel(...) runs one delegate across many items in the current value.

13) Fork + merge (custom reducer)

Use Merge(...) immediately after Fork(...) when the pipeline should continue with a single value. With the default CustomReducer behavior, failed branches are ignored and the merge delegate receives only successful branch values. If every branch failed, the merge step throws an AggregateException containing the branch failures.

var pipeline = Pipeline
    .For<decimal>()
    .Fork<decimal>(
        (amount, _) => ValueTask.FromResult(amount + 5m),
        (amount, _) => ValueTask.FromResult(amount * 1.08m),
        (amount, _) => ValueTask.FromResult(amount - 3m))
    .Merge<decimal, decimal>((results, _) => ValueTask.FromResult(results.Sum()), MergeStepOptions.CustomReducer())
    .Build("Price workflow");

var finalPrice = await pipeline.RunAsync(120m);
Console.WriteLine(finalPrice);
// 371.6

14) Built-in merge strategy: throw on any failure

Use MergeStepOptions.ThrowOnAnyFailure() when every branch is required. If any branch failed, the merge step throws an AggregateException; otherwise it returns all successful values in branch order.

var pipeline = Pipeline
    .For<int>()
    .Fork<int>(
        (value, _) => ValueTask.FromResult(value + 1),
        (_, _) => ValueTask.FromException<int>(new InvalidOperationException("branch failed")),
        (value, _) => ValueTask.FromResult(value + 3))
    .Merge<int, IReadOnlyList<int>>(
        (results, _) => ValueTask.FromResult<IReadOnlyList<int>>(results),
        MergeStepOptions.ThrowOnAnyFailure())
    .Build();

await Assert.ThrowsAsync<AggregateException>(() => pipeline.RunAsync(10));

15) Built-in merge strategy: ignore failures

Use MergeStepOptions.IgnoreFailures() when branches are optional. Failed branches are dropped, and the output is the successful values in branch order.

var pipeline = Pipeline
    .For<int>()
    .Fork<int>(
        (value, _) => ValueTask.FromResult(value + 1),
        (_, _) => ValueTask.FromException<int>(new InvalidOperationException("branch failed")),
        (value, _) => ValueTask.FromResult(value + 3))
    .Merge<int, IReadOnlyList<int>>(
        (results, _) => ValueTask.FromResult<IReadOnlyList<int>>(results),
        MergeStepOptions.IgnoreFailures())
    .Build();

var results = await pipeline.RunAsync(10);
// [11, 13]

16) Built-in merge strategy: take first

Use MergeStepOptions.TakeFirst() when branches represent fallback providers or competing strategies. The first successful branch in registration order becomes the merge output. If no branch succeeds, the merge step throws.

var pipeline = Pipeline
    .For<int>()
    .Fork<int>(
        (value, _) => ValueTask.FromResult(value + 1),
        (value, _) => ValueTask.FromResult(value + 2))
    .Merge<int, int>(
        (results, _) => ValueTask.FromResult(results[0]),
        MergeStepOptions.TakeFirst())
    .Build();

var first = await pipeline.RunAsync(10);
// 11

17) Parallel projection

var pipeline = Pipeline
    .For<int[]>()
    .ThenParallel<int, int>(
        (value, _) => ValueTask.FromResult(value * value),
        ParallelStepOptions.Create(4))
    .Build();

var squares = await pipeline.RunAsync([1, 2, 3, 4]);
// [1, 4, 9, 16]

Batch execution

Batch from memory

var pipeline = Pipeline
    .For<string>()
    .Then<int>(Convert.ToInt32)
    .Then<int>(value => value + 1)
    .Build();

var results = await pipeline.RunBatchAsync(new[] { "1", "2", "3" });
// [2, 3, 4]

Batch from IAsyncEnumerable<T>

static async IAsyncEnumerable<string> GetInputsAsync()
{
    yield return "10";
    await Task.Yield();
    yield return "20";
}

var pipeline = Pipeline
    .For<string>()
    .Then<int>(Convert.ToInt32)
    .Then<int>(value => value + 5)
    .Build();

await foreach (var item in pipeline.RunBatchAsync(GetInputsAsync()))
{
    Console.WriteLine(item);
}

Checkpoints and persistence

Checkpoints let a request-response pipeline persist the current value at explicit points in the workflow. They are useful for debugging production failures, auditing important intermediate states, preserving expensive transformation results, and preparing for manual recovery scenarios.

Checkpointing is opt-in:

  • add one or more .Checkpoint(...) calls,
  • configure persistence with .WithCheckpointing(...),
  • make sure the value at each checkpoint can be serialized with System.Text.Json.

Pipelines without checkpoints are unaffected. Only the value flowing through an explicit checkpoint needs to be JSON-serializable.

Checkpointing is currently for Pipeline.For<TInput>() request-response pipelines. Stream pipelines are excluded from durable execution v1.

In-memory checkpoints

InMemoryPipelineCheckpointStore is useful for tests, local diagnostics, and short-lived processes.

var store = new InMemoryPipelineCheckpointStore();

var pipeline = Pipeline
    .For<string>()
    .Then("Parse", int.Parse)
    .Checkpoint("After parse")
    .Then("Increment", value => value + 1)
    .WithCheckpointing(store)
    .Build("Checkpointed parse workflow");

var result = await pipeline.RunAsync("41");
// 42

var checkpoints = await store.LoadByPipelineAsync(pipeline.Id);

foreach (var checkpoint in checkpoints)
{
    Console.WriteLine($"{checkpoint.CheckpointName}: {checkpoint.PayloadJson}");
}

File-backed checkpoints

FilePipelineCheckpointStore stores checkpoint records as JSON files.

var store = new FilePipelineCheckpointStore("./checkpoints");

var pipeline = Pipeline
    .For<Order>()
    .Then("Validate", ValidateOrder)
    .Checkpoint("After validation")
    .ThenAsync("Submit to ERP", SubmitToErpAsync)
    .WithCheckpointing(store)
    .Build("Order submission", "2.4.0");

await pipeline.RunAsync(order, cancellationToken);

Each saved checkpoint includes:

  • run ID,
  • pipeline ID, name, and version,
  • checkpoint name,
  • checkpoint node ID,
  • payload type,
  • JSON payload,
  • creation timestamp,
  • pipeline version.

You can load checkpoints for a specific run or all checkpoints for a pipeline:

IReadOnlyList<PipelineCheckpoint> runCheckpoints =
    await store.LoadAsync(runId, cancellationToken);

IReadOnlyList<PipelineCheckpoint> pipelineCheckpoints =
    await store.LoadByPipelineAsync(pipeline.Id, cancellationToken);

Checkpoint failure behavior

Checkpoint persistence failures fail the pipeline run by default. This is the safest behavior for workflows where a checkpoint is part of the reliability contract.

var pipeline = Pipeline
    .For<Order>()
    .Then("Validate", ValidateOrder)
    .Checkpoint("After validation")
    .WithCheckpointing(
        store,
        PipelineCheckpointFailureBehavior.FailRun)
    .Build();

If checkpoint persistence is best-effort for your workflow, configure the pipeline to continue when checkpoint storage fails:

var pipeline = Pipeline
    .For<Order>()
    .Then("Validate", ValidateOrder)
    .Checkpoint("After validation")
    .WithCheckpointing(
        store,
        PipelineCheckpointFailureBehavior.Continue)
    .Build();

Custom checkpoint stores

Implement IPipelineCheckpointStore to persist checkpoints in a database, blob store, document database, queue, or another durable system.

public sealed class SqlPipelineCheckpointStore : IPipelineCheckpointStore
{
    private readonly string connectionString;

    public SqlPipelineCheckpointStore(string connectionString)
    {
        this.connectionString = connectionString;
    }

    public async ValueTask SaveAsync(
        PipelineCheckpoint checkpoint,
        CancellationToken cancellationToken = default)
    {
        // Insert checkpoint.RunId, checkpoint.PipelineId,
        // checkpoint.CheckpointName, checkpoint.PayloadType,
        // checkpoint.PayloadJson, and checkpoint.CreatedAt.
        await SaveCheckpointRowAsync(connectionString, checkpoint, cancellationToken);
    }

    public async ValueTask<IReadOnlyList<PipelineCheckpoint>> LoadAsync(
        string runId,
        CancellationToken cancellationToken = default)
    {
        // Return all checkpoints for one pipeline run, ordered by creation time.
        return await LoadCheckpointRowsForRunAsync(connectionString, runId, cancellationToken);
    }

    public async ValueTask<IReadOnlyList<PipelineCheckpoint>> LoadByPipelineAsync(
        string pipelineId,
        CancellationToken cancellationToken = default)
    {
        // Return all checkpoints for one pipeline definition, ordered by creation time.
        return await LoadCheckpointRowsForPipelineAsync(connectionString, pipelineId, cancellationToken);
    }
}

Then use it like any built-in store:

var store = new SqlPipelineCheckpointStore(connectionString);

var pipeline = Pipeline
    .For<OrderImport>()
    .Then("Normalize", NormalizeImport)
    .Checkpoint("After normalization")
    .ThenAsync("Persist", PersistImportAsync)
    .WithCheckpointing(PipelineCheckpointOptions.FailRun(store))
    .Build("Import workflow");

Checkpoint compatibility

Pipeline definitions are versioned. Pass a version when building a pipeline if checkpoints may outlive a deployment or need compatibility checks before recovery.

var pipeline = Pipeline
    .For<Order>()
    .Then("Validate", ValidateOrder)
    .Checkpoint("After validation")
    .ThenAsync("Submit", SubmitOrderAsync)
    .WithCheckpointing(store)
    .Build("Order submission", "2.4.0");

var checkpoints = await store.LoadByPipelineAsync(pipeline.Id, cancellationToken);
var latest = checkpoints.Last();

PipelineCheckpointCompatibilityReport report =
    pipeline.Describe().ValidateCheckpointCompatibility(latest);

if (!report.IsCompatible)
{
    foreach (var issue in report.Issues)
    {
        Console.WriteLine(issue);
    }
}

Compatibility validation checks the pipeline ID, pipeline version, checkpoint node ID, and checkpoint payload type against the current pipeline definition.

Stream execution with backpressure

Stream pipelines are for IAsyncEnumerable<T> sources where items arrive over time instead of as one request payload. They are useful when you want to keep processing code strongly typed while consuming a live or long-running source.

Typical use cases:

  • process events from a queue, broker, websocket, or change feed,
  • normalize telemetry or metrics as they arrive,
  • import large files without loading everything into memory,
  • enrich incoming records one item at a time,
  • group small events into database-friendly batches,
  • compute rolling time-window summaries,
  • apply bounded backpressure when consumers are slower than producers.

Stream pipelines are intentionally separate from durable checkpoint execution v1. Use them for live transformation and bounded buffering; use request-response pipelines with checkpoints when you need persisted recovery points.

Stream builder quick start

static async IAsyncEnumerable<string> ReadLinesAsync(
    [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
    while (!cancellationToken.IsCancellationRequested)
    {
        var line = await ReadNextLineAsync(cancellationToken);
        if (line is null)
            yield break;

        yield return line;
    }
}

var streamPipeline = Pipeline
    .StreamFor<string>()
    .Then("Trim", line => line.Trim())
    .Then("Parse", int.Parse)
    .ThenAsync("Enrich", async (value, cancellationToken) =>
    {
        var multiplier = await GetCurrentMultiplierAsync(cancellationToken);
        return value * multiplier;
    })
    .Build("Live line processor");

await foreach (var item in streamPipeline.RunStreamAsync(ReadLinesAsync(), cancellationToken))
{
    Console.WriteLine(item);
}

Each item flows through the configured steps independently. The pipeline starts yielding transformed output as soon as output is available; it does not wait for the source sequence to finish.

Event ingestion example

This shape works well for queue or broker consumers where each event should be validated, normalized, and sent downstream.

var ingestionPipeline = Pipeline
    .StreamFor<OrderCreated>()
    .WithBackpressure(BackpressureOptions.Create(512, BackpressureMode.Wait))
    .Then("Validate", ValidateOrderCreated)
    .Then("Normalize", NormalizeOrderCreated)
    .ThenAsync("Publish read model", PublishReadModelAsync)
    .Build("Order event ingestion");

await foreach (var published in ingestionPipeline.RunStreamAsync(orderEvents, cancellationToken))
{
    logger.LogInformation("Published read model {OrderId}", published.OrderId);
}

Use this when every item matters and the producer should slow down if the consumer cannot keep up.

Configure bounded channel backpressure

Backpressure controls what happens when the internal channel reaches capacity.

var streamPipeline = Pipeline
    .StreamFor<SensorReading>()
    .WithBackpressure(BackpressureOptions.Create(
        capacity: 1_000,
        mode: BackpressureMode.Wait))
    .Then("Normalize", NormalizeReading)
    .Build("Sensor normalization");

Available modes:

  • BackpressureMode.Wait: wait for capacity; best when every item must be processed.
  • BackpressureMode.DropNewest: reject the newest write when full; useful for lossy real-time feeds.
  • BackpressureMode.DropOldest: discard the oldest buffered item; useful when the newest value is most important.
  • BackpressureMode.DropWrite: drop the current write; useful when producers should never block.

Choose the mode based on data value:

  • orders, payments, audit events: usually Wait,
  • UI cursor positions, live gauges, frequent sensor updates: often DropOldest or DropWrite,
  • noisy telemetry where recent data matters more than completeness: often DropNewest or DropOldest.

Batch stream items for efficient writes

Batch(size, maxDelay) groups items by count and can also flush a partial batch after a maximum delay. This is useful for database inserts, API bulk calls, search indexing, or file writes.

var batchedPipeline = Pipeline
    .StreamFor<OrderCreated>()
    .Then("Normalize", NormalizeOrderCreated)
    .Batch(size: 100, maxDelay: TimeSpan.FromSeconds(5))
    .ThenAsync("Bulk import", async (orders, cancellationToken) =>
    {
        await orderRepository.InsertManyAsync(orders, cancellationToken);
        return new ImportResult(orders.Count);
    })
    .Build("Order import batches");

await foreach (var result in batchedPipeline.RunStreamAsync(events, cancellationToken))
{
    Console.WriteLine($"Imported {result.ImportedCount} orders");
}

Batching behavior:

  • if size is reached first, the batch flushes immediately,
  • if maxDelay is reached first, the current partial batch flushes,
  • when the source completes, any remaining items flush as the final batch.

Window stream items for time-based summaries

Window(duration) emits the items collected during each time window. This is useful for metrics, telemetry, monitoring, rolling summaries, and compacting high-volume feeds.

var metricPipeline = Pipeline
    .StreamFor<MetricPoint>()
    .Window(TimeSpan.FromSeconds(10))
    .Then(window => new MetricSummary(
        Count: window.Count,
        Average: window.Count == 0 ? 0 : window.Average(point => point.Value),
        Max: window.Count == 0 ? 0 : window.Max(point => point.Value)))
    .Build("Metric summaries");

await foreach (var summary in metricPipeline.RunStreamAsync(metricPoints, cancellationToken))
{
    await dashboard.PublishAsync(summary, cancellationToken);
}

Windowing is time-first. It is a good fit when the question is �what happened during this period?� rather than �do I have exactly N items?�

Large file import example

Streaming lets you process large files without materializing all rows first.

static async IAsyncEnumerable<CustomerRow> ReadCustomerRowsAsync(
    string path,
    [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
    await foreach (var row in Csv.ReadAsync<CustomerRow>(path, cancellationToken))
        yield return row;
}

var importPipeline = Pipeline
    .StreamFor<CustomerRow>()
    .Then("Validate", ValidateCustomerRow)
    .Then("Map", MapCustomer)
    .Batch(size: 500, maxDelay: TimeSpan.FromSeconds(2))
    .ThenAsync("Upsert", async (customers, cancellationToken) =>
    {
        await customerRepository.UpsertManyAsync(customers, cancellationToken);
        return customers.Count;
    })
    .Build("Customer CSV import");

var imported = 0;
await foreach (var count in importPipeline.RunStreamAsync(ReadCustomerRowsAsync(path), cancellationToken))
    imported += count;

This keeps memory bounded and makes the batch size explicit.

Cancellation

Stream execution observes cancellation while reading, transforming, batching, and yielding items.

using var cancellationTokenSource = new CancellationTokenSource(TimeSpan.FromMinutes(5));

await foreach (var item in streamPipeline.RunStreamAsync(source, cancellationTokenSource.Token))
{
    await ProcessAsync(item, cancellationTokenSource.Token);
}

Cancel the token to stop consuming the source and stop pending asynchronous steps.

MergeReducers helper

Use MergeReducers directly when you need custom aggregation over detailed branch outcomes (ForkResult<T>). Reducers operate on the branch result list produced by Fork(...), so they can preserve branch order and skip failed branches consistently:

var forkPipeline = Pipeline
    .For<int>()
    .Fork<int>(
        (value, _) => ValueTask.FromResult(value + 1),
        (_, _) => ValueTask.FromException<int>(new InvalidOperationException("branch failed")),
        (value, _) => ValueTask.FromResult(value + 3))
    .Build();

var forkExecution = await forkPipeline.RunAsync(10);

var reduced = await MergeReducers.ReduceAsync(
    forkExecution.BranchResults,
    0,
    (acc, value, _) => ValueTask.FromResult(acc + value));

Console.WriteLine(reduced);
// 24

Dependency injection friendly operations

var pipeline = Pipeline
    .For<int>()
    .Then<MyServiceStep, int>(() => new MyServiceStep(serviceProvider.GetRequiredService<MyService>()))
    .Build();

public sealed class MyServiceStep(MyService service) : IPipelineStep<int, int>
{
    public ValueTask<int> ExecuteAsync(int input, CancellationToken cancellationToken = default) =>
        ValueTask.FromResult(service.Transform(input));
}

Pipeline descriptions and visualization

Built pipelines expose structural metadata through Describe(). Named overloads make exported graphs readable for documentation, diagnostics, pull request notes, architecture diagrams, or UI rendering.

PipelineDefinition contains:

  • pipeline ID, name, and version,
  • graph nodes,
  • graph edges,
  • node kinds,
  • step input/output types.

Request-response pipelines and stream pipelines both expose Describe().

var pipeline = Pipeline
    .For<Order>()
    .Then("Validate", ValidateOrder)
    .ThenAsync("Price", PriceOrderAsync)
    .Branch(
        "Route by risk",
        order => order.RiskScore > 80,
        high => high with { ReviewRequired = true },
        low => low)
    .Build("Order workflow", "2.4.0");

var definition = pipeline.Describe();

Export as JSON

Use ToJson() when you want to store, inspect, diff, or render the pipeline definition in another tool.

string json = definition.ToJson();

Console.WriteLine(json);

Example output:

{
  "id": "4f4e7b8f-2c2e-4f2f-9f5d-4b90f9f2c99f",
  "name": "Order workflow",
  "version": "2.4.0",
  "nodes": [
    {
      "id": "input",
      "name": "Input",
      "kind": "Input",
      "inputType": "Order",
      "outputType": "Order"
    },
    {
      "id": "step_1",
      "name": "Validate",
      "kind": "Step",
      "inputType": "Order",
      "outputType": "ValidatedOrder"
    }
  ],
  "edges": [
    {
      "from": "input",
      "to": "step_1",
      "label": null
    }
  ]
}

You can pass custom JsonSerializerOptions if you want different formatting:

string compactJson = definition.ToJson(
    new JsonSerializerOptions { WriteIndented = false });

Export as Mermaid

Use ToMermaid() when you want Markdown-friendly diagrams for GitHub, documentation sites, or generated architecture notes.

string mermaid = definition.ToMermaid();

Console.WriteLine(mermaid);

Example output:

flowchart TD
    input["Input<br/>Input<br/>Order -> Order"]
    step_1["Validate<br/>Step<br/>Order -> ValidatedOrder"]
    step_2["Price<br/>Step<br/>ValidatedOrder -> PricedOrder"]
    input --> step_1
    step_1 --> step_2

You can paste the generated Mermaid into Markdown renderers that support Mermaid diagrams.

Export as Graphviz DOT

Use ToDot() when you want to render the pipeline graph with Graphviz or other tooling that understands DOT files.

string dot = definition.ToDot();

await File.WriteAllTextAsync("order-workflow.dot", dot, cancellationToken);

Example output:

digraph "Order workflow" {
    "input" [label="Input\nInput\nOrder -> Order"];
    "step_1" [label="Validate\nStep\nOrder -> ValidatedOrder"];
    "step_2" [label="Price\nStep\nValidatedOrder -> PricedOrder"];
    "input" -> "step_1";
    "step_1" -> "step_2";
}

A DOT file can be rendered with Graphviz:

dot -Tpng order-workflow.dot -o order-workflow.png

Observers and step attempts

Use IPipelineObserver when you need live run and step events without adding logging or metrics code inside every step. Observers can collect profiling data, audit records, operational logs, or custom telemetry.

public sealed class LoggingPipelineObserver : IPipelineObserver
{
    public ValueTask OnStepCompletedAsync(
        PipelineStepCompleted completed,
        CancellationToken cancellationToken = default)
    {
        var attempt = completed.Attempt;
        Console.WriteLine(
            $"{attempt.StepName} completed in {attempt.Duration.TotalMilliseconds}ms");

        return ValueTask.CompletedTask;
    }

    public ValueTask OnStepFailedAsync(
        PipelineStepFailed failed,
        CancellationToken cancellationToken = default)
    {
        var attempt = failed.Attempt;
        Console.WriteLine(
            $"{attempt.StepName} failed on attempt {attempt.AttemptNumber}: {attempt.ExceptionMessage}");

        return ValueTask.CompletedTask;
    }
}

var pipeline = Pipeline
    .For<Order>()
    .WithObserver(new LoggingPipelineObserver())
    .Then("Validate", ValidateOrder)
    .ThenAsync("Submit", SubmitOrderAsync)
    .Build("Observed order workflow", "2.4.0");

PipelineStepAttempt records include run ID, pipeline ID, pipeline name, pipeline version, step ID, step name, step kind, attempt number, status, start/end timestamps, duration, exception type, and exception message.

Step tracing

Use RunWithTrace(...) or RunWithTraceAsync(...) to execute a pipeline and capture per-step timing metadata for fluent steps.

var run = await pipeline.RunWithTraceAsync(input, cancellationToken);

Console.WriteLine(run.Result);

foreach (var step in run.Trace.Steps)
{
    Console.WriteLine($"{step.Name}: {step.Duration.TotalMilliseconds}ms");
}

Trace entries include the step name, kind, input/output types, duration, success flag, and exception type when captured around a failing step.

Dry-run validation

Use DryRun() to validate the captured pipeline structure without executing any step delegates or side effects.

var report = pipeline.DryRun();

if (!report.IsValid)
{
    foreach (var issue in report.Issues)
    {
        Console.WriteLine($"{issue.Severity}: {issue.Code} - {issue.Message}");
    }
}

Dry-run validation checks graph consistency, missing edge endpoints, duplicate node IDs, and unreachable nodes.

Observability

OperationPipeline emits:

  • ActivitySource spans using source name Pipeliner.Net,
  • Meter metrics:
    • pipeliner.pipeline.runs
    • pipeliner.pipeline.failures
    • pipeliner.pipeline.duration.ms
    • pipeliner.pipeline.operation.duration.ms

This integrates cleanly with OpenTelemetry collectors and exporters.

Error handling model

  • Per-operation exception handling via onExceptionHandler.
  • Unhandled exceptions bubble to caller.
  • Retry and similar behavior can be applied via IPipelineExecutionPolicy.

Testing guidance

Pipelines are easy to unit test because:

  • each step is just a delegate or IPipelineStep,
  • full flow can be executed in-memory,
  • branch/fork/merge behavior is deterministic and explicit.

Roadmap alignment

Current API includes:

  • async-first delegates with ValueTask,
  • typed fluent composition,
  • batch APIs,
  • streaming APIs with channel-backed backpressure,
  • policy hooks,
  • branch/fork/merge capabilities,
  • built-in merge conflict strategies and reducers,
  • instrumentation support.

Future phases can expand higher-level integration helpers and additional execution policies without breaking request-response usage.

Product Compatible and additional computed target framework versions.
.NET 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 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. 
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
2.5.2 135 7/7/2026
2.5.1 130 7/7/2026
2.5.0 126 7/7/2026
2.3.0 125 7/6/2026
2.2.0 119 7/6/2026
2.1.0 115 7/3/2026
2.0.2 117 7/2/2026
2.0.1 122 7/2/2026
1.0.3 224 5/10/2024
1.0.2 246 4/28/2024
1.0.1 221 4/26/2024
1.0.0 223 4/26/2024