LowCodeHub.InboxOutbox
0.0.1
See the version list below for details.
dotnet add package LowCodeHub.InboxOutbox --version 0.0.1
NuGet\Install-Package LowCodeHub.InboxOutbox -Version 0.0.1
<PackageReference Include="LowCodeHub.InboxOutbox" Version="0.0.1" />
<PackageVersion Include="LowCodeHub.InboxOutbox" Version="0.0.1" />
<PackageReference Include="LowCodeHub.InboxOutbox" />
paket add LowCodeHub.InboxOutbox --version 0.0.1
#r "nuget: LowCodeHub.InboxOutbox, 0.0.1"
#:package LowCodeHub.InboxOutbox@0.0.1
#addin nuget:?package=LowCodeHub.InboxOutbox&version=0.0.1
#tool nuget:?package=LowCodeHub.InboxOutbox&version=0.0.1
LowCodeHub.InboxOutbox
LowCodeHub.InboxOutbox is a DI-first .NET library that implements reliable Inbox/Outbox processing for ASP.NET Core applications.
Why Inbox/Outbox Pattern
Distributed systems are usually at-least-once for message delivery. That means duplicates, ack failures, and partial failures are normal.
Inbox pattern purpose: safely consume incoming messages with idempotency.
- Prevents duplicate business execution when the same external message is delivered again.
- In this library, deduplication is based on
(Source, ExternalMessageId).
Outbox pattern purpose: safely publish outgoing messages from your local transaction boundary.
- Prevents losing events when DB commit succeeds but broker/API publish fails.
- Outgoing messages are persisted first, then dispatched by worker with retry.
Together, they make message processing resilient across crashes, restarts, deployment rollouts, and network failures.
What You Get
- Inbox worker that continuously polls and handles new messages.
- Outbox worker that continuously publishes pending messages.
- Repository abstractions for storage-layer independence.
- Handler/sender abstractions for business logic and transport integration.
- Retry/backoff behavior with configurable limits.
Install
dotnet add package LowCodeHub.InboxOutbox
Register Services
using LowCodeHub.InboxOutbox.Extensions;
builder.Services.AddInboxOutbox(options =>
{
options.Inbox.PollInterval = TimeSpan.FromSeconds(1);
options.Outbox.PollInterval = TimeSpan.FromSeconds(1);
options.Inbox.DrainOnShutdown = true;
options.Inbox.DrainTimeout = TimeSpan.FromSeconds(30);
options.Inbox.HealthStaleAfter = TimeSpan.FromMinutes(2);
});
or through configuration:
builder.Services.AddInboxOutbox(builder.Configuration);
Section name: LowCodeHub:InboxOutbox.
Worker Operations (Shutdown + Health)
Each worker now supports:
- graceful drain on shutdown (
DrainOnShutdown,DrainTimeout) - liveness heartbeat checks (
HealthStaleAfter)
Drain Concept (Clean Shutdown)
When Kubernetes terminates a pod, the app usually gets a short grace period before force kill.
DrainOnShutdown = true:- worker stops claiming new messages
- worker finishes the currently running batch (best effort)
- reduces interrupted processing during deployments/restarts
DrainTimeout:- maximum wait time for that in-flight batch to finish
- if exceeded, worker cancels remaining in-flight work and exits
Why this adds value:
- fewer unnecessary retries after rollout
- less duplicate/replayed work noise
- cleaner operational behavior under frequent pod restarts
Recommended baseline:
- keep
DrainOnShutdown = true - set
DrainTimeoutto match your platform shutdown grace period (for example 20-30 seconds)
Built-in health checks registered by AddInboxOutbox:
inbox-worker(liveness)outbox-worker(liveness)
Provider registrations add readiness checks:
AddInboxOutboxSqlServer⇒inboxoutbox-sqlserverAddInboxOutboxPostgreSql⇒inboxoutbox-postgresql
Expose them from your app:
app.MapHealthChecks("/health/live", new HealthCheckOptions
{
Predicate = r => r.Tags.Contains("liveness")
});
app.MapHealthChecks("/health/ready", new HealthCheckOptions
{
Predicate = r => r.Tags.Contains("readiness")
});
Implement Required Contracts
You must register:
IInboxRepositoryIOutboxRepositoryIInboxMessageHandler(one handler per message type)IOutboxMessageSender
The workers resolve repositories and handlers through scoped DI on each polling cycle.
Add Messages (Enqueue)
You can enqueue messages directly through repositories:
public sealed class OrderApplicationService
{
private readonly IInboxRepository _inboxRepository;
private readonly IOutboxRepository _outboxRepository;
public OrderApplicationService(
IInboxRepository inboxRepository,
IOutboxRepository outboxRepository)
{
_inboxRepository = inboxRepository;
_outboxRepository = outboxRepository;
}
public async Task EnqueueAsync(CancellationToken ct = default)
{
var inboxInserted = 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);
var outboxInserted = await _outboxRepository.TryAddAsync(new OutboxMessage
{
Destination = "orders.exchange",
MessageType = "order.created",
Payload = "{\"orderId\":\"123\"}",
CorrelationId = Guid.NewGuid().ToString("N")
}, ct);
// false means duplicate (Source, ExternalMessageId)
_ = inboxInserted;
_ = outboxInserted;
}
}
Workers will pick these records automatically.
Recommended Runtime Flow (Important)
Inbox: what RabbitMQ consumer should do
Your RabbitMQ consumer should stay thin and only do ingestion:
- Read broker message metadata and payload.
- Map it to inbox fields:
Source,ExternalMessageId,MessageType,Payload,CorrelationId. - Call
IInboxRepository.TryAddAsync(...). ackonly after DB call succeeds.true: new row inserted →ack.false: duplicate(Source, ExternalMessageId)→ack(already ingested before).
- Do not execute heavy business logic in consumer; worker/handler does that.
If DB insert fails, do not ack so broker can redeliver.
Inbox: what worker does
- Worker claims pending messages.
- Dispatches to
IInboxMessageHandler. - On success → marks processed.
- On failure → marks failed and schedules retry (until max retries).
Outbox: what application code should do
Inside your business transaction:
- Save domain state changes.
- Add outbox row via
IOutboxRepository.TryAddAsync(...). - Commit transaction.
Then outbox worker publishes later:
- Worker claims pending outbox rows.
- Calls
IOutboxMessageSender. - On success → marks processed.
- On failure → retry with backoff.
This is the key separation:
- Consumer/API path = persist intent only.
- Worker path = execute side effects with retries.
SQL Server Repositories (Built-in)
Built-in repositories are available for SQL Server:
SqlServerInboxRepositorySqlServerOutboxRepository
Register them with:
using LowCodeHub.InboxOutbox.Extensions;
builder.Services
.AddInboxOutbox(builder.Configuration)
.AddInboxOutboxSqlServer(builder.Configuration);
or code-based:
builder.Services
.AddInboxOutbox()
.AddInboxOutboxSqlServer(options =>
{
options.ConnectionString = builder.Configuration.GetConnectionString("DefaultConnection")!;
options.LeaseDuration = TimeSpan.FromMinutes(2);
options.InboxSchema = "dbo";
options.InboxTable = "InboxMessages";
options.OutboxSchema = "dbo";
options.OutboxTable = "OutboxMessages";
});
Configuration section:
{
"LowCodeHub": {
"InboxOutbox": {
"SqlServer": {
"ConnectionString": "...",
"InboxSchema": "dbo",
"InboxTable": "InboxMessages",
"OutboxSchema": "dbo",
"OutboxTable": "OutboxMessages",
"LeaseDuration": "00:02:00"
}
}
}
}
Schema script:
PostgreSQL Repositories (Built-in)
Built-in repositories are available for PostgreSQL:
PostgreSqlInboxRepositoryPostgreSqlOutboxRepository
Register them with:
using LowCodeHub.InboxOutbox.Extensions;
builder.Services
.AddInboxOutbox(builder.Configuration)
.AddInboxOutboxPostgreSql(builder.Configuration);
or code-based:
builder.Services
.AddInboxOutbox()
.AddInboxOutboxPostgreSql(options =>
{
options.ConnectionString = builder.Configuration.GetConnectionString("PostgreSql")!;
options.LeaseDuration = TimeSpan.FromMinutes(2);
options.InboxSchema = "public";
options.InboxTable = "inbox_messages";
options.OutboxSchema = "public";
options.OutboxTable = "outbox_messages";
});
Configuration section:
{
"LowCodeHub": {
"InboxOutbox": {
"PostgreSql": {
"ConnectionString": "...",
"InboxSchema": "public",
"InboxTable": "inbox_messages",
"OutboxSchema": "public",
"OutboxTable": "outbox_messages",
"LeaseDuration": "00:02:00"
}
}
}
}
Schema script:
Minimal Example
public sealed class OrderCreatedInboxHandler : IInboxMessageHandler
{
public bool CanHandle(string messageType)
=> messageType.Equals("order.created", StringComparison.OrdinalIgnoreCase);
public Task HandleAsync(InboxMessage message, CancellationToken cancellationToken = default)
{
// process incoming payload
return Task.CompletedTask;
}
}
public sealed class RabbitOutboxSender : IOutboxMessageSender
{
public Task SendAsync(OutboxMessage message, CancellationToken cancellationToken = default)
{
// publish message to broker / API
return Task.CompletedTask;
}
}
Processing Notes
- Repository
ClaimPendingAsyncshould atomically claim rows to avoid duplicate processing across instances. MaxRetriescontrols retry budget per message; once exhausted,nextAttemptAtUtcis set tonull.- Workers are safe to run in multi-instance deployments when repository claim logic uses locking/concurrency guarantees.
- Inbox deduplication is enforced using unique
(Source, ExternalMessageId). - SQL Server implementation uses atomic claim with
UPDLOCK+READPAST, and reclaims expired leases for crashed pods. - PostgreSQL implementation uses atomic claim with
FOR UPDATE SKIP LOCKED, and reclaims expired leases for crashed pods.
SQL Starters
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 CONSTRAINT [DF_InboxMessages_RetryCount] DEFAULT (0),
[MaxRetries] INT NOT NULL CONSTRAINT [DF_InboxMessages_MaxRetries] DEFAULT (5),
[State] INT NOT NULL CONSTRAINT [DF_InboxMessages_State] DEFAULT (0),
[LockedUntilUtc] DATETIMEOFFSET NULL,
[ProcessedAtUtc] DATETIMEOFFSET NULL,
[LastError] NVARCHAR(MAX) NULL
);
GO
CREATE UNIQUE INDEX [UX_InboxMessages_Source_ExternalMessageId]
ON [dbo].[InboxMessages] ([Source], [ExternalMessageId]);
GO
CREATE INDEX [IX_InboxMessages_Polling]
ON [dbo].[InboxMessages] ([State], [NextAttemptAtUtc], [CreatedAtUtc]);
GO
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 CONSTRAINT [DF_OutboxMessages_RetryCount] DEFAULT (0),
[MaxRetries] INT NOT NULL CONSTRAINT [DF_OutboxMessages_MaxRetries] DEFAULT (5),
[State] INT NOT NULL CONSTRAINT [DF_OutboxMessages_State] DEFAULT (0),
[LockedUntilUtc] DATETIMEOFFSET NULL,
[ProcessedAtUtc] DATETIMEOFFSET NULL,
[LastError] NVARCHAR(MAX) NULL
);
GO
CREATE INDEX [IX_OutboxMessages_Polling]
ON [dbo].[OutboxMessages] ([State], [NextAttemptAtUtc], [CreatedAtUtc]);
GO
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);
License
MIT
| 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 (>= 6.1.4)
- Microsoft.Extensions.Configuration.Binder (>= 10.0.3)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 10.0.3)
- Microsoft.Extensions.Diagnostics.HealthChecks (>= 10.0.3)
- Microsoft.Extensions.Hosting.Abstractions (>= 10.0.3)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.3)
- Microsoft.Extensions.Options (>= 10.0.3)
- Microsoft.Extensions.Options.ConfigurationExtensions (>= 10.0.3)
- Npgsql (>= 10.0.1)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.