A modern, high-performance CQS/CQRS framework for .NET with pipeline behaviors, notifications, streaming, sagas, and distributed system support. MIT-licensed alternative to MediatR.
β
Commands & Queries - Clear separation of read and write operations
β
Pipeline Behaviors - Cross-cutting concerns (logging, validation, etc.)
β
Notifications - Publish-subscribe pattern with multiple handlers
β
Streaming - IAsyncEnumerable support for large datasets
β
Sagas - Multi-step transactions with compensation
β
Outbox Pattern - Reliable message publishing
β
Job Scheduling - Deferred command execution
β
Paged Queries - Built-in pagination support
β
Modern .NET - Targets .NET 6, 8, and 9
Package
Description
NuGet
eQuantic.Core.CQS
Core framework
eQuantic.Core.CQS.Abstractions
Interfaces and contracts
eQuantic.Core.CQS.Generators
Source generators
Package
Provider
Features
eQuantic.Core.CQS.Redis
Redis
Saga Repository, Outbox, Job Scheduler
eQuantic.Core.CQS.MongoDb
MongoDB
Saga Repository, Outbox, Job Scheduler
eQuantic.Core.CQS.PostgreSql
PostgreSQL
Saga Repository, Outbox, Job Scheduler
eQuantic.Core.CQS.EntityFramework
EF Core
Saga Repository, Outbox, Job Scheduler
Package
Provider
Features
eQuantic.Core.CQS.Azure
Azure Service Bus
Outbox Publisher (Queue/Topic)
eQuantic.Core.CQS.AWS
Amazon SQS
Outbox Publisher
Package
Provider
Features
eQuantic.Core.CQS.OpenTelemetry
OpenTelemetry
Distributed tracing, metrics
eQuantic.Core.CQS.ApplicationInsights
Azure App Insights
Distributed tracing, metrics
eQuantic.Core.CQS.Datadog
Datadog APM
Distributed tracing
Package
Provider
Features
eQuantic.Core.CQS.Resilience
Default
Saga timeout, compensation
eQuantic.Core.CQS.Polly
Polly
Retry, circuit breaker
eQuantic.Core.CQS.Resilience.Redis
Redis
Dead letter queue
eQuantic.Core.CQS.Resilience.ServiceBus
Azure Service Bus
Dead letter queue
# Core package
dotnet add package eQuantic.Core.CQS
# Optional: Provider packages
dotnet add package eQuantic.Core.CQS.Redis
dotnet add package eQuantic.Core.CQS.Azure
// Basic setup
services . AddCQS ( options => options
. FromAssemblyContaining < Program > ( ) ) ;
// With providers
services . AddCQS ( options => options
. FromAssemblyContaining < Program > ( )
. UsePreProcessor = true ;
. UseRedis ( redis => redis . ConnectionString = "localhost:6379" )
. UseAzureServiceBus ( sb => {
sb . ConnectionString = "Endpoint=sb://..." ;
sb . QueueOrTopicName = "outbox" ;
} ) ) ;
public record GetUserByIdQuery ( Guid Id ) : IQuery < UserDto > ;
public class GetUserByIdHandler : IQueryHandler < GetUserByIdQuery , UserDto >
{
public async Task < UserDto > Execute ( GetUserByIdQuery query , CancellationToken ct )
{
return new UserDto ( query . Id , "John Doe" ) ;
}
}
public class UsersController : ControllerBase
{
private readonly IMediator _mediator ;
[ HttpGet ( "{id}" ) ]
public async Task < UserDto > Get ( Guid id )
=> await _mediator . ExecuteAsync ( new GetUserByIdQuery ( id ) ) ;
}
// Command without result
public record DeleteUserCommand ( Guid Id ) : ICommand ;
// Command with result
public record CreateUserCommand ( string Name ) : ICommand < Guid > ;
public record UserCreatedNotification ( Guid UserId ) : INotification ;
public class SendWelcomeEmailHandler : INotificationHandler < UserCreatedNotification >
{
public async Task Handle ( UserCreatedNotification notification , CancellationToken ct )
{
// Send email
}
}
// Publish to all handlers
await _notificationPublisher . Publish ( new UserCreatedNotification ( userId ) ) ;
public record GetAllUsersStreamQuery : IStreamQuery < UserDto > ;
await foreach ( var user in _mediator . ExecuteStreamAsync ( new GetAllUsersStreamQuery ( ) ) )
{
Console . WriteLine ( user . Name ) ;
}
Sagas (Multi-step Transactions)
public class OrderSaga : Saga < OrderSagaData >
{
protected override void ConfigureSteps ( )
{
Step ( "ProcessPayment" ,
execute : async ( data , ct ) => { /* charge */ } ,
compensate : async ( data , ct ) => { /* refund */ } ) ;
Step ( "ReserveInventory" ,
execute : async ( data , ct ) => { /* reserve */ } ,
compensate : async ( data , ct ) => { /* release */ } ) ;
Step ( "Ship" ,
execute : async ( data , ct ) => { /* ship */ } ,
compensate : async ( data , ct ) => { /* cancel shipment */ } ) ;
}
}
// Execute - automatic compensation on failure
var result = await saga . Execute ( new OrderSagaData { OrderId = orderId } ) ;
if ( ! result . IsSuccess )
{
Console . WriteLine ( $ "Saga failed: { result . Error ? . Message } ") ;
}
Azure Service Bus Integration
services . AddCQSAzureServiceBus ( options =>
{
options . ConnectionString = "Endpoint=sb://..." ;
options . QueueOrTopicName = "outbox-messages" ;
options . UseTopic = false ;
} ) ;
// Publish outbox messages
await _outboxPublisher . PublishAsync ( message ) ;
await _outboxPublisher . PublishBatchAsync ( messages ) ;
services . AddCQSAwsSqs ( options =>
{
options . QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789/my-queue" ;
options . Region = "us-east-1" ;
} ) ;
public class LoggingBehavior < TRequest , TResponse > : IPipelineBehavior < TRequest , TResponse >
{
public async Task < TResponse > Execute ( TRequest request , CancellationToken ct , HandlerDelegate < TResponse > next )
{
_logger . LogInformation ( "Handling {Request}" , typeof ( TRequest ) . Name ) ;
var response = await next ( ) ;
_logger . LogInformation ( "Handled {Request}" , typeof ( TRequest ) . Name ) ;
return response ;
}
}
Feature
eQuantic.Core.CQS
MediatR
License
MIT β
Commercial π°
Sagas
Built-in β
β
Outbox Pattern
Built-in β
β
Cloud Messaging
Azure/AWS β
β
Paged Queries
Built-in β
Manual
Streaming
IAsyncEnumerable β
IAsyncEnumerable
Notifications
β
β
Pipeline Behaviors
β
β
MIT License - See LICENSE for details.