LowCodeHub.InboxOutbox
0.0.2
See the version list below for details.
dotnet add package LowCodeHub.InboxOutbox --version 0.0.2
NuGet\Install-Package LowCodeHub.InboxOutbox -Version 0.0.2
<PackageReference Include="LowCodeHub.InboxOutbox" Version="0.0.2" />
<PackageVersion Include="LowCodeHub.InboxOutbox" Version="0.0.2" />
<PackageReference Include="LowCodeHub.InboxOutbox" />
paket add LowCodeHub.InboxOutbox --version 0.0.2
#r "nuget: LowCodeHub.InboxOutbox, 0.0.2"
#:package LowCodeHub.InboxOutbox@0.0.2
#addin nuget:?package=LowCodeHub.InboxOutbox&version=0.0.2
#tool nuget:?package=LowCodeHub.InboxOutbox&version=0.0.2
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
{
"LowCodeHub": {
"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);
{
"LowCodeHub": {
"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);
{
"LowCodeHub": {
"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
SQL Server Schema
CREATE TABLE [dbo].[InboxMessages]
(
[Id] UNIQUEIDENTIFIER NOT NULL PRIMARY KEY,
[Source] NVARCHAR(200) NOT NULL,
[ExternalMessageId] NVARCHAR(300) NOT NULL,
[MessageType] NVARCHAR(200) NOT NULL,
[Payload] NVARCHAR(MAX) NOT NULL,
[Headers] NVARCHAR(MAX) NULL,
[CorrelationId] NVARCHAR(200) NULL,
[CreatedAtUtc] DATETIMEOFFSET NOT NULL,
[NextAttemptAtUtc] DATETIMEOFFSET NULL,
[RetryCount] INT NOT NULL DEFAULT (0),
[MaxRetries] INT NOT NULL DEFAULT (5),
[State] INT NOT NULL DEFAULT (0),
[LockedUntilUtc] DATETIMEOFFSET NULL,
[ProcessedAtUtc] DATETIMEOFFSET NULL,
[LastError] NVARCHAR(MAX) NULL
);
CREATE UNIQUE INDEX [UX_InboxMessages_Source_ExternalMessageId]
ON [dbo].[InboxMessages] ([Source], [ExternalMessageId]);
CREATE INDEX [IX_InboxMessages_Polling]
ON [dbo].[InboxMessages] ([State], [NextAttemptAtUtc], [CreatedAtUtc]);
CREATE TABLE [dbo].[OutboxMessages]
(
[Id] UNIQUEIDENTIFIER NOT NULL PRIMARY KEY,
[Destination] NVARCHAR(300) NOT NULL,
[MessageType] NVARCHAR(200) NOT NULL,
[Payload] NVARCHAR(MAX) NOT NULL,
[Headers] NVARCHAR(MAX) NULL,
[CorrelationId] NVARCHAR(200) NULL,
[CreatedAtUtc] DATETIMEOFFSET NOT NULL,
[NextAttemptAtUtc] DATETIMEOFFSET NULL,
[RetryCount] INT NOT NULL DEFAULT (0),
[MaxRetries] INT NOT NULL DEFAULT (5),
[State] INT NOT NULL DEFAULT (0),
[LockedUntilUtc] DATETIMEOFFSET NULL,
[ProcessedAtUtc] DATETIMEOFFSET NULL,
[LastError] NVARCHAR(MAX) NULL
);
CREATE INDEX [IX_OutboxMessages_Polling]
ON [dbo].[OutboxMessages] ([State], [NextAttemptAtUtc], [CreatedAtUtc]);
PostgreSQL Schema
CREATE TABLE IF NOT EXISTS public.inbox_messages
(
id UUID PRIMARY KEY,
source VARCHAR(200) NOT NULL,
external_message_id VARCHAR(300) NOT NULL,
message_type VARCHAR(200) NOT NULL,
payload TEXT NOT NULL,
headers JSONB NULL,
correlation_id VARCHAR(200) NULL,
created_at_utc TIMESTAMPTZ NOT NULL,
next_attempt_at_utc TIMESTAMPTZ NULL,
retry_count INTEGER NOT NULL DEFAULT 0,
max_retries INTEGER NOT NULL DEFAULT 5,
state INTEGER NOT NULL DEFAULT 0,
locked_until_utc TIMESTAMPTZ NULL,
processed_at_utc TIMESTAMPTZ NULL,
last_error TEXT NULL
);
CREATE UNIQUE INDEX IF NOT EXISTS ux_inbox_messages_source_external_message_id
ON public.inbox_messages (source, external_message_id);
CREATE INDEX IF NOT EXISTS ix_inbox_messages_polling
ON public.inbox_messages (state, next_attempt_at_utc, created_at_utc);
CREATE TABLE IF NOT EXISTS public.outbox_messages
(
id UUID PRIMARY KEY,
destination VARCHAR(300) NOT NULL,
message_type VARCHAR(200) NOT NULL,
payload TEXT NOT NULL,
headers JSONB NULL,
correlation_id VARCHAR(200) NULL,
created_at_utc TIMESTAMPTZ NOT NULL,
next_attempt_at_utc TIMESTAMPTZ NULL,
retry_count INTEGER NOT NULL DEFAULT 0,
max_retries INTEGER NOT NULL DEFAULT 5,
state INTEGER NOT NULL DEFAULT 0,
locked_until_utc TIMESTAMPTZ NULL,
processed_at_utc TIMESTAMPTZ NULL,
last_error TEXT NULL
);
CREATE INDEX IF NOT EXISTS ix_outbox_messages_polling
ON public.outbox_messages (state, next_attempt_at_utc, created_at_utc);
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.0)
- Microsoft.Extensions.Configuration.Binder (>= 10.0.6)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 10.0.6)
- Microsoft.Extensions.Diagnostics.HealthChecks (>= 10.0.6)
- Microsoft.Extensions.Hosting.Abstractions (>= 10.0.6)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.6)
- Microsoft.Extensions.Options (>= 10.0.6)
- Microsoft.Extensions.Options.ConfigurationExtensions (>= 10.0.6)
- Npgsql (>= 10.0.2)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.