EventHorizon.EventStreaming.InMemory 1.6.0-preview.36

This is a prerelease version of EventHorizon.EventStreaming.InMemory.
There is a newer prerelease version of this package available.
See the version list below for details.
dotnet add package EventHorizon.EventStreaming.InMemory --version 1.6.0-preview.36
                    
NuGet\Install-Package EventHorizon.EventStreaming.InMemory -Version 1.6.0-preview.36
                    
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="EventHorizon.EventStreaming.InMemory" Version="1.6.0-preview.36" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="EventHorizon.EventStreaming.InMemory" Version="1.6.0-preview.36" />
                    
Directory.Packages.props
<PackageReference Include="EventHorizon.EventStreaming.InMemory" />
                    
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 EventHorizon.EventStreaming.InMemory --version 1.6.0-preview.36
                    
#r "nuget: EventHorizon.EventStreaming.InMemory, 1.6.0-preview.36"
                    
#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 EventHorizon.EventStreaming.InMemory@1.6.0-preview.36
                    
#: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=EventHorizon.EventStreaming.InMemory&version=1.6.0-preview.36&prerelease
                    
Install as a Cake Addin
#tool nuget:?package=EventHorizon.EventStreaming.InMemory&version=1.6.0-preview.36&prerelease
                    
Install as a Cake Tool

EventHorizon

CI License: MIT

EventHorizon is a .NET framework for Event Sourcing and Event Streaming, providing a clean abstraction layer over multiple storage and streaming backends. Build event-driven applications with pluggable persistence (MongoDB, Elasticsearch, Apache Ignite, In-Memory) and messaging (Apache Pulsar, Kafka, In-Memory).

Table of Contents

Features

  • Event Sourcing — Snapshot and View stores with automatic event application
  • Event Streaming — Publish/subscribe with topic-based routing
  • CQRS — Commands, Events, Requests, and Responses as first-class citizens
  • Pluggable Backends — Swap storage and streaming providers without changing business logic
  • Aggregate Pattern — Built-in aggregate lifecycle management with locking
  • Middleware — Extensible pipeline for aggregate processing
  • Multi-Stream Subscriptions — Subscribe to multiple event streams in a single consumer
  • Migration Support — Built-in tooling for migrating between state schemas

Supported Platforms

.NET Version Support Level End of Support
.NET 10 ✅ LTS November 2028
.NET 9 ✅ STS May 2026
.NET 8 ✅ LTS November 2026

NuGet Packages

All packages are published to NuGet.org with the Cts. prefix.

Package Description
Cts.EventHorizon.Abstractions Core interfaces, models, and attributes
Cts.EventHorizon.EventStore Event store abstractions (CRUD stores, locks)
Cts.EventHorizon.EventStore.InMemory In-memory event store (great for testing)
Cts.EventHorizon.EventStore.MongoDb MongoDB-backed event store
Cts.EventHorizon.EventStore.ElasticSearch Elasticsearch-backed event store
Cts.EventHorizon.EventStore.Ignite Apache Ignite-backed event store
Cts.EventHorizon.EventStreaming Event streaming abstractions
Cts.EventHorizon.EventStreaming.InMemory In-memory streaming (great for testing)
Cts.EventHorizon.EventStreaming.Pulsar Apache Pulsar streaming provider
Cts.EventHorizon.EventSourcing Event sourcing orchestration (aggregates, senders, subscriptions)

Quick Start

1. Install packages

# Core + In-Memory (for getting started / testing)
dotnet add package Cts.EventHorizon.EventSourcing
dotnet add package Cts.EventHorizon.EventStore.InMemory
dotnet add package Cts.EventHorizon.EventStreaming.InMemory

2. Define your state

using EventHorizon.Abstractions.Attributes;
using EventHorizon.Abstractions.Interfaces;
using EventHorizon.Abstractions.Interfaces.Actions;
using EventHorizon.Abstractions.Interfaces.Handlers;

