Msdiator.Core 1.0.0

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

Msdiator

A high-performance CQRS mediator library for .NET — built from the ground up with compiled expression trees, aggressive caching, and zero-overhead abstractions.

NuGet .NET


Why Msdiator?

Msdiator MediatR
Void Command 58 ns / 40 B 92 ns / 128 B
Command with Response 69 ns / 184 B 96 ns / 272 B
Publish (3 handlers) 120 ns / 224 B 227 ns / 664 B
100 Notifications 11,505 ns 21,962 ns
1000 Sequential Requests 60,960 ns / 136 KB 85,360 ns / 224 KB
Pipeline (1 Behavior) 167 ns 198 ns
Parallel Notification (4 handlers) 15.3 ms 60.5 ms (sequential only)

Up to 1.9x faster for notifications, 1.6x faster for commands, and 3.9x faster with parallel notification execution — all with significantly lower memory allocations.


Features

  • Full CQRS Support — Commands, Queries, Notifications with dedicated interfaces
  • Compiled Expression Trees — Handler factories and invokers compiled at startup, not reflected at runtime
  • Handler Caching — ConcurrentDictionary-based pipeline descriptor caching with warmup support
  • Pipeline Behaviors — Request and notification pipeline behaviors with priority ordering
  • Streaming — First-class IAsyncEnumerable<T> streaming support for commands and queries
  • Batch Processing — Send multiple requests in parallel with configurable concurrency, timeouts, and failure policies
  • Parallel Notifications — [Parallel] attribute for concurrent notification handler execution
  • Priority Ordering — [Priority] attribute to control handler and behavior execution order
  • Execution Groups — [HandlerExecutionGroup] to batch handlers into sequential groups with internal parallelism
  • Fire-and-Forget — PublishSync for channel-based background notification dispatch
  • Cache Warmup — IHostedService that pre-warms handler caches at application startup
  • Performance Tracking — Optional resolution time tracking with HandlerCacheStatistics
  • Multi-Target — Supports .NET 8, .NET 9, and .NET 10
  • Minimal Dependencies — Only Microsoft.Extensions.DependencyInjection.Abstractions, Hosting.Abstractions, and Logging.Abstractions

Installation

dotnet add package Msdiator.Core

The Msdiator.Contracts package is included automatically as a dependency.


Quick Start

1. Register Msdiator

builder.Services.AddMsdiator(msdiator =>
{
    msdiator.AddHandlersFromAssembly(typeof(Program).Assembly);
});

2. Define a Command

public record CreateUserCommand(string Name, string Email) : ICommand<int>;

public class CreateUserCommandHandler : ICommandHandler<CreateUserCommand, int>
{
    public async Task<int> Handle(CreateUserCommand request, CancellationToken cancellationToken)
    {
        // Save user to database...
        return userId;
    }
}

3. Send It

public class UsersController : ControllerBase
{
    private readonly IMsdiator _msdiator;

    public UsersController(IMsdiator msdiator) => _msdiator = msdiator;

    [HttpPost]
    public async Task<IActionResult> Create(CreateUserCommand command)
    {
        var userId = await _msdiator.Send(command);
        return Ok(userId);
    }
}

Usage Guide

Commands

Commands represent actions that change state. They can return a value or be void.

Command with response:

public record CreateOrderCommand(string Product, int Quantity) : ICommand<OrderResult>;

public class CreateOrderHandler : ICommandHandler<CreateOrderCommand, OrderResult>
{
    public async Task<OrderResult> Handle(CreateOrderCommand request, CancellationToken cancellationToken)
    {
        var order = new Order(request.Product, request.Quantity);
        await _repository.SaveAsync(order, cancellationToken);
        return new OrderResult(order.Id);
    }
}

Void command (no return value):

public record DeleteOrderCommand(Guid OrderId) : ICommand;

public class DeleteOrderHandler : ICommandHandler<DeleteOrderCommand>
{
    public async Task Handle(DeleteOrderCommand request, CancellationToken cancellationToken)
    {
        await _repository.DeleteAsync(request.OrderId, cancellationToken);
    }
}

Queries

Queries represent read operations that return data without side effects.

public record GetOrderByIdQuery(Guid OrderId) : IQuery<OrderDto>;

public class GetOrderByIdHandler : IQueryHandler<GetOrderByIdQuery, OrderDto>
{
    public async Task<OrderDto> Handle(GetOrderByIdQuery request, CancellationToken cancellationToken)
    {
        var order = await _repository.GetByIdAsync(request.OrderId, cancellationToken);
        return order.ToDto();
    }
}

