LowCodeHub.InboxOutbox 0.0.15

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

LowCodeHub.InboxOutbox

A DI-first inbox/outbox processing library for ASP.NET Core with resilient background workers, retry policies, graceful shutdown, and clean repository abstractions. Inbox deduplicates incoming messages; outbox guarantees outgoing messages survive crashes — both backed by SQL Server or PostgreSQL.

NuGet License: MIT

Why This Library?

Feature LowCodeHub.InboxOutbox Manual Implementation MassTransit Outbox
Inbox deduplication Built-in (Source, ExternalMessageId) Manual unique constraint Separate config
Outbox guarantee Persisted before publish Build from scratch EF Core integration
Retry strategy Exponential backoff (configurable) Manual Built-in
Graceful shutdown Drain + timeout Manual CancellationToken Built-in
Worker health checks Auto-registered liveness Manual Manual
Storage backends SQL Server + PostgreSQL One custom backend EF Core only
Infrastructure Your existing DB Your existing DB EF Core required
Multi-instance safe Pessimistic locking Manual EF Core row locking

Installation

dotnet add package LowCodeHub.InboxOutbox

Quick Start

builder.Services
    .AddInboxOutbox(builder.Configuration)
    .AddInboxOutboxSqlServer(builder.Configuration);
// Enqueue an inbox message (from your RabbitMQ consumer)
await inboxRepository.TryAddAsync(new InboxMessage
{
    Source = "orders-service",
    ExternalMessageId = "rabbitmq:order-123-created",
    MessageType = "order.validate",
    Payload = "{\"orderId\":\"123\"}"
}, ct);

// Enqueue an outbox message (from your business transaction)
await outboxRepository.TryAddAsync(new OutboxMessage
{
    Destination = "orders.exchange",
    MessageType = "order.created",
    Payload = "{\"orderId\":\"123\"}"
}, ct);

That's it. Background workers poll for pending messages, dispatch them to your handlers/senders with automatic retries and exponential backoff, and mark them as processed — surviving crashes, restarts, and deployment rollouts.


Table of Contents


How It Works

┌──────────────────────────────────────────────────────────────┐
│  INBOX FLOW                                                  │
│                                                              │
│  RabbitMQ Consumer / API                                     │
│  └── inboxRepository.TryAddAsync(message)                    │
│       ├── true  → new row inserted → ack                     │
│       └── false → duplicate (Source, ExternalMessageId) → ack│
└─────────────────────┬────────────────────────────────────────┘
                      │
                      ▼
┌──────────────────────────────────────────────────────────────┐
│  InboxWorker (BackgroundService, polls every PollInterval)   │
│  1. Claim batch of pending messages (pessimistic lock)       │
│  2. Dispatch to IInboxMessageHandler (by MessageType)        │
│     ├── Success → MarkProcessed                              │
│     └── Failure → MarkFailed + exponential backoff retry     │
│         └── Max retries exhausted → permanently Failed       │
└──────────────────────────────────────────────────────────────┘

┌──────────────────────────────────────────────────────────────┐
│  OUTBOX FLOW                                                 │
│                                                              │
│  Your Business Transaction                                   │
│  └── outboxRepository.TryAddAsync(message)                   │
│      (same DB transaction as domain state changes)           │
└─────────────────────┬────────────────────────────────────────┘
                      │
                      ▼
┌──────────────────────────────────────────────────────────────┐
│  OutboxWorker (BackgroundService, polls every PollInterval)  │
│  1. Claim batch of pending messages (pessimistic lock)       │
│  2. Send via IOutboxMessageSender                            │
│     ├── Success → MarkProcessed                              │
│     └── Failure → MarkFailed + exponential backoff retry     │
│         └── Max retries exhausted → permanently Failed       │
└──────────────────────────────────────────────────────────────┘

Key design principles:

  • Inbox — deduplicates incoming messages by (Source, ExternalMessageId). Your RabbitMQ consumer stays thin: persist → ack. The worker handles business logic with retries.
  • Outbox — persists outgoing messages inside your business transaction. The worker publishes later, ensuring no event is lost even if the process crashes after DB commit.
  • Pessimistic locking — workers claim rows atomically using UPDLOCK, READPAST (SQL Server) or FOR UPDATE SKIP LOCKED (PostgreSQL). Safe for multi-instance deployments.