[SnapshotStore("my_app_accounts")]
[Stream("$type")]
public class Account : IState,
    IHandleCommand<CreateAccount>,
    IApplyEvent<AccountCreated>
{
    public string Id { get; set; }
    public string Name { get; set; }
    public int Balance { get; set; }

    public void Handle(CreateAccount command, AggregateContext context)
    {
        context.AddEvent(new AccountCreated(command.Name, command.InitialBalance));
    }

    public void Apply(AccountCreated @event)
    {
        Name = @event.Name;
        Balance = @event.Balance;
    }
}

3. Define actions

using EventHorizon.Abstractions.Interfaces.Actions;

public record CreateAccount(string Name, int InitialBalance) : ICommand<Account>;
public record AccountCreated(string Name, int Balance) : IEvent<Account>;

4. Register services

using EventHorizon.Abstractions.Extensions;
using EventHorizon.EventSourcing.Extensions;
using EventHorizon.EventStore.InMemory.Extensions;
using EventHorizon.EventStreaming.InMemory.Extensions;

services.AddEventHorizon(x =>
{
    x.AddEventSourcing()
        .AddInMemorySnapshotStore()
        .AddInMemoryViewStore()
        .AddInMemoryEventStream()
        .ApplyCommandsToSnapshot<Account>();
});

5. Use the client

using EventHorizon.EventSourcing;

public class AccountService
{
    private readonly EventSourcingClient<Account> _client;

    public AccountService(EventSourcingClient<Account> client)
    {
        _client = client;
    }

    public async Task CreateAccountAsync(string name, int balance)
    {
        await _client.CreateSender()
            .Send(new CreateAccount(name, balance))
            .ExecuteAsync();
    }
}

Core Concepts

State (IState)

The root entity that represents the current state of your domain object. Must implement IState with an Id property.

Actions

EventHorizon uses a CQRS-style action hierarchy:

Action Interface Purpose
Command ICommand<T> Mutates state — handled by IHandleCommand<T> on the state class
Event IEvent<T> Records what happened — applied by IApplyEvent<T> on the state class
Request IRequest<T, TResponse> Query or operation that returns a response
Response IResponse<T> Result of a request

Snapshots and Views

  • Snapshot (Snapshot<T>) — The authoritative persisted state, rebuilt by replaying events
  • View (View<T>) — A read-optimized projection derived from events, can combine data from multiple streams

Aggregates

The AggregateBuilder manages the lifecycle of loading state from a store, applying actions, and persisting results. It handles optimistic concurrency via sequence IDs and distributed locking.

Subscriptions

SubscriptionBuilder<T> creates durable consumers that process messages from one or more streams. Implement IStreamConsumer<T> to handle batches of messages:

public class MyConsumer : IStreamConsumer<Event>
{
    public Task OnBatch(SubscriptionContext<Event> context)
    {
        foreach (var message in context.Messages)
        {
            var payload = message.Data.GetPayload();
            // Process event...
        }
        return Task.CompletedTask;
    }
}

Register with:

x.AddSubscription<MyConsumer, Event>(s => s.AddStream<Account>());

Middleware

Aggregate processing supports middleware for cross-cutting concerns:

x.ApplyEventsToView<MyView>(h => h.UseMiddleware<MyMiddleware>());

Architecture

┌─────────────────────────────────────────────────────────────┐
│                    Your Application                         │
│  ┌──────────────────┐  ┌──────────────────┐                │
│  │ EventSourcingClient│  │  StreamingClient  │               │
│  └────────┬─────────┘  └────────┬─────────┘                │
├───────────┼──────────────────────┼──────────────────────────┤
│           │   EventHorizon Core  │                          │
│  ┌────────▼─────────┐  ┌────────▼─────────┐                │
│  │  AggregateBuilder │  │ SubscriptionBuilder│               │
│  │  SenderBuilder    │  │ PublisherBuilder   │               │
│  │  ICrudStore<T>    │  │ ReaderBuilder      │               │
│  └────────┬─────────┘  └────────┬─────────┘                │
├───────────┼──────────────────────┼──────────────────────────┤
│  ┌────────▼─────────┐  ┌────────▼─────────┐                │
│  │   Event Stores    │  │ Event Streaming   │                │
│  │  ┌─────────────┐  │  │ ┌─────────────┐  │                │
│  │  │  MongoDB     │  │  │ │  Pulsar     │  │                │
│  │  │  Elastic     │  │  │ │  Kafka      │  │                │
│  │  │  Ignite      │  │  │ │  In-Memory  │  │                │
│  │  │  In-Memory   │  │  │ └─────────────┘  │                │
│  │  └─────────────┘  │  └──────────────────┘                │
│  └──────────────────┘                                       │
└─────────────────────────────────────────────────────────────┘

