WalFlow.Abstractions
0.1.0-alpha.alpha.20260708090623
dotnet add package WalFlow.Abstractions --version 0.1.0-alpha.alpha.20260708090623
NuGet\Install-Package WalFlow.Abstractions -Version 0.1.0-alpha.alpha.20260708090623
<PackageReference Include="WalFlow.Abstractions" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageVersion Include="WalFlow.Abstractions" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageReference Include="WalFlow.Abstractions" />
paket add WalFlow.Abstractions --version 0.1.0-alpha.alpha.20260708090623
#r "nuget: WalFlow.Abstractions, 0.1.0-alpha.alpha.20260708090623"
#:package WalFlow.Abstractions@0.1.0-alpha.alpha.20260708090623
#addin nuget:?package=WalFlow.Abstractions&version=0.1.0-alpha.alpha.20260708090623&prerelease
#tool nuget:?package=WalFlow.Abstractions&version=0.1.0-alpha.alpha.20260708090623&prerelease
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 fileWalFlow.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 storageWalFlow.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 transactionIncremental- Debezium-style chunked snapshot with live WAL mergeBestEffort- 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:
ICdcPayloadMiddlewareCdcPayloadMiddlewareBuilderCdcPayloadMiddlewarePipelineCdcPayloadTransformingSinkCdcPayloadMiddlewareContext.Stage(GlobalorSink)
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 sinksSinkFilteringSink- applies per-sink filteringSinkRegistration- 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- transformsDebeziumPayloadtoSinkTransformedMessageITransformedSinkTarget- target for transformed messagesSinkPayloadTransformAdapter- bridges WalFlowISinkto transformed targetsDelegateSinkPayloadTransformer/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 | Versions 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. |
-
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 |