RtFlow.Pipelines.Hosting 1.0.0

The owner has unlisted this package. This could mean that the package is deprecated, has security vulnerabilities or shouldn't be used anymore.
dotnet add package RtFlow.Pipelines.Hosting --version 1.0.0
                    
NuGet\Install-Package RtFlow.Pipelines.Hosting -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="RtFlow.Pipelines.Hosting" Version="1.0.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="RtFlow.Pipelines.Hosting" Version="1.0.0" />
                    
Directory.Packages.props
<PackageReference Include="RtFlow.Pipelines.Hosting" />
                    
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 RtFlow.Pipelines.Hosting --version 1.0.0
                    
#r "nuget: RtFlow.Pipelines.Hosting, 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 RtFlow.Pipelines.Hosting@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=RtFlow.Pipelines.Hosting&version=1.0.0
                    
Install as a Cake Addin
#tool nuget:?package=RtFlow.Pipelines.Hosting&version=1.0.0
                    
Install as a Cake Tool

RtFlow.Pipelines.Hosting

NuGet

RtFlow.Pipelines.Hosting provides seamless integration with ASP.NET Core and .NET Generic Host, enabling lifecycle management of data processing pipelines as hosted services.

Installation

# Install core package first
dotnet add package RtFlow.Pipelines.Core

# Then add hosting integration
dotnet add package RtFlow.Pipelines.Hosting

Key Features

  • 🏠 ASP.NET Core Integration - Register pipelines as hosted services
  • 🔄 Lifecycle Management - Automatic startup, shutdown, and graceful termination
  • ⚙️ Configuration Support - Configure pipelines via appsettings.json
  • 📊 Health Checks - Built-in health check integration
  • 🔍 Observability - Logging and metrics integration with .NET diagnostics
  • 🎛️ Background Processing - Long-running pipeline services

Quick Start

Basic Hosted Pipeline

using RtFlow.Pipelines.Hosting;

var builder = WebApplication.CreateBuilder(args);

// Register a pipeline as a hosted service
builder.Services.AddPipelineHostedService<DataProcessor>(serviceProvider =>
{
    return PipelineBuilder
        .Create<DataItem>()
        .Transform(item => ProcessItem(item))
        .ForEach(result => SaveResult(result))
        .Build();
});

var app = builder.Build();
app.Run();

Pipeline Factory Service

public class DataProcessingService : PipelineHostedService<DataItem>
{
    private readonly ILogger<DataProcessingService> _logger;
    private readonly IDataSource _dataSource;

    public DataProcessingService(
        ILogger<DataProcessingService> logger,
        IDataSource dataSource)
    {
        _logger = logger;
        _dataSource = dataSource;
    }

    protected override IPipeline<DataItem> CreatePipeline()
    {
        return PipelineBuilder
            .Create<DataItem>()
            .Transform(async item => await EnrichData(item))
            .Filter(item => item.IsValid)
            .Batch(100, TimeSpan.FromSeconds(30))
            .Transform(batch => ProcessBatch(batch))
            .SideEffect(result => _logger.LogInformation("Processed batch: {Count}", result.Count))
            .Build();
    }

    protected override async IAsyncEnumerable<DataItem> GetDataSource(
        [EnumeratorCancellation] CancellationToken cancellationToken)
    {
        await foreach (var item in _dataSource.StreamAsync(cancellationToken))
        {
            yield return item;
        }
    }
}

// Register in Program.cs
builder.Services.AddHostedService<DataProcessingService>();

Configuration

appsettings.json Configuration

{
  "PipelineOptions": {
    "MaxConcurrency": 10,
    "BoundedCapacity": 1000,
    "EnableMetrics": true,
    "GracefulShutdownTimeout": "00:00:30"
  }
}

Dependency Injection

// Configure pipeline options
builder.Services.Configure<PipelineOptions>(
    builder.Configuration.GetSection("PipelineOptions"));

// Register dependencies
builder.Services.AddScoped<IDataProcessor, DataProcessor>();
builder.Services.AddScoped<IDataRepository, DataRepository>();

// Register pipeline with dependencies
builder.Services.AddPipelineHostedService<OrderData>(serviceProvider =>
{
    var processor = serviceProvider.GetRequiredService<IDataProcessor>();
    var repository = serviceProvider.GetRequiredService<IDataRepository>();
    
    return PipelineBuilder
        .Create<OrderData>()
        .Transform(order => processor.ProcessOrder(order))
        .ForEach(result => repository.SaveResult(result))
        .Build();
});

Health Checks

builder.Services.AddHealthChecks()
    .AddPipelineHealthCheck<DataProcessingService>("data-pipeline");

// Health check endpoint
app.MapHealthChecks("/health");

Graceful Shutdown

The hosting integration automatically handles graceful shutdown:

public class DataProcessingService : PipelineHostedService<DataItem>
{
    protected override async Task OnStoppingAsync(CancellationToken cancellationToken)
    {
        _logger.LogInformation("Pipeline is shutting down gracefully...");
        
        // Custom cleanup logic
        await FlushPendingData();
        
        await base.OnStoppingAsync(cancellationToken);
    }
}

Background Processing Patterns

Continuous Processing

public class ContinuousProcessorService : PipelineHostedService<QueueMessage>
{
    protected override async IAsyncEnumerable<QueueMessage> GetDataSource(
        [EnumeratorCancellation] CancellationToken cancellationToken)
    {
        while (!cancellationToken.IsCancellationRequested)
        {
            var messages = await _messageQueue.ReceiveAsync(cancellationToken);
            foreach (var message in messages)
            {
                yield return message;
            }
        }
    }
}

Scheduled Processing

public class ScheduledProcessorService : BackgroundService
{
    private readonly IPipeline<DataItem> _pipeline;

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        using var timer = new PeriodicTimer(TimeSpan.FromHours(1));
        
        while (!stoppingToken.IsCancellationRequested && 
               await timer.WaitForNextTickAsync(stoppingToken))
        {
            var data = await LoadBatchData();
            await _pipeline.ExecuteAsync(data, stoppingToken);
        }
    }
}

Requirements

  • .NET 8.0 or later
  • RtFlow.Pipelines.Core package
  • Microsoft.Extensions.Hosting package

Documentation

For complete documentation and examples, visit the main project repository.

License

This project is licensed under the MIT License - see the LICENSE file for details.

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

Hosting package for ASP.NET Core integration with dependency injection and lifecycle management.