Configuration

appsettings.json

{
  "InboxOutbox": {
    "Inbox": {
      "Enabled": true,
      "BatchSize": 50,
      "MaxParallelism": 4,
      "PollInterval": "00:00:02",
      "InitialRetryDelay": "00:00:05",
      "MaxRetryDelay": "00:05:00",
      "DrainOnShutdown": true,
      "DrainTimeout": "00:00:30",
      "HealthStaleAfter": "00:02:00"
    },
    "Outbox": {
      "Enabled": true,
      "BatchSize": 50,
      "MaxParallelism": 4,
      "PollInterval": "00:00:02",
      "InitialRetryDelay": "00:00:05",
      "MaxRetryDelay": "00:05:00",
      "DrainOnShutdown": true,
      "DrainTimeout": "00:00:30",
      "HealthStaleAfter": "00:02:00"
    }
  }
}

Worker Options (shared by Inbox and Outbox)

Option Default Description
Enabled true Enable/disable the worker
BatchSize 50 Messages claimed per poll cycle
MaxParallelism 4 Concurrent handler/sender invocations per batch
PollInterval 2s Time between polling cycles
InitialRetryDelay 5s Delay before first retry
MaxRetryDelay 5min Maximum delay between retries
DrainOnShutdown true Wait for in-flight batch on shutdown
DrainTimeout 30s Maximum drain wait time
HealthStaleAfter 2min Worker heartbeat staleness threshold

Code-Based Configuration

builder.Services.AddInboxOutbox(options =>
{
    options.Inbox.BatchSize = 100;
    options.Inbox.MaxParallelism = 8;
    options.Outbox.PollInterval = TimeSpan.FromSeconds(1);
});

Implementing Contracts

Inbox Message Handler

Implement IInboxMessageHandler to process incoming messages. Register multiple handlers — each decides which message types it can handle:

public sealed class OrderCreatedInboxHandler : IInboxMessageHandler
{
    public bool CanHandle(string messageType)
        => messageType.Equals("order.created", StringComparison.OrdinalIgnoreCase);

    public async Task HandleAsync(InboxMessage message, CancellationToken cancellationToken)
    {
        var order = JsonSerializer.Deserialize<OrderPayload>(message.Payload);
        // Execute business logic...
    }
}

Outbox Message Sender

Implement IOutboxMessageSender to publish outgoing messages to your broker or API:

public sealed class RabbitOutboxSender : IOutboxMessageSender
{
    private readonly IMessagePublisher _publisher;

    public RabbitOutboxSender(IMessagePublisher publisher) => _publisher = publisher;

    public async Task SendAsync(OutboxMessage message, CancellationToken cancellationToken)
    {
        await _publisher.PublishAsync(message.Destination, message.MessageType,
            message.Payload, cancellationToken: cancellationToken);
    }
}

Enqueuing Messages

Inbox (from your message consumer)

var inserted = await inboxRepository.TryAddAsync(new InboxMessage
{
    Source = "orders-service",
    ExternalMessageId = "rabbitmq:order-123-created",
    MessageType = "order.validate",
    Payload = "{\"orderId\":\"123\"}",
    CorrelationId = Guid.NewGuid().ToString("N")
}, ct);

// false = duplicate (Source, ExternalMessageId) — already ingested

Outbox (from your business transaction)

await outboxRepository.TryAddAsync(new OutboxMessage
{
    Destination = "orders.exchange",
    MessageType = "order.created",
    Payload = "{\"orderId\":\"123\"}",
    CorrelationId = Guid.NewGuid().ToString("N")
}, ct);

Inbox: RabbitMQ consumer stays thin

  1. Read broker message metadata and payload.
  2. Map to inbox fields: Source, ExternalMessageId, MessageType, Payload, CorrelationId.
  3. Call TryAddAsync(...).
  4. ack after DB call succeeds (true = new, false = duplicate — both ack).
  5. Do not execute business logic in the consumer — the worker/handler does that.