Storage Backends

MongoDB

x.AddMongoDbSnapshotStore(config.GetSection("MongoDb").Bind)
 .AddMongoDbViewStore(config.GetSection("MongoDb").Bind);
{
  "MongoDb": {
    "ConnectionString": "mongodb://localhost:27017",
    "Database": "my_database",
    "IgnoreExtraElements": true
  }
}

IgnoreExtraElements (default true) lets snapshot, view and lock documents load after a property was removed or renamed on the state class (or on any class it contains), instead of throwing FormatException. It applies only to the types EventHorizon stores, not to other MongoDB collections in the application. Set it to false to keep the driver's strict default.

Elasticsearch

x.AddElasticSnapshotStore(config.GetSection("ElasticSearch").Bind)
 .AddElasticViewStore(config.GetSection("ElasticSearch").Bind);
{
  "ElasticSearch": {
    "Uri": "http://localhost:9200"
  }
}

Apache Ignite

x.AddIgniteSnapshotStore(config.GetSection("Ignite").Bind)
 .AddIgniteViewStore(config.GetSection("Ignite").Bind);

In-Memory

x.AddInMemorySnapshotStore()
 .AddInMemoryViewStore();

Best suited for unit/integration testing. No external dependencies required.

Streaming Backends

Apache Pulsar

x.AddPulsarEventStream(config.GetSection("Pulsar").Bind);
{
  "Pulsar": {
    "ServiceUrl": "pulsar://localhost:6650"
  }
}

In-Memory

x.AddInMemoryEventStream();

Best suited for unit/integration testing. No external dependencies required.

Configuration

Attributes

Attribute Target Purpose
[SnapshotStore("bucket_id")] Class Configures the snapshot store bucket/collection name
[ViewStore("database")] Class Configures the view store database/index name
[Stream("topic")] Class Maps a type to a streaming topic
[StreamPartitionKey] Property Designates the property used for stream partitioning
[StoreField(FieldIntent...)] Property Declares how a state field is queried so the store can map/index it efficiently

Field Mapping Intents

States and views can declare how their fields are queried without coupling to any specific store. Each store translates the intent into its native mapping or indexing; stores with no equivalent ignore it, so state classes stay swappable between backends.

[ViewStore("my_app_search_products")]
public class ProductSearchView : IState
{
    public string Id { get; set; }

    [StoreField(FieldIntent.FullText)]      // Elastic: text. Mongo: text index.
    public string Description { get; set; }

    [StoreField(FieldIntent.FullText | FieldIntent.ExactMatch)]  // Elastic: text + .keyword. Mongo: text + ascending index.
    public string Name { get; set; }

    [StoreField(FieldIntent.ExactMatch)]    // Elastic: keyword. Mongo: ascending index.
    public string Sku { get; set; }

    [StoreField(FieldIntent.Sortable)]      // Elastic: native type. Mongo: ascending index.
    public decimal Price { get; set; }

    [StoreField(FieldIntent.NotQueried)]    // Elastic: not indexed. Mongo: no index.
    public ProductDetails Details { get; set; }
}
Intent Elasticsearch MongoDB (view stores only)
ExactMatch keyword (ignore_above: 8191) ascending index
FullText text text index (all FullText fields combined)
Sortable keyword for strings/enums/Guids, native type otherwise ascending index
FullText \| ExactMatch (or \| Sortable) text with a keyword sub-field named keyword text index + ascending index
NotQueried index: false / enabled: false none
[StoreField] without an intent inferred from the CLR type none
(not annotated) dynamic mapping (see below) none

