Donakunn.MessagingOverQueue 0.2.5-alpha

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

MessagingOverQueue

GitHub Repository .NET 10 License

Async messaging library for .NET 10 with Redis Streams. Implement IMessageHandler<T> and the library creates streams, consumer groups, and consumers automatically.

  • Handler auto-discovery via assembly scanning
  • Reflection-free handler dispatch with O(1) lookup
  • Middleware pipeline (retry, circuit breaker, timeout, idempotency)
  • Outbox pattern with SQL Server provider
  • Delayed message delivery — schedule events and commands for future processing
  • Scoped DI — each message gets its own handler instance

Installation

dotnet add package Donakunn.MessagingOverQueue.RedisStreams

Quick Start

1. Define a message

using Donakunn.MessagingOverQueue.Abstractions.Messages;

public record OrderCreatedEvent : Event
{
    public Guid OrderId { get; init; }
    public string CustomerId { get; init; } = string.Empty;
}

2. Create a handler

using Donakunn.MessagingOverQueue.Abstractions.Consuming;
using Donakunn.MessagingOverQueue.Abstractions.Messages;

public class OrderCreatedHandler : IMessageHandler<OrderCreatedEvent>
{
    private readonly ILogger<OrderCreatedHandler> _logger;

    public OrderCreatedHandler(ILogger<OrderCreatedHandler> logger)
    {
        _logger = logger;
    }

    public Task HandleAsync(
        OrderCreatedEvent message,
        IMessageContext context,
        CancellationToken cancellationToken = default)
    {
        _logger.LogInformation("Order {OrderId} created", message.OrderId);
        return Task.CompletedTask;
    }
}

3. Configure services

using Donakunn.MessagingOverQueue.RedisStreams.DependencyInjection;
using Donakunn.MessagingOverQueue.Topology.DependencyInjection;

services.AddRedisStreamsMessaging(options => options
    .UseConnectionString("localhost:6379")
    .WithStreamPrefix("myapp"))
    .AddTopology(topology => topology
        .WithServiceName("order-service")
        .ScanAssemblyContaining<OrderCreatedHandler>())
    .AddRedisStreamsConsumerHostedService();

The library automatically discovers your handlers, creates Redis streams and consumer groups, registers handlers in DI with scoped lifetime, and starts consuming.

Publishing

Inject IEventPublisher or ICommandSender:

public class OrderController(IEventPublisher publisher) : ControllerBase
{
    [HttpPost]
    public async Task<IActionResult> CreateOrder(CreateOrderRequest request)
    {
        await publisher.PublishAsync(new OrderCreatedEvent
        {
            OrderId = Guid.NewGuid(),
            CustomerId = request.CustomerId
        });

        return Accepted();
    }
}

Delayed Publishing

Schedule a message for future delivery by passing a TimeSpan delay. Requires the outbox to be configured.

// Deliver in 10 minutes
await publisher.PublishAsync(new OrderReminderEvent { OrderId = id }, TimeSpan.FromMinutes(10));

// Send a command after a delay
await sender.SendAsync(new ExpireSessionCommand { UserId = userId }, TimeSpan.FromHours(1));

Calling the delay overload without the outbox configured throws NotSupportedException immediately, making misconfiguration explicit at the call site.

Features

Handler Discovery — TopologyScanner finds all IMessageHandler<T> implementations at startup. Streams, consumer groups, and consumers are created based on conventions or attributes. Override the consumer group with [RedisConsumerGroup("custom-group")].

Reflection-Free Dispatch — HandlerInvoker<T> instances are created once at startup and cached in a ConcurrentDictionary. Runtime dispatch is a dictionary lookup + direct method call.

Middleware Pipeline — Extensible pipeline for both publishing and consuming. Built-in consume middlewares (in execution order):

Order Middleware Purpose
100 CircuitBreaker Fail fast when downstream is unhealthy
200 Retry Automatic retry with exponential backoff
300 Timeout Cancel long-running handlers
400 Logging Structured logging
500 Idempotency Duplicate detection via inbox pattern
600 Deserialization JSON to strongly-typed message