Sending a query:

var order = await _msdiator.Send(new GetOrderByIdQuery(orderId));

Notifications

Notifications are dispatched to all registered handlers — ideal for domain events.

public record OrderCreatedEvent(Guid OrderId, DateTime CreatedAt) : INotification;

public class SendEmailOnOrderCreated : INotificationHandler<OrderCreatedEvent>
{
    public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
    {
        await _emailService.SendOrderConfirmation(notification.OrderId);
    }
}

public class UpdateAnalyticsOnOrderCreated : INotificationHandler<OrderCreatedEvent>
{
    public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
    {
        await _analytics.TrackOrderCreated(notification.OrderId);
    }
}

Publishing a notification (await all handlers):

await _msdiator.Publish(new OrderCreatedEvent(order.Id, DateTime.UtcNow));

Fire-and-forget (returns immediately, runs in background):

await _msdiator.PublishSync(new OrderCreatedEvent(order.Id, DateTime.UtcNow));

PublishSync uses System.Threading.Channels internally. Notifications are enqueued onto a singleton channel and processed by a dedicated BackgroundService inside a fresh DI scope. This avoids thread-pool starvation, DI scope-capture issues, and ensures graceful shutdown — all without any API or configuration change on your side.

Parallel Notifications

Mark a notification with [Parallel] to execute all handlers concurrently.

[Parallel(MaxDegreeOfParallelism = 4)]
public record OrderCreatedEvent(Guid OrderId) : INotification;

Without [Parallel], handlers execute sequentially. With it, all handlers run in parallel — up to 3.9x faster in benchmarks.

Parallel notifications and thread safety

[Parallel] runs all handlers for that notification at the same time, using the same DI scope as the current Publish / PublishSync execution. That is great for independent I/O, but unsafe if handlers depend on a non–thread-safe scoped service that is shared across those handlers — the usual example is a single EF Core DbContext per scope: two handlers must not use the same instance concurrently.

Safer options when you need parallelism and database (or other scoped) work:

  • Prefer sequential handlers (do not use [Parallel]) if handlers must share one scoped unit of work.
  • Use IDbContextFactory<TContext> and create a separate context per handler (or per operation).
  • Inject IServiceScopeFactory and create a new scope inside each handler so parallel handlers do not share scoped instances.
  • Keep [Parallel] only for handlers that do not share mutable scoped state (e.g. isolated HTTP calls via IHttpClientFactory).

This is a general .NET / DI concurrency concern, not a bug in Msdiator: parallel mode changes when handlers run, not how scoped services are shared.

Does this apply to PublishSync?

Yes. PublishSync eventually runs the same Publish pipeline (including [Parallel]) inside a background scope. Parallel handlers still share one scope for that single notification, so the same rules apply: do not use [Parallel] with handlers that unsafely share one DbContext (or similar) unless you use a factory or per-handler scope as above.

Priority Ordering

Control the order in which notification handlers execute using [Priority]. Higher values run first.

[Priority(100)]
public class CriticalHandler : INotificationHandler<OrderCreatedEvent>
{
    public Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
    {
        // Runs first
    }
}

[Priority(1)]
public class LowPriorityHandler : INotificationHandler<OrderCreatedEvent>
{
    public Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
    {
        // Runs last
    }
}

Execution Groups

Group notification handlers so groups execute sequentially while handlers within each group run in parallel.

[HandlerExecutionGroup("Validation")]
public class ValidateOrderHandler : INotificationHandler<OrderCreatedEvent> { ... }

[HandlerExecutionGroup("Validation")]
public class ValidateInventoryHandler : INotificationHandler<OrderCreatedEvent> { ... }

[HandlerExecutionGroup("Processing")]
public class ProcessPaymentHandler : INotificationHandler<OrderCreatedEvent> { ... }

[HandlerExecutionGroup("Processing")]
public class ReserveInventoryHandler : INotificationHandler<OrderCreatedEvent> { ... }

In this example, both Validation handlers run in parallel first, then both Processing handlers run in parallel.

Streaming

Stream results back using IAsyncEnumerable<T> for large datasets or real-time data.

Stream query:

public record GetLargeDatasetQuery(int PageSize) : IStreamQuery<DataRow>;

public class GetLargeDatasetHandler : IStreamHandler<GetLargeDatasetQuery, DataRow>
{
    public async IAsyncEnumerable<DataRow> Handle(
        GetLargeDatasetQuery request,
        [EnumeratorCancellation] CancellationToken cancellationToken)
    {
        await foreach (var row in _repository.StreamAllAsync(cancellationToken))
        {
            yield return row;
        }
    }
}