Intents are flags and can be combined, except NotQueried, which must stand alone. The combined FullText | ExactMatch mapping matches Elasticsearch's dynamic mapping for strings, so queries, sorts and aggregations on field.keyword keep working.

Elasticsearch. Mappings are generated when an index is first created; existing indices are never altered (reindex to apply a new mapping). The Mapping property of [ElasticIndex] selects the behavior:

  • MappingBehavior.Auto (default): only annotated fields (and the objects containing them) are mapped statically. Every other field keeps dynamic mapping, so an unannotated string is still text with a .keyword sub-field. A state without annotations gets a fully dynamic index.
  • MappingBehavior.Static: every statically knowable field is mapped from the CLR type using conventions (strings become keyword with ignore_above: 8191, numbers and dates their native types), with intents applied where declared.
  • MappingBehavior.Dynamic: intents are ignored and Elasticsearch infers everything.

In every mode, fields whose shape is not statically knowable stay dynamic: object-typed and interface-typed properties, non-generic collections, dictionaries, framework and driver types (System.*, Microsoft.*, MongoDB.*, Elastic.*, e.g. BsonDocument, JsonElement), types with a [JsonConverter], and recursive or very deep object graphs. Unannotated numeric arrays (float[], List<double>, ...) are left to index templates or dynamic mapping, so a template can still map them as dense_vector. Note that a create-index request's mappings take precedence over composable index templates for the fields it maps.

MongoDB. Intents become indexes on view collections only. Snapshot collections are write-heavy and read by id, so secondary indexes there would cost writes without serving queries. Element names are resolved through the driver's class maps, so [BsonElement], registered class maps and conventions are honored (register custom class maps before stores are set up). Existing indexes are compared by key specification rather than name: an equivalent index under another name is reused, and conflicts (an index name reused for different keys, an existing text index over different fields, a changed TTL) are logged as warnings instead of failing setup. MongoDB allows one text index per collection, so changing the set of FullText fields requires dropping the old text index.

When a store's native features are needed, an escape hatch is available at registration. The hook receives the create-index request already populated with the generated mapping and the [ElasticIndex] settings; adjust them in place (assigning a new Mappings or Settings object discards the generated one):

x.AddElasticViewStore(cfg =>
{
    config.GetSection("ElasticSearch").Bind(cfg);
    cfg.ConfigureIndex<ProductSearchView>(request =>
    {
        request.Settings.NumberOfShards = 4;
        (request.Mappings.Properties ??= new Properties())
            .Add("embedding", new DenseVectorProperty { Dims = 384 });
    });
});

Docker Compose

Development infrastructure is provided in the compose/ directory:

# Start MongoDB
docker compose -f compose/MongoDb/docker-compose.yml up -d

# Start Elasticsearch
docker compose -f compose/ElasticSearch/docker-compose.yml up -d

# Start Pulsar
docker compose -f compose/Pulsar/docker-compose.yml up -d

# Start Ignite
docker compose -f compose/Ignite/docker-compose.yml up -d

Testing

The test suite uses xUnit with Bogus for data generation.

# Run unit tests only
dotnet test --filter "Category!=Integration"

# Run all tests (requires Docker services)
dotnet test

Writing Tests

Use the in-memory providers for fast, isolated unit tests:

services.AddEventHorizon(x =>
{
    x.AddInMemorySnapshotStore()
     .AddInMemoryViewStore()
     .AddInMemoryEventStream()
     .AddEventSourcing();
});

Integration tests use [Collection("Integration")] and require running Docker Compose services.

CI/CD

This project uses GitHub Actions (.github/workflows/ci.yml) with GitVersion for automatic semantic versioning based on the GitFlow branching model.

Versions are derived from git history and tags — no manual version bumping required after initial setup.