Resilience — Retry, circuit breaker, and timeout powered by Polly v8. Configure via UseResilience:

builder.UseResilience(r => r
    .WithRetry(opts => opts.MaxRetryAttempts = 5)
    .WithCircuitBreaker(opts => opts.FailureThreshold = 10)
    .WithTimeout(TimeSpan.FromSeconds(30)));

Outbox Pattern — Reliable message delivery with SQL Server (ADO.NET, no EF Core dependency). Supports partition-based horizontal scaling with multiple workers, and delayed delivery via ScheduledAt. The EnsureSchemaAsync migration is idempotent — existing tables gain the ScheduledAt column automatically. Configure via UsePersistence:

builder.UsePersistence(p => p
    .WithOutbox(opts => opts.BatchSize = 50)
        .UseSqlServer(connectionString)
    .WithIdempotency());

Dead Letter Handling — Messages exceeding max delivery attempts are moved to DLQ streams. Configurable per-stream or per-consumer-group.

Stream Retention — Time-based (MINID trimming) or count-based (MAXLEN trimming).

Health Checks — Built-in ASP.NET Core health check. Reports connection status, latency, and Redis server version.

Configuration — Fluent API, appsettings.json (section: "RedisStreams"), or both.

Full Configuration Example

services.AddMessaging()
    .UseRedisStreamsQueues(queues => queues
        .WithConnection(opts => opts
            .UseConnectionString("localhost:6379")
            .WithStreamPrefix("myapp")
            .ConfigureConsumer(batchSize: 20, maxPendingMessages: 1000)
            .ConfigureClaiming(claimIdleTime: TimeSpan.FromMinutes(5))
            .WithTimeBasedRetention(TimeSpan.FromDays(7))
            .WithDeadLetterPerConsumerGroup(maxDeliveryAttempts: 5))
        .WithTopology(t => t
            .WithServiceName("order-service")
            .ScanAssemblyContaining<OrderCreatedHandler>())
        .WithHealthChecks())
    .UseResilience(r => r
        .WithRetry(opts => opts.MaxRetryAttempts = 5)
        .WithCircuitBreaker()
        .WithTimeout(TimeSpan.FromSeconds(30)))
    .UsePersistence(p => p
        .WithOutbox(opts => opts.BatchSize = 50)
            .UseSqlServer(connectionString)
        .WithIdempotency());

Requirements

  • .NET 10+
  • Redis 6.2+ (for XAUTOCLAIM)
  • SQL Server 2016+ (outbox provider, optional)

License

Apache 2.0

Contributing

Contributions are welcome. Visit the GitHub repository to report issues, submit pull requests, or request features.

Product 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. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

NuGet packages (1)

Showing the top 1 NuGet packages that depend on Donakunn.MessagingOverQueue:

Package Downloads
Donakunn.MessagingOverQueue.RedisStreams

Redis Streams messaging provider for MessagingOverQueue. Provides high-performance message streaming using Redis Streams with consumer groups, automatic message claiming, and dead letter handling.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
0.2.5-alpha 106 3/2/2026
0.2.4-alpha 78 3/2/2026
0.2.3-alpha 91 2/27/2026
0.2.2-alpha 91 2/27/2026
0.2.1-alpha 86 2/27/2026
0.2.0-alpha 82 2/26/2026
0.1.3-alpha 276 1/27/2026
0.1.2-alpha 85 1/27/2026
0.1.1-alpha 119 1/23/2026
0.0.10-alpha 135 1/18/2026
0.0.9-alpha 91 1/15/2026
0.0.8-alpha 92 1/15/2026
0.0.7-alpha 94 1/14/2026
0.0.6-alpha 91 1/10/2026
0.0.5-alpha 96 1/9/2026
0.0.4-alpha 97 1/8/2026
0.0.3-alpha 93 1/8/2026
0.0.2-alpha 85 1/7/2026
0.0.1-alpha 84 1/7/2026