LowCodeHub.InboxOutbox
0.0.15
dotnet add package LowCodeHub.InboxOutbox --version 0.0.15
NuGet\Install-Package LowCodeHub.InboxOutbox -Version 0.0.15
<PackageReference Include="LowCodeHub.InboxOutbox" Version="0.0.15" />
<PackageVersion Include="LowCodeHub.InboxOutbox" Version="0.0.15" />
<PackageReference Include="LowCodeHub.InboxOutbox" />
paket add LowCodeHub.InboxOutbox --version 0.0.15
#r "nuget: LowCodeHub.InboxOutbox, 0.0.15"
#:package LowCodeHub.InboxOutbox@0.0.15
#addin nuget:?package=LowCodeHub.InboxOutbox&version=0.0.15
#tool nuget:?package=LowCodeHub.InboxOutbox&version=0.0.15
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.
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
- Configuration
- Implementing Contracts
- Enqueuing Messages
- Recommended Runtime Flow
- Storage Backends
- Message Lifecycle
- Graceful Shutdown
- Health Checks
- Database Migrations
- Custom Repository Implementations
- Requirements
- License
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) orFOR 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);
Recommended Runtime Flow
Inbox: RabbitMQ consumer stays thin
- Read broker message metadata and payload.
- Map to inbox fields:
Source,ExternalMessageId,MessageType,Payload,CorrelationId. - Call
TryAddAsync(...). ackafter DB call succeeds (true= new,false= duplicate — both ack).- Do not execute business logic in the consumer — the worker/handler does that.
Outbox: persist inside business transaction
- Save domain state changes.
- Add outbox row via
TryAddAsync(...). - Commit transaction.
- 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:
- Workers stop accepting new batches
- If
DrainOnShutdownistrue(default), in-flight batches are allowed to complete - If the drain exceeds
DrainTimeout(default: 30s), remaining work is cancelled - 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 (viaNpgsql) — or your own repository implementation
License
MIT © Ahmed Abuelnour
| Product | Versions 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. |
-
net10.0
- Microsoft.Data.SqlClient (>= 7.0.2)
- Microsoft.Extensions.Diagnostics.HealthChecks (>= 10.0.9)
- Microsoft.Extensions.Options.ConfigurationExtensions (>= 10.0.9)
- Npgsql (>= 10.0.3)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.