Branch/Tag Pre-release Label Example Version
v* tag (stable) 1.3.0
master / main (stable) 1.3.0
release/* rc 1.3.0-rc.3
hotfix/* hf 1.3.1-hf.1
develop preview 1.4.0-preview.12
feature/* {branch} 1.4.0-my-feature.1

How versioning works

  • Tag a release on main/master (e.g., v1.3.0) to set the version baseline
  • All subsequent commits on branches derive their version from git tags and merge history
  • Commit messages with +semver: major, +semver: minor, or +semver: fix control version increments
  • Configuration lives in GitVersion.yml at the repo root

All packages are published with the Cts.* prefix (e.g., Cts.EventHorizon.Abstractions).

Trusted Publishing

NuGet packages are published using trusted publishing via GitHub's OIDC tokens — no API keys or secrets required. The trusted publisher is configured on nuget.org to trust the ci.yml workflow in this repository.

Samples

Working examples are in the samples/ directory:

  • EventHorizon.EventSourcing.Samples — Full event sourcing example with accounts, commands, events, views, and subscriptions using MongoDB + Elasticsearch + Pulsar
  • EventHorizon.EventStreaming.Samples — Standalone streaming example with multi-topic subscription and publishing

Run samples with:

# Start required infrastructure
docker compose -f compose/MongoDb/docker-compose.yml up -d
docker compose -f compose/ElasticSearch/docker-compose.yml up -d
docker compose -f compose/Pulsar/docker-compose.yml up -d

# Run the event sourcing sample
dotnet run --project samples/EventHorizon.EventSourcing.Samples

Project Structure

EventHorizon/
├── src/
│   ├── EventHorizon.Abstractions/          # Core interfaces, models, attributes
│   ├── EventHorizon.EventStore/            # Store abstractions (ICrudStore, Lock)
│   ├── EventHorizon.EventStore.InMemory/   # In-memory store implementation
│   ├── EventHorizon.EventStore.MongoDb/    # MongoDB store implementation
│   ├── EventHorizon.EventStore.ElasticSearch/ # Elasticsearch store implementation
│   ├── EventHorizon.EventStore.Ignite/     # Apache Ignite store implementation
│   ├── EventHorizon.EventStreaming/        # Streaming abstractions
│   ├── EventHorizon.EventStreaming.InMemory/ # In-memory streaming
│   ├── EventHorizon.EventStreaming.Pulsar/ # Apache Pulsar streaming
│   └── EventHorizon.EventSourcing/        # Event sourcing orchestration
├── test/                                   # Unit and integration tests
├── samples/                                # Working example applications
├── benchmark/                              # Performance benchmarks
├── compose/                                # Docker Compose files for local dev
└── charts/                                 # Helm charts for Kubernetes deployment

Contributing

  1. Fork the repository
  2. Create a feature branch (feature/my-feature)
  3. Commit changes with clear messages
  4. Open a pull request against develop

License

This project is licensed under the MIT License.

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 is compatible.  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. 
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
1.6.0-preview.37 33 10/9/2026
1.6.0-preview.36 51 10/8/2026
1.6.0-preview.15 41 10/7/2026
1.6.0-preview.14 42 10/7/2026
1.6.0-preview.11 75 7/7/2026
1.6.0-preview.10 76 7/7/2026
1.6.0-preview.8 64 7/7/2026
1.6.0-preview.6 75 7/7/2026
1.6.0-preview.3 65 7/5/2026
1.6.0-preview.2 72 7/5/2026
1.6.0-preview.1 81 5/6/2026
1.6.0-preview.0 71 5/6/2026
1.5.0 1,331 5/6/2026
1.5.0-preview.1 74 4/28/2026
1.4.0 191 4/28/2026
1.4.0-preview.4 70 4/28/2026
1.4.0-preview.3 88 4/13/2026
1.4.0-preview.1 96 2/6/2026
1.4.0-preview.0 83 2/6/2026
1.3.2 1,591 2/6/2026