Consuming the stream:

await foreach (var row in _msdiator.SendStream(new GetLargeDatasetQuery(100)))
{
    Process(row);
}

Stream command:

public record ImportDataCommand(string FilePath) : IStreamCommand<ImportResult>;

Batch Processing

Send multiple requests in parallel with fine-grained control.

var requests = Enumerable.Range(1, 100)
    .Select(i => new ProcessItemCommand(i))
    .ToList();

var results = await _msdiator.SendBatch<ProcessItemCommand, ItemResult>(
    requests,
    new BatchOptions
    {
        MaxDegreeOfParallelism = 10,
        ContinueOnFailure = true,
        OperationTimeout = TimeSpan.FromSeconds(30),
        BatchTimeout = TimeSpan.FromMinutes(5)
    });

foreach (var result in results)
{
    if (result.IsSuccess)
        Console.WriteLine($"Item {result.Index}: {result.Response}");
    else
        Console.WriteLine($"Item {result.Index} failed: {result.Exception?.Message}");
}

Pipeline Behaviors

Pipeline behaviors wrap every request, enabling cross-cutting concerns like logging, validation, and transactions.

Request pipeline behavior:

public class LoggingBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
    where TRequest : IRequest<TResponse>
{
    private readonly ILogger<LoggingBehavior<TRequest, TResponse>> _logger;

    public LoggingBehavior(ILogger<LoggingBehavior<TRequest, TResponse>> logger)
        => _logger = logger;

    public async Task<TResponse> Handle(
        TRequest request,
        RequestHandlerDelegate<TResponse> next,
        CancellationToken cancellationToken)
    {
        _logger.LogInformation("Handling {RequestType}", typeof(TRequest).Name);
        var response = await next();
        _logger.LogInformation("Handled {RequestType}", typeof(TRequest).Name);
        return response;
    }
}

Notification pipeline behavior:

public class NotificationLoggingBehavior<TNotification> : INotificationPipelineBehavior<TNotification>
    where TNotification : INotification
{
    public async Task Handle(
        TNotification notification,
        NotificationHandlerDelegate next,
        CancellationToken cancellationToken)
    {
        Console.WriteLine($"Before: {typeof(TNotification).Name}");
        await next();
        Console.WriteLine($"After: {typeof(TNotification).Name}");
    }
}

Registering behaviors:

builder.Services.AddMsdiator(msdiator =>
{
    msdiator.AddHandlersFromAssembly(typeof(Program).Assembly);
    msdiator.AddPipelineBehavior<LoggingBehavior<,>>();
    msdiator.AddNotificationBehavior<NotificationLoggingBehavior<>>();
});

Priority on behaviors:

[Priority(100)]
public class ValidationBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
    where TRequest : IRequest<TResponse>
{
    // Runs before lower-priority behaviors
}

Configuration

Full Builder Options

builder.Services.AddMsdiator(msdiator =>
{
    // Register handlers from assemblies
    msdiator.AddHandlersFromAssembly(typeof(Program).Assembly);
    msdiator.AddHandlersFromAssemblies(assembly1, assembly2);

    // Pipeline behaviors
    msdiator.AddPipelineBehavior<LoggingBehavior<,>>();
    msdiator.AddPipelineBehavior<ValidationBehavior<,>>();
    msdiator.AddNotificationBehavior<NotificationLoggingBehavior<>>();

    // Core options
    msdiator.ConfigureCore(options =>
    {
        options.HandlerLifetime = ServiceLifetime.Scoped;
        options.EnableHandlerCaching = true;
        options.EnablePerformanceTracking = false;
        options.MaxConcurrentRequests = 100;
        options.MaxConcurrentNotifications = 50;
        options.DefaultTimeout = TimeSpan.FromSeconds(30);
    });
});

Configuration Options

Option Default Description
HandlerLifetime Transient DI lifetime for handlers (Transient, Scoped, Singleton)
EnableHandlerCaching true Cache compiled handler factories and pipeline descriptors
EnablePerformanceTracking false Track handler resolution times for diagnostics
MaxConcurrentRequests 100 Max concurrent batch request executions
MaxConcurrentNotifications 50 Max concurrent parallel notification handlers
DefaultTimeout 30s Default timeout for operations

Cache Management

// Pre-warm caches for specific types
await _msdiator.WarmHandlerCache<CreateUserCommand, int>();
await _msdiator.WarmHandlerCache(typeof(OrderCreatedEvent));

