Pipeliner.Net
2.5.2
dotnet add package Pipeliner.Net --version 2.5.2
NuGet\Install-Package Pipeliner.Net -Version 2.5.2
<PackageReference Include="Pipeliner.Net" Version="2.5.2" />
<PackageVersion Include="Pipeliner.Net" Version="2.5.2" />
<PackageReference Include="Pipeliner.Net" />
paket add Pipeliner.Net --version 2.5.2
#r "nuget: Pipeliner.Net, 2.5.2"
#:package Pipeliner.Net@2.5.2
#addin nuget:?package=Pipeliner.Net&version=2.5.2
#tool nuget:?package=Pipeliner.Net&version=2.5.2
Pipeliner.Net <img src="Pipeliner.Net/Icon.png" alt="Pipeliner.Net Icon" width="40" />
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(...)andRunAsync(...)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
DropOldestorDropWrite, - noisy telemetry where recent data matters more than completeness: often
DropNewestorDropOldest.
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
sizeis reached first, the batch flushes immediately, - if
maxDelayis 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:
ActivitySourcespans using source namePipeliner.Net,Metermetrics:pipeliner.pipeline.runspipeliner.pipeline.failurespipeliner.pipeline.duration.mspipeliner.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 | Versions 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. |
-
net10.0
- Microsoft.Extensions.Logging.Abstractions (>= 8.0.3)
- System.Threading.RateLimiting (>= 8.0.0)
-
net8.0
- Microsoft.Extensions.Logging.Abstractions (>= 8.0.3)
- System.Threading.RateLimiting (>= 8.0.0)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.