WalFlow.Abstractions 0.1.0-alpha.alpha.20260708090623

This is a prerelease version of WalFlow.Abstractions.
dotnet add package WalFlow.Abstractions --version 0.1.0-alpha.alpha.20260708090623
                    
NuGet\Install-Package WalFlow.Abstractions -Version 0.1.0-alpha.alpha.20260708090623
                    
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="WalFlow.Abstractions" Version="0.1.0-alpha.alpha.20260708090623" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="WalFlow.Abstractions" Version="0.1.0-alpha.alpha.20260708090623" />
                    
Directory.Packages.props
<PackageReference Include="WalFlow.Abstractions" />
                    
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 WalFlow.Abstractions --version 0.1.0-alpha.alpha.20260708090623
                    
#r "nuget: WalFlow.Abstractions, 0.1.0-alpha.alpha.20260708090623"
                    
#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 WalFlow.Abstractions@0.1.0-alpha.alpha.20260708090623
                    
#: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=WalFlow.Abstractions&version=0.1.0-alpha.alpha.20260708090623&prerelease
                    
Install as a Cake Addin
#tool nuget:?package=WalFlow.Abstractions&version=0.1.0-alpha.alpha.20260708090623&prerelease
                    
Install as a Cake Tool

WalFlow.Abstractions

Core abstractions and interfaces for the WalFlow Change Data Capture library.

Overview

This package provides the foundational types, interfaces, and models used across all WalFlow components. It has zero external dependencies, making it lightweight and suitable for building custom CDC components.

Interfaces

ISink

Defines the contract for CDC event destinations.

public interface ISink
{
    Task SendAsync(DebeziumPayload payload, CancellationToken ct = default);
    Task FlushAsync(CancellationToken ct = default);
}

Implementations:

  • WalFlow.Sinks.File - Write events to file
  • WalFlow.Sinks.Kafka - Publish to Kafka topics
  • Build your own!

IReplicationStore

Defines the contract for persisting replication state/checkpoints.

public interface IReplicationStore
{
    Task<ReplicationState?> LoadReplicationStateAsync(string slotName, CancellationToken ct = default);
    Task SaveReplicationStateAsync(ReplicationState state, CancellationToken ct = default);
}

Implementations:

  • WalFlow.Stores.File - JSON file storage
  • WalFlow.Stores.Redis - Redis persistence
  • Build your own!

ISchemaHistoryStore

Defines the contract for persisting and replaying schema-change history independent of replication state storage.

public interface ISchemaHistoryStore
{
    Task<SchemaHistoryStoreReadiness> CheckReadinessAsync(CancellationToken ct = default);
    IAsyncEnumerable<SchemaChangeMessage> ReadAllAsync(CancellationToken ct = default);
    Task AppendAsync(SchemaChangeMessage message, CancellationToken ct = default);
}

ISnapshotStrategy

Defines the contract for snapshot execution strategies.

public interface ISnapshotStrategy
{
    Task<ulong?> ExecuteAsync(CancellationToken ct);
}

Built-in strategies:

  • TransactionalExported - Consistent snapshot with exported transaction
  • Incremental - Debezium-style chunked snapshot with live WAL merge
  • BestEffort - Simple read without blocking

Models

DebeziumPayload

Debezium-compatible CDC event structure with key, before, after, source, operation, and transaction metadata.

ReplicationState

Tracks CDC processor state including LSN position, snapshot progress, and incremental cursors.

CdcProcessorOptions

Configuration for CDC processors including connection strings, snapshot modes, and publication settings.

CDC Payload Middleware (SMT)

Reusable middleware contracts for source-agnostic payload transforms:

  • ICdcPayloadMiddleware
  • CdcPayloadMiddlewareBuilder
  • CdcPayloadMiddlewarePipeline
  • CdcPayloadTransformingSink
  • CdcPayloadMiddlewareContext.Stage (Global or Sink)

Built-in transform middleware:

  • FilterCdcPayloadMiddleware (operation/schema/table filtering)
  • ReplaceFieldCdcPayloadMiddleware (include/exclude/rename)
  • MaskFieldCdcPayloadMiddleware (field masking)

Sink Composition

Reusable sink composition utilities:

  • FanOutSink - sends each payload to all configured sinks
  • SinkFilteringSink - applies per-sink filtering
  • SinkRegistration - describes one sink endpoint with optional per-sink filters and middleware

Sink Payload Transform Adapter

For custom sink protocols that require strings/bytes instead of DebeziumPayload objects:

  • ISinkPayloadTransformer - transforms DebeziumPayload to SinkTransformedMessage
  • ITransformedSinkTarget - target for transformed messages
  • SinkPayloadTransformAdapter - bridges WalFlow ISink to transformed targets
  • DelegateSinkPayloadTransformer / DelegateTransformedSinkTarget - delegate-based helpers

Use this when your destination expects a non-Debezium wire format and you want full control over final payload shape. This adapter is transport-focused (shape/encoding), not filtering. Use connector/sink filters for subset routing. Most built-in sinks also accept an optional ISinkPayloadTransformer constructor/factory argument to override transport encoding directly.

Usage

dotnet add package WalFlow.Abstractions
using WalFlow.Abstractions;
using WalFlow.Abstractions.Models;

public class MyCustomSink : ISink
{
    public async Task SendAsync(DebeziumPayload payload, CancellationToken ct)
    {
        // Your custom logic
    }

    public async Task FlushAsync(CancellationToken ct)
    {
        // Flush buffered data
    }
}
using WalFlow.Abstractions.Transforms;

var builder = new CdcPayloadMiddlewareBuilder()
    .Use((payload, context, next, ct) =>
    {
        payload.After ??= new Dictionary<string, object?>();
        payload.After["processed_by"] = context.ConnectorName;
        return next(payload, context, ct);
    });

var pipeline = builder.Build();
var transformingSink = new CdcPayloadTransformingSink(innerSink, pipeline);

License

MIT

Product Compatible and additional computed target framework versions.
.NET 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.
  • net10.0

    • No dependencies.

NuGet packages (36)

Showing the top 5 NuGet packages that depend on WalFlow.Abstractions:

Package Downloads
WalFlow.Sources.Postgres

PostgreSQL Change Data Capture source for WalFlow. Captures database changes using logical replication (WAL) with support for multiple snapshot strategies.

WalFlow.Transforms.DependencyInjection

Dependency injection extensions and ordered pipeline builder for WalFlow CDC payload middleware.

WalFlow.Stores.Memory

In-memory store for WalFlow CDC processor. Stores replication state in memory (non-persistent).

WalFlow.Failover.Postgres

PostgreSQL lease store and worker adapter plugin for WalFlow failover.

WalFlow.Failover

Lease-based active-passive failover coordinator for WalFlow workers.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
0.1.0-alpha.alpha... 262 7/8/2026
0.1.0-alpha.alpha... 143 3/30/2026