// Get cache statistics
var stats = _msdiator.GetHandlerCacheStatistics();
Console.WriteLine($"Cached entries: {stats.CachedEntries}");
Console.WriteLine($"Avg resolution: {stats.AverageResolutionTimeMicroseconds} µs");

// Clear all caches
await _msdiator.ClearHandlerCache();

When EnableHandlerCaching is true, a hosted service automatically warms all handler caches at application startup.


Benchmarks

All benchmarks run with BenchmarkDotNet on .NET 10, comparing Msdiator against MediatR.

Request / Command Dispatch

Benchmark Msdiator MediatR Faster Less Alloc
Void Command 58.49 ns / 40 B 92.17 ns / 128 B 1.58x 3.2x
Command with Response 69.07 ns / 184 B 95.74 ns / 272 B 1.39x 1.5x
Simple Request 77.65 ns / 208 B 99.92 ns / 272 B 1.29x 1.3x
Complex Request 105.58 ns / 296 B 126.63 ns / 360 B 1.20x 1.2x
1000 Sequential Requests 60,960 ns / 136 KB 85,360 ns / 224 KB 1.40x 1.6x

Notification Dispatch

Benchmark Msdiator MediatR Faster Less Alloc
Publish (3 handlers) 120.0 ns / 224 B 227.2 ns / 664 B 1.89x 2.96x
Fire-and-Forget Publish 119.9 ns / 232 B N/A — —
Priority Ordered Notification 168.1 ns / 376 B N/A — —
100 Notifications 11,505 ns 21,962 ns 1.91x —
Parallel (4 handlers) 15.3 ms / 1,584 B 60.5 ms / 2,250 B 3.95x 1.4x

Pipeline Behavior Overhead

Benchmark Msdiator MediatR Faster Less Alloc
No Pipeline 79.14 ns / 232 B 104.41 ns / 296 B 1.32x 1.28x
1 Behavior Pipeline 167.57 ns / 608 B 198.56 ns / 552 B 1.19x —
Pipeline Overhead (1000 calls) 140,572 ns / 480 KB 183,431 ns / 5,360 KB 1.30x 11.2x

Streaming (IAsyncEnumerable)

Benchmark (10,000 items) Msdiator MediatR Faster Less Alloc
Stream Early Cancel 1.997 µs / 456 B 3.529 µs / 864 B 1.77x 1.9x
Stream Large Objects 151.56 µs / 1,128 KB 153.78 µs / 1,128 KB 1.01x ~1.0x
Stream Integers 192.58 µs / 520 B 421.77 µs / 928 B 2.19x 1.8x

How Msdiator Achieves Better Performance

1. Compiled Expression Trees

MediatR resolves handlers through reflection and Activator.CreateInstance at runtime. Msdiator compiles handler factories and invokers into strongly-typed delegates using System.Linq.Expressions at startup, eliminating per-request reflection overhead.

2. Typed Invokers (No Boxing)

Handler invocations return Task<TResponse> directly — no intermediate Task<object> boxing/unboxing. Stream handlers return IAsyncEnumerable<TResponse> natively.

3. Aggressive Pipeline Descriptor Caching

The entire pipeline (handler factory, invoker, behavior chain) is resolved once and cached in ConcurrentDictionary. Subsequent calls hit a direct dictionary lookup with zero allocation.

4. Hot-Path Optimizations

  • [MethodImpl(MethodImplOptions.AggressiveInlining)] on critical dispatch paths
  • Fast task.IsCompletedSuccessfully checks to skip async state machines when possible
  • Separate no-behavior and with-behavior execution paths to avoid unnecessary delegate allocations
  • Pre-computed sorted indices and execution group batches for notifications

5. Startup Cache Warmup

An IHostedService pre-warms all handler caches during application startup, so the first request pays zero compilation cost.


Msdiator vs MediatR — Feature Comparison

Feature Msdiator MediatR
Commands / Queries / Requests Yes Yes
Notifications Yes Yes
Pipeline Behaviors Yes Yes
Notification Pipeline Behaviors Yes No
Streaming (IAsyncEnumerable) Yes Yes
Batch Processing Yes No
Parallel Notifications Yes ([Parallel]) No
Priority Ordering Yes ([Priority]) No
Execution Groups Yes ([HandlerExecutionGroup]) No
Fire-and-Forget Publish Yes (PublishSync) No
Handler Cache Warmup Yes (automatic IHostedService) No
Cache Statistics / Diagnostics Yes (HandlerCacheStatistics) No
Compiled Expression Trees Yes No
Typed Invokers (no boxing) Yes No
Performance Tracking Yes (opt-in) No
.NET 8 / 9 / 10 Yes Yes