Outbox: persist inside business transaction

  1. Save domain state changes.
  2. Add outbox row via TryAddAsync(...).
  3. Commit transaction.
  4. The outbox worker publishes later — guaranteed delivery.

Storage Backends

SQL Server

builder.Services
    .AddInboxOutbox(builder.Configuration)
    .AddInboxOutboxSqlServer(builder.Configuration);
{
  "InboxOutbox": {
    "SqlServer": {
      "ConnectionString": "Server=...;Database=...;",
      "InboxSchema": "dbo",
      "InboxTable": "InboxMessages",
      "OutboxSchema": "dbo",
      "OutboxTable": "OutboxMessages",
      "LeaseDuration": "00:02:00"
    }
  }
}

PostgreSQL

builder.Services
    .AddInboxOutbox(builder.Configuration)
    .AddInboxOutboxPostgreSql(builder.Configuration);
{
  "InboxOutbox": {
    "PostgreSql": {
      "ConnectionString": "Host=...;Database=...;",
      "InboxSchema": "public",
      "InboxTable": "inbox_messages",
      "OutboxSchema": "public",
      "OutboxTable": "outbox_messages",
      "LeaseDuration": "00:02:00"
    }
  }
}

Message Lifecycle

Each message progresses through these states:

Pending → Processing → Processed
                    └→ Pending (retry scheduled)
                    └→ Failed  (max retries exhausted)
State Description
Pending Waiting to be picked up by the worker
Processing Claimed by a worker (locked via lease)
Processed Successfully handled/sent
Failed Permanently failed after exhausting all retries

Lease-Based Recovery

A LockedUntilUtc lease prevents stuck messages. If a worker crashes mid-processing, the lease expires (default: 2 minutes) and another worker reclaims the message.

Retry Backoff

Retry 1: +5s     (InitialRetryDelay)
Retry 2: +10s    (×2)
Retry 3: +20s    (×2)
Retry 4: +40s    (×2)
Retry 5: +80s    (×2, capped at MaxRetryDelay)

Graceful Shutdown

When the application is stopping:

  1. Workers stop accepting new batches
  2. If DrainOnShutdown is true (default), in-flight batches are allowed to complete
  3. If the drain exceeds DrainTimeout (default: 30s), remaining work is cancelled
  4. Unfinished messages remain in the database — their lease expires and they're reclaimed on next startup

Health Checks

Health checks are registered automatically:

Check Tags
inbox-worker liveness
outbox-worker liveness
inboxoutbox-sqlserver readiness
inboxoutbox-postgresql readiness
app.MapHealthChecks("/health/live", new() { Predicate = r => r.Tags.Contains("liveness") });
app.MapHealthChecks("/health/ready", new() { Predicate = r => r.Tags.Contains("readiness") });

Worker health checks report unhealthy if the worker is not running or its heartbeat is older than HealthStaleAfter.


Database Migrations

LowCodeHub.InboxOutbox does not create or migrate its own database schema during service registration. The schema scripts ship as embedded resources in the package, but the consuming application owns when and how they are applied.

Embedded resource prefixes:

Provider Embedded resource prefix
SQL Server LowCodeHub.InboxOutbox.Repositories.SqlServer.Scripts.
PostgreSQL LowCodeHub.InboxOutbox.Repositories.PostgreSql.Scripts.

The scripts are idempotent: they create the inbox and outbox message tables when absent and add missing worker lease/result columns to an existing schema. They also ensure the required polling and deduplication indexes exist.

Existing journaled scripts are not automatically re-executed when a package changes. To repair an older schema, explicitly execute the embedded resources again from an application-owned migration, or configure the directory as an always-execute directory.

With LowCodeHub.Migration.SqlServer

using LowCodeHub.InboxOutbox.Migrations;
using LowCodeHub.Migration.SqlServer.Extensions;

