Winche.Events
6.1.0
dotnet add package Winche.Events --version 6.1.0
NuGet\Install-Package Winche.Events -Version 6.1.0
<PackageReference Include="Winche.Events" Version="6.1.0" />
<PackageVersion Include="Winche.Events" Version="6.1.0" />
<PackageReference Include="Winche.Events" />
paket add Winche.Events --version 6.1.0
#r "nuget: Winche.Events, 6.1.0"
#:package Winche.Events@6.1.0
#addin nuget:?package=Winche.Events&version=6.1.0
#tool nuget:?package=Winche.Events&version=6.1.0
Winche.Events
A Marten-backed event sourcing library for .NET 10. Provides typed projections, an explicit unit-of-work session, optimistic concurrency, post-commit notifications, and an optional command-dispatch layer — all without exposing Marten types to your domain code.
Packages
| Package | Purpose |
|---|---|
Winche.Events.Abstractions |
Base types: IAggregate, Aggregate, IEvent, Event, ICommand<TAggregate>, Command<TAggregate>, EventEnvelope<TEvent>, StreamEnvelope<TAggregate> |
Winche.Events |
Core: event store, sessions, projections, notifiers |
Winche.Events.Commands |
Optional: command handlers and dispatcher |
Winche.Events.Grpc |
Optional: gRPC transport — exposes the event store over Dispatch (unary) + WatchStream (server streaming) |
Winche.Events.WebSocket |
Optional: WebSocket transport — single multiplexed JSON connection; browser-compatible on a single HTTP/1.1 port |
Getting started
1. Define your domain model
using Winche.Events.Abstractions;
record Order(string Status, decimal Total) : Aggregate
{
public static Order Empty => new("none", 0);
}
record OrderPlaced(string OrderId, decimal Total) : Event;
record OrderShipped(string OrderId) : Event;
record OrderCancelled(string OrderId) : Event;
2. Define a projection
Projections fold events into an aggregate document. Define one public Apply method per event type. The aggregate type is inferred from the base class — no second type parameter needed. Unregistered event types are silently ignored.
using Winche.Events.Projection;
class OrderProjection : InlineProjection<Order>
{
public Order Apply(Order state, EventEnvelope<OrderPlaced> e) =>
state with { Status = "placed", Total = e.Data.Total };
public Order Apply(Order state, EventEnvelope<OrderShipped> e) =>
state with { Status = "shipped" };
public Order Apply(Order state, EventEnvelope<OrderCancelled> e) =>
state with { Status = "cancelled" };
public override Order Create(string id) => Order.Empty with { Id = id };
}
Each Apply method receives EventEnvelope<TEvent> — access the event via e.Data and stream metadata via e.Version, e.Timestamp, e.StreamId, e.Id, e.Sequence.
3. Register services
using Winche.Events.DependencyInjection;
services.AddWincheEvents(opts =>
{
opts.ConnectionString = "Host=localhost;Database=mydb;Username=postgres;Password=...";
opts.AddEvent<OrderPlaced>();
opts.AddEvent<OrderShipped>();
opts.AddEvent<OrderCancelled>();
opts.AddProjection<OrderProjection>(); // TAggregate inferred from base class
});
4. Use the event store
await using var session = await store.OpenSessionAsync();
await session.AppendStreamAsync("orders/123", [new OrderPlaced("orders/123", 49.99m)]);
await session.SaveChangesAsync();
var order = await session.GetStateAsync<Order>("orders/123");
// order.Status == "placed"
Projection modes
The lifecycle is determined by the base class — no mode parameter needed.
| Base class | When document is built | Method name | GetStateAsync returns |
|---|---|---|---|
InlineProjection<T> |
Same transaction as the append | Apply |
Always fresh after SaveChangesAsync |
AsyncProjection<T> |
Background daemon after commit | ApplyAsync |
Eventually consistent |
When at least one AsyncProjection<T> is registered, AddWincheEvents automatically enables the Marten async daemon in HotCold mode (leader election via PostgreSQL advisory locks — safe for multi-instance deployments).
InlineProjection<T> — for projections that need to be immediately consistent. Handlers run inside the open PostgreSQL transaction — keep them fast and free of external I/O.
class OrderProjection : InlineProjection<Order>
{
public Order Apply(Order state, EventEnvelope<OrderPlaced> e) =>
state with { Status = "placed", Total = e.Data.Total };
public Order Apply(Order state, EventEnvelope<OrderShipped> e) =>
state with { Status = "shipped" };
public override Order Create(string id) => Order.Empty with { Id = id };
}
AsyncProjection<T> — for projections that need to call external APIs or secondary databases. The daemon runs in the background; GetStateAsync may return stale state until it catches up.
class OrderReportProjection : AsyncProjection<OrderReport>
{
public async Task<OrderReport> ApplyAsync(OrderReport state, EventEnvelope<OrderShipped> e)
{
var details = await _warehouse.GetAsync(e.Data.OrderId);
return state with { LastShipped = details.ProductName };
}
public override OrderReport Create(string id) => new OrderReport { Id = id };
}
opts.AddProjection<OrderProjection>(); // inline
opts.AddProjection<OrderReportProjection>(); // async daemon — inferred from base class
IEventSession
IEventSession is a unit of work scoped to a single PostgreSQL connection. Always dispose with await using.
AppendStreamAsync
await session.AppendStreamAsync("orders/123", [new OrderPlaced("orders/123", 49.99m)]);
Optimistic concurrency — pass expectedVersion to reject concurrent writes:
await session.AppendStreamAsync("orders/123", events, expectedVersion: 3);
// Marten throws if the stream's current version doesn't match.
GetStateAsync
Reads the stored aggregate document. For Inline projections this is always fresh after commit. For Async projections it reflects the last time the daemon ran.
var order = await session.GetStateAsync<Order>("orders/123");
QueryStatesAsync
Queries stored aggregate documents using a LINQ compose function.
var placed = await session.QueryStatesAsync<Order>(
q => q.Where(o => o.Status == "placed"));
var recent = await session.QueryStatesAsync<Order>(
q => q.OrderByDescending(o => o.Id).Take(10));
Pass q => q to return all documents of that type.
GetStreamAsync
Returns a StreamEnvelope<TAggregate> combining the mt_streams row with the current projected aggregate. Returns null if the stream does not exist.
var envelope = await session.GetStreamAsync<Order>("orders/123");
if (envelope is not null)
{
// envelope.Id — stream identifier
// envelope.Aggregate — current projected document (null if no projection stored yet)
// envelope.Version — total events appended
// envelope.Created — when the stream was first written
// envelope.LastModified — when the stream last received an event
// envelope.IsArchived — whether the stream is archived
}
GetEventsAsync
Returns events for a stream in order, each wrapped as EventEnvelope<IEvent>. Pass fromVersion to fetch only events at or after that version.
var events = await session.GetEventsAsync("orders/123");
var recent = await session.GetEventsAsync("orders/123", fromVersion: 5);
Use OfEventType<TEvent>() to filter and cast to a specific event type:
var placed = events.OfEventType<OrderPlaced>();
// Each element is EventEnvelope<OrderPlaced> with strongly-typed e.Data
SaveChangesAsync
Commits all buffered appends to PostgreSQL, then fires registered notifiers.
await session.SaveChangesAsync();
EventEnvelope<TEvent>
Projection handlers and GetEventsAsync work with EventEnvelope<TEvent>. It carries the full row from mt_events.
public sealed record EventEnvelope<TEvent>(
string Id, // event UUID as string
string StreamId, // stream identifier
TEvent Data, // strongly-typed event payload
long Version, // 1-based position within the stream
DateTimeOffset Timestamp, // when the event was committed (UTC)
long Sequence, // global sequence number across all streams
string TypeAlias, // value in the mt_events.type column
string DotNetType // value in the mt_events.dotnet_type column
) where TEvent : IEvent;
Post-commit notifications
Implement IAppendNotifier to receive a callback after each successful commit:
using Winche.Events.Notification;
class OrderNotifier : IAppendNotifier
{
public Task NotifyAsync(string streamId, IReadOnlyList<IEvent> events,
CancellationToken ct = default)
{
// Runs after the PostgreSQL transaction commits.
// Events are already persisted — this cannot roll them back.
return Task.CompletedTask;
}
}
opts.AddNotifier<OrderNotifier>();
Multiple notifiers can be registered. An exception in one is logged and swallowed; it does not affect the others or the caller.
Event observers (EventObserver)
Event observers react to events through the Marten async daemon — ordered, checkpointed, and at-least-once — without writing any aggregate document. Use them for relational DB writes, audit trails, external API calls, or any side effect that needs retry-on-failure.
Define Handle (sync) or HandleAsync (async) methods per event type. Both can coexist for the same type — sync is called first.
using Winche.Events.Projection;
class OrderAuditObserver : EventObserver
{
private readonly IAuditDb _db;
public OrderAuditObserver(IAuditDb db) => _db = db;
public void Handle(EventEnvelope<OrderPlaced> e) =>
_db.Record(e.StreamId, "placed", e.Timestamp);
public async Task HandleAsync(EventEnvelope<OrderShipped> e) =>
await _db.RecordAsync(e.StreamId, "shipped", e.Timestamp);
}
services.AddWincheEvents(opts =>
{
opts.AddEventObserver<OrderAuditObserver>();
});
AddEventObserver automatically enables the Marten async daemon if it is not already enabled. Observers are registered as singletons and support constructor injection.
vs IAppendNotifier:
IAppendNotifier |
EventObserver |
|
|---|---|---|
| Retry on crash | No | Yes (daemon checkpoint) |
| Ordered across restarts | No | Yes |
| Can rebuild from event 0 | No | Yes |
| Timing | Inline post-commit | Eventually consistent |
Commands (Winche.Events.Commands)
The commands package adds a load → handle → append → commit dispatch loop on top of IEventSession.
1. Define commands
Commands extend Command<TAggregate> (or implement ICommand<TAggregate>). The aggregate type is encoded in the type.
using Winche.Events.Abstractions;
record PlaceOrderCommand(string OrderId, decimal Total) : Command<Order>;
record ShipOrderCommand(string OrderId) : Command<Order>;
Command<TAggregate> carries one built-in property:
long? ExpectedVersion { get; init; } // null = no version check
2. Define a command handler
All commands for an aggregate live in one CommandHandler<TAggregate> class. Define one public Handle or HandleAsync method per command type — the aggregate type is inferred from the base class.
using Winche.Events.Commands;
class OrderCommandHandler : CommandHandler<Order>
{
public IEnumerable<IEvent> Handle(Order? state, PlaceOrderCommand cmd)
{
if (state is { Status: not "none" })
throw new InvalidOperationException("Order already exists.");
return [new OrderPlaced(cmd.OrderId, cmd.Total)];
}
public IEnumerable<IEvent> Handle(Order? state, ShipOrderCommand cmd)
{
if (state is null or { Status: "none" })
throw new InvalidOperationException("Order does not exist.");
return [new OrderShipped(cmd.OrderId)];
}
}
state is the current aggregate loaded from the store (null if the stream does not exist). Throw to reject the command — no events will be appended.
Async handlers — use HandleAsync when the handler needs to do async work (external API calls, secondary DB lookups) before producing events. Sync and async handlers can coexist in the same class:
class OrderCommandHandler : CommandHandler<Order>
{
private readonly IInventoryService _inventory;
public OrderCommandHandler(IInventoryService inventory)
=> _inventory = inventory;
public IEnumerable<IEvent> Handle(Order? state, PlaceOrderCommand cmd)
{
if (state is { Status: not "none" })
throw new InvalidOperationException("Order already exists.");
return [new OrderPlaced(cmd.OrderId, cmd.Total)];
}
public async Task<IEnumerable<IEvent>> HandleAsync(Order? state, ShipOrderCommand cmd, CancellationToken ct)
{
var stock = await _inventory.ReserveStockAsync(cmd.OrderId, ct);
return [new OrderShipped(cmd.OrderId, reservedStock: stock)];
}
}
3. Register
using Winche.Events.Commands.DependencyInjection;
services.AddWincheEventsCommands(cmds =>
{
cmds.AddCommandHandler<OrderCommandHandler>(); // TAggregate inferred from base class
});
4. Dispatch
var dispatcher = provider.GetRequiredService<ICommandDispatcher>();
var result = await dispatcher.DispatchAsync("orders/123", new PlaceOrderCommand("orders/123", 49.99m));
DispatchAsync returns a DispatchResult with the appended events and new stream version:
result.Version // new stream version after commit
result.Events // full EventEnvelope<IEvent> for each appended event
result.Events[0].Id // server-assigned event UUID
result.Events[0].Data // the domain event (e.g. OrderPlaced)
Dispatch flow:
- Open a session
- Load current stream via
GetStreamAsync— captures aggregate state and version - Call the registered handler → produce events
- Append events and commit
- Fetch the newly appended events with full server metadata
- Return
DispatchResult
Transaction isolation
OpenSessionAsync accepts an optional IsolationLevel:
await using var session = await store.OpenSessionAsync(IsolationLevel.Serializable);
Default is ReadCommitted.
Connection string configuration
Supply the connection string at registration time:
services.AddWincheEvents(opts =>
{
opts.ConnectionString = builder.Configuration.GetConnectionString("Postgres")
?? throw new InvalidOperationException("Postgres connection string not configured.");
});
For integration tests read from an environment variable so the value is never committed:
private static readonly string ConnectionString =
Environment.GetEnvironmentVariable("WINCHE_TEST_CONN")
?? throw new InvalidOperationException("Set WINCHE_TEST_CONN before running integration tests.");
$env:WINCHE_TEST_CONN = "Host=localhost;Database=winche_events_test;Username=postgres;Password=..."
dotnet test
Web API sample (samples/Winche.Events.WebApi)
A minimal ASP.NET Core API demonstrating ICommandDispatcher and IEventSession over HTTP. Uses a Notes domain (create / update / delete).
Endpoints
| Method | Path | Description |
|---|---|---|
POST |
/api/dispatch |
Dispatch any command. streamId lives in the body so slash-prefixed IDs (e.g. notes/uuid) work without routing issues. |
GET |
/api/notes/{noteId} |
Current note state. Returns 404 if deleted. |
GET |
/api/notes |
All non-deleted notes. |
Running
$env:ConnectionStrings__Postgres = "Host=localhost;Database=mydb;Username=postgres;Password=..."
dotnet run --project samples/Winche.Events.WebApi
# API available at http://localhost:5000
gRPC transport (Winche.Events.Grpc)
Adds a gRPC endpoint on top of the existing event store and command dispatcher. Dart (and any other gRPC client) can dispatch commands, fetch stream history, and receive real-time event pushes per stream — without polling.
Four RPCs
| RPC | Type | Description |
|---|---|---|
Dispatch |
Unary | Send a command → returns DispatchResponse with version + events |
GetStream |
Server streaming, terminates | Fetch historical events for a stream — connection closes when history is exhausted |
WatchStream |
Server streaming, persistent | Subscribe to a specific stream; optional catch-up from a version before going live |
WatchEvents |
Server streaming, persistent | Subscribe to one or more event types across any stream (or a specific stream); live-only |
WatchStream catch-up
WatchStream accepts stream_id and from_version. Setting from_version = 0 gives live events only. Setting from_version = N replays events from version N from the database first, then continues live. The subscribe-first design guarantees no gap — events that arrive while history is being replayed are buffered and deduplicated.
WatchEvents filtering
WatchEvents accepts a list of event type names and an optional stream_id. With no stream_id it delivers matching events from any stream; with a stream_id it narrows to that stream only. Live-only — no catch-up.
Register
Call AddWincheEventsCommands first (registers ICommandDispatcher and your handlers), then AddWincheEventsGrpc (registers the transport). Command types are discovered automatically from method signatures — no manual mapping needed.
builder.Services.AddWincheEvents(opts =>
{
opts.AddEvent<NoteCreated>("NoteCreated");
opts.AddProjection<NoteProjection>();
});
builder.Services.AddWincheEventsCommands(opts =>
{
opts.AddCommandHandler<NoteCommandHandler>(); // command types discovered automatically
});
builder.Services.AddWincheEventsGrpc(); // pure transport — no handler config needed
app.UseGrpcWeb(new GrpcWebOptions { DefaultEnabled = true });
app.MapWincheEventsGrpc(); // registers all four RPCs
Run the gRPC sample
$env:ConnectionStrings__Postgres = "Host=localhost;..."
dotnet run --project samples/Winche.Events.GrpcSample
# gRPC service at http://localhost:5000 (plain HTTP/2, no TLS — dev only)
The gRPC sample uses HttpProtocols.Http2 without TLS so Dart clients can connect on plain HTTP/2 locally. For production, configure HTTPS and update the Dart transport accordingly.
WebSocket transport (Winche.Events.WebSocket)
An alternative to the gRPC transport for deployments where plain HTTP/1.1 is preferred or where gRPC-Web is inconvenient. A single persistent WebSocket connection per client carries all operations as JSON messages — no proto files, no code generation, browser-compatible out of the box on a single port.
Ten message types
| Client → Server | Description |
|---|---|
DispatchRequest |
Send a command |
GetStreamRequest |
One-shot history fetch — terminates with GetStreamEnd |
WatchStreamRequest |
Subscribe to a stream; fromVersion = 0 = live only, N = catch-up first |
WatchEventsRequest |
Subscribe to event types across streams; live-only |
UnsubscribeRequest |
Cancel an active subscription |
| Server → Client | Description |
|---|---|
DispatchResponse |
Command committed — version + events |
GetStreamResponse |
One event per message, followed by GetStreamEnd |
WatchStreamResponse |
Pushed event for a stream subscription |
WatchEventsResponse |
Pushed event matching the event-type filter |
ErrorResponse |
Failure for a request — code + message |
Setup
builder.Services.AddWincheEvents(opts =>
{
opts.AddEvent<NoteCreated>("NoteCreated");
opts.AddProjection<NoteProjection>();
});
builder.Services.AddWincheEventsCommands(opts =>
{
opts.AddCommandHandler<NoteCommandHandler>();
});
builder.Services.AddWincheEventsWebSocket(); // registers broadcaster + options
app.UseCors(); // required for browser clients
app.UseWebSockets(); // required before MapWincheEventsWebSocket
app.MapWincheEventsWebSocket("/ws");
AddWincheEventsWebSocket must be called after AddWincheEventsCommands. The endpoint path defaults to /ws.
Authentication
Connection-time (recommended): Protect the HTTP upgrade with standard ASP.NET Core auth. The WebSocket connection is rejected before it opens — zero code in the transport:
app.MapWincheEventsWebSocket().RequireAuthorization();
Per-message expiry check (optional): Pass an IsAuthorized delegate that is evaluated before every incoming message. Cheap checks only — reads a claim, never re-verifies a signature:
builder.Services.AddWincheEventsWebSocket(opts =>
{
// Close the connection mid-session when the token expires
opts.IsAuthorized = ctx =>
{
var exp = ctx.User.FindFirst("exp")?.Value;
return exp is not null &&
DateTimeOffset.FromUnixTimeSeconds(long.Parse(exp)) > DateTimeOffset.UtcNow;
};
});
Browser clients cannot set Authorization headers on WebSocket upgrades. The standard workaround is a token in the query string (ws://host/ws?token=...) coupled with a custom IAuthenticationHandler that reads it:
// In your auth handler:
var token = context.Request.Query["token"].FirstOrDefault()
?? context.Request.Headers["Authorization"].ToString().Replace("Bearer ", "");
Native clients (Flutter mobile/desktop) can pass the Authorization: Bearer <token> header normally.
Run the WebSocket sample
$env:ConnectionStrings__Postgres = "Host=localhost;..."
dotnet run --project samples/Winche.Events.WebSocketSample
# WebSocket endpoint at ws://localhost:5002/ws (plain HTTP/1.1 — dev only)
Port 5002 is used so the WebSocket and gRPC samples can run simultaneously. For production, configure HTTPS — the WebSocket upgrade negotiates via TLS the same as any other request.
Requirements
- .NET 10
- PostgreSQL (via Marten / Npgsql)
| 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
- Marten (>= 9.3.5)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 10.0.8)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.8)
- Winche.Events.Abstractions (>= 6.1.0)
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 |
|---|