API Reference

IMsdiator

// Send a request/command/query and get a typed response
Task<TResponse> Send<TResponse>(IRequest<TResponse> request, CancellationToken ct = default);

// Send a void request/command
Task Send(IRequest request, CancellationToken ct = default);

// Stream results via IAsyncEnumerable
IAsyncEnumerable<TResponse> SendStream<TResponse>(IStreamRequest<TResponse> request, CancellationToken ct = default);

// Batch-send requests with parallel execution
Task<IEnumerable<BatchResult<TResponse>>> SendBatch<TRequest, TResponse>(
    IEnumerable<TRequest> requests, BatchOptions? options = null, CancellationToken ct = default)
    where TRequest : IRequest<TResponse>;

// Publish notification — awaits all handlers
Task Publish<TNotification>(TNotification notification, CancellationToken ct = default)
    where TNotification : INotification;

// Publish notification — fire-and-forget (returns immediately)
Task PublishSync<TNotification>(TNotification notification, CancellationToken ct = default)
    where TNotification : INotification;

// Cache management
Task WarmHandlerCache<TRequest, TResponse>() where TRequest : IRequest<TResponse>;
Task ClearHandlerCache();
HandlerCacheStatistics GetHandlerCacheStatistics();

Contract Interfaces

Interface Purpose
ICommand Void command (state mutation, no return value)
ICommand<TResponse> Command with a return value
IQuery<TResponse> Query (read-only, always returns data)
INotification Notification / domain event
IStreamQuery<TResponse> Streaming query returning IAsyncEnumerable<TResponse>
IStreamCommand<TResponse> Streaming command returning IAsyncEnumerable<TResponse>
ICommandHandler<TCommand> Handler for void commands
ICommandHandler<TCommand, TResponse> Handler for commands with response
IQueryHandler<TQuery, TResponse> Handler for queries
INotificationHandler<TNotification> Handler for notifications
IStreamHandler<TRequest, TResponse> Handler for streaming requests
IPipelineBehavior<TRequest, TResponse> Request pipeline behavior
INotificationPipelineBehavior<TNotification> Notification pipeline behavior

Attributes

Attribute Target Purpose
[Parallel(MaxDegreeOfParallelism, ...)] Notification class Execute handlers in parallel
[Priority(int)] Handler / Behavior class Control execution order (higher = first)
[HandlerExecutionGroup(string)] Handler class Group handlers for sequential group execution

Extensions

Msdiator exposes an internal extension point that allows external packages to register additional services during the AddMsdiator setup. This is done through the IMsdiatorBuilderInternal interface and its AddExtensions method.

How It Works

When AddMsdiator runs, it executes all registered feature extensions after core services and handlers are registered. This gives extensions full access to the IServiceCollection while preserving the standard registration order.

Writing an Extension Package

Create an extension method on IMsdiatorBuilder that casts to IMsdiatorBuilderInternal to register services:

public static class MsdiatorRedisExtensions
{
    public static IMsdiatorBuilder AddRedisCaching(
        this IMsdiatorBuilder builder,
        Action<RedisCacheOptions> configure)
    {
        if (builder is IMsdiatorBuilderInternal internalBuilder)
        {
            internalBuilder.AddExtensions("RedisCaching", services =>
            {
                var options = new RedisCacheOptions();
                configure(options);
                services.AddSingleton(options);
                services.AddSingleton<IRedisCacheService, RedisCacheService>();
            });
        }
        return builder;
    }
}

Consuming an Extension

Extensions integrate seamlessly into the builder chain:

builder.Services.AddMsdiator(msdiator =>
{
    msdiator.AddHandlersFromAssembly(typeof(Program).Assembly);

    // Extension packages plug in naturally
    msdiator.AddRedisCaching(options =>
    {
        options.ConnectionString = "localhost:6379";
        options.DefaultExpiration = TimeSpan.FromMinutes(5);
    });
});

Extension Guidelines

  • Each extension must provide a unique feature name (first argument to AddExtensions). Duplicate names will throw at startup.
  • Extensions receive the full IServiceCollection, so they can register any services, hosted services, or options they need.
  • The IMsdiatorBuilderInternal cast is intentional — it keeps the core IMsdiatorBuilder API clean while giving extension authors full flexibility.

License

Copyright 2026 Mohamed Sayed

Licensed under the Apache License, Version 2.0


Author

Mohamed Sayed


Built for developers who demand the fastest possible mediator for production .NET applications.

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.0.0 128 4/20/2026