builder.Services.AddSqlServerMigrations(migrations =>
{
    migrations.AddTarget<IInboxOutboxScriptScanner>(o =>
    {
        o.ConnectionString = builder.Configuration.GetConnectionString("InboxOutbox")!;
        o.Directories = ["LowCodeHub.InboxOutbox.Repositories.SqlServer.Scripts."];
    });
});

await app.RunSqlServerMigrationAsync();

With LowCodeHub.Migration.PostgreSql

using LowCodeHub.InboxOutbox.Migrations;
using LowCodeHub.Migration.PostgreSql.Extensions;

builder.Services.AddPostgreSqlMigrations(migrations =>
{
    migrations.AddTarget<IInboxOutboxScriptScanner>(o =>
    {
        o.ConnectionString = builder.Configuration.GetConnectionString("InboxOutbox")!;
        o.Directories = ["LowCodeHub.InboxOutbox.Repositories.PostgreSql.Scripts."];
    });
});

await app.RunPostgreSqlMigrationAsync();

Do not scan the whole LowCodeHub.InboxOutbox assembly without a directory filter, because the package contains scripts for both providers.

With EF Core Migrations

EF Core will not discover these embedded scripts automatically. Create an application-owned EF migration and execute the embedded scripts from Up.

using LowCodeHub.InboxOutbox.Migrations;
using Microsoft.EntityFrameworkCore.Migrations;

public partial class AddLowCodeHubInboxOutboxSchema : Migration
{
    protected override void Up(MigrationBuilder migrationBuilder)
    {
        foreach (var resource in SqlInboxOutbox.SqlServerResources)
        {
            migrationBuilder.Sql(SqlInboxOutbox.ReadFromResource(resource));
        }
    }

    protected override void Down(MigrationBuilder migrationBuilder)
    {
        // Drop InboxOutbox objects here if your application's migration policy requires reversible migrations.
    }
}

For PostgreSQL, loop over SqlInboxOutbox.PostgreSqlResources instead.


Custom Repository Implementations

The library uses two repository interfaces. Implement them to use a different database:

public interface IInboxRepository
{
    Task<bool> TryAddAsync(InboxMessage message, CancellationToken ct);
    Task<IReadOnlyList<InboxMessage>> ClaimPendingAsync(int batchSize, DateTimeOffset utcNow, CancellationToken ct);
    Task MarkProcessedAsync(Guid id, DateTimeOffset processedAtUtc, CancellationToken ct);
    Task MarkFailedAsync(Guid id, string error, int retryCount, DateTimeOffset? nextAttemptAtUtc, CancellationToken ct);
}

public interface IOutboxRepository
{
    Task<bool> TryAddAsync(OutboxMessage message, CancellationToken ct);
    Task<IReadOnlyList<OutboxMessage>> ClaimPendingAsync(int batchSize, DateTimeOffset utcNow, CancellationToken ct);
    Task MarkProcessedAsync(Guid id, DateTimeOffset processedAtUtc, CancellationToken ct);
    Task MarkFailedAsync(Guid id, string error, int retryCount, DateTimeOffset? nextAttemptAtUtc, CancellationToken ct);
}

Register your implementations before calling AddInboxOutbox:

builder.Services.AddScoped<IInboxRepository, MyMongoInboxRepository>();
builder.Services.AddScoped<IOutboxRepository, MyMongoOutboxRepository>();
builder.Services.AddInboxOutbox(builder.Configuration);
// No need to call AddInboxOutboxSqlServer/AddInboxOutboxPostgreSql

The library uses TryAddScoped, so your registrations take precedence.


Requirements

  • .NET 10 or later
  • SQL Server (via Microsoft.Data.SqlClient) or PostgreSQL (via Npgsql) — or your own repository implementation

License

MIT © Ahmed Abuelnour

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

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
0.0.15 128 8/10/2026
0.0.14 103 8/6/2026
0.0.13 155 7/16/2026
0.0.12 114 7/9/2026
0.0.11 121 6/23/2026
0.0.10 124 6/21/2026
0.0.5 118 5/18/2026
0.0.4 115 5/13/2026
0.0.3 113 5/12/2026
0.0.2 128 4/23/2026
0.0.1 119 3/26/2026