OrionPatch
OrionPatch
Transactional outbox primitive for .NET. Enqueue inside SaveChanges, dispatch at-least-once through a pluggable sink.
What it does
OrionPatch is a transactional outbox primitive for .NET. You enqueue a message inside an EF Core SaveChanges call; it commits in the same transaction as your domain data; a background dispatcher hands it to a pluggable IOutboxSink at-least-once.
The current release is 0.4.2. The sections below describe the original v0.1.0 surface as historical record; capabilities added since then — the inbox / dedup table, the concrete RabbitMQ / Azure Service Bus / Kafka broker sinks, and the dead-letter store and archival APIs — are listed in the package table below and detailed in the CHANGELOG.
The core package is deliberately small and ships no broker itself. Concrete broker sinks for RabbitMQ, Azure Service Bus, and Kafka ship as separate opt-in sub-packages (OrionPatch.RabbitMQ, OrionPatch.AzureServiceBus, OrionPatch.Kafka); a NATS sink remains on the roadmap. The core also ships ChannelOutboxSink (in-process System.Threading.Channels, zero external dependency, useful for monoliths and tests).
At its core it owns one thing well: getting a message from "I just did a domain mutation" to "the sink received it — at least once per row, even if my process crashes between commit and send." (Delivery is at-least-once; sinks must be idempotent.) Inbox idempotency / dedup and the broker sinks build outward from that core in the sub-packages above.
How it works
A domain event is enqueued by application code, persisted by the EF Core interceptor inside the same transaction as your data, then handed to the sink asynchronously by a hosted dispatcher.
The diagram shows the at-least-once contract clearly: the outbox row and the domain rows commit together, but the sink call happens outside the transaction.
Why OrionPatch?
| Feature | OrionPatch | DIY interceptor | MassTransit | Wolverine |
|---|---|---|---|---|
| Transactional enqueue | Yes | Yes | Yes | Yes |
| At-least-once dispatch | Yes | Maybe | Yes | Yes |
| EF Core SaveChangesInterceptor | Yes | Yes | Optional | - |
| Multi-provider claim (SQL Server/Postgres/MySQL/SQLite) | Yes (native SKIP LOCKED on Postgres/MySQL) | Maybe | Optional | Yes |
| Pluggable sink (no broker bundled) | Yes | - | Bundled | Bundled |
| Built-in retry + dead-letter | Yes | Maybe | Yes | Yes |
| Dead-letter store (route exhausted rows out of the hot outbox) | Yes (v0.3) | Maybe | Yes | Yes |
| Outbox archival / retention (reap processed rows) | Yes (v0.3) | Maybe | Optional | Optional |
| OpenTelemetry | Yes | Maybe | Yes | Yes |
| In-process test sink | Yes | - | Yes | Yes |
| Saga / process manager | No (out of scope) | - | Yes | Yes |
| Standalone primitive (no framework) | Yes | Yes | No | No |
OrionPatch is a primitive, not a framework. If you want sagas, request/response, or a built-in mediator, reach for MassTransit or Wolverine. If you want transactional outbox without adopting a messaging framework, OrionPatch is the package.
30-second quick start
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Moongazing.OrionPatch.Abstractions;
using Moongazing.OrionPatch.DependencyInjection;
using Moongazing.OrionPatch.EntityFrameworkCore;
using Moongazing.OrionPatch.EntityFrameworkCore.DependencyInjection;
services.AddDbContext<AppDbContext>((sp, options) =>
{
options.UseNpgsql(connectionString);
options.UseOrionPatch(sp);
});
services.AddOrionPatch()
.UseEntityFrameworkCore<AppDbContext>()
.UseSink<MyKafkaSink>(); // or .UseChannelSink() for in-process
Apply the entity configuration in OnModelCreating:
protected override void OnModelCreating(ModelBuilder modelBuilder) =>
modelBuilder.ApplyOrionPatchConfiguration();
Enqueue from your service code:
public class OrderService
{
private readonly AppDbContext _db;
private readonly IOutbox _outbox;
public async Task ConfirmOrderAsync(Guid orderId, CancellationToken ct)
{
var order = await _db.Orders.FindAsync([orderId], ct);
order.Confirm();
_outbox.Enqueue(new OrderConfirmed(order.Id, order.TotalCents));
await _db.SaveChangesAsync(ct); // outbox row + order update commit together
}
}
Implement a sink:
public sealed class MyKafkaSink : IOutboxSink
{
public async Task SendAsync(OutboxEnvelope envelope, CancellationToken ct)
{
// External publish — keep this the last statement of the implementation so
// a failure after publish does not silently lose acknowledgement.
await _producer.ProduceAsync(envelope.MessageType, envelope.Payload, ct);
}
}
That's it. The dispatcher runs as a hosted service; messages flow from your transaction into the sink.
Packages
| Package | Description |
|---|---|
OrionPatch | Core: IOutbox, IOutboxSink, IOutboxStorage, dispatcher hosted service, telemetry, options. Includes ChannelOutboxSink. |
OrionPatch.EntityFrameworkCore | EF Core storage backend: OrionPatch_Outbox table, provider-aware claim (native FOR UPDATE SKIP LOCKED on Postgres/MySQL; SQLite + unknown providers use a portable compare-and-swap fallback), SaveChangesInterceptor for transactional enqueue, and an inbox / dedup table. |
OrionPatch.Testing | Test helpers: in-memory storage, deterministic dispatcher, capturing sink, test clock, fluent assertions. Zero EF Core dependency. |
OrionPatch.RabbitMQ | RabbitMQ broker sink (AddOrionPatchRabbitMqSink). |
OrionPatch.AzureServiceBus | Azure Service Bus broker sink (AddOrionPatchAzureServiceBusSink). |
OrionPatch.Kafka | Kafka broker sink (AddOrionPatchKafkaSink) plus a Kafka inbox (AddOrionPatchKafkaInbox). |
What the core does NOT do
- No broker bundled in the core — RabbitMQ, Azure Service Bus, and Kafka sinks ship as opt-in sub-packages; a NATS sink is still on the roadmap.
- No saga / process manager (that is OrionSaga territory).
- No distributed transactions across heterogeneous sinks.
- No push-based dispatch (PostgreSQL
LISTEN/NOTIFY, SQL Server Service Broker) — v0.3+ work.
At-least-once contract
OrionPatch guarantees at-least-once delivery. Duplicates occur in two known scenarios:
- The sink succeeds but the subsequent
CompleteAsyncwrite fails or the process crashes before it runs. The row stays Claimed, the lease expires, another dispatcher re-delivers. - The sink call exceeds
OrionPatchOptions.LeaseDuration(default 2 minutes). Another dispatcher may claim and re-deliver the row mid-flight.
Consumer sinks MUST be idempotent. Typical patterns: deduplicate at the destination on OutboxEnvelope.Id, or use upserts. Keep the external publish the last statement of the sink implementation so a failure after publish does not silently lose acknowledgement.
Dead-letter store and archival (v0.3.0)
v0.3.0 adds two outbox maintenance capabilities. Both are SPIs on the storage backend, not separate services: a storage type opts in by implementing the interface, and the dispatcher uses it when present.
Dead-letter store (IDeadLetterStore)
When a row exhausts OrionPatchOptions.MaxAttempts, the dispatcher prefers to route it OUT of the hot outbox into a dedicated dead-letter store instead of flipping it to DeadLettered in place. Routing removes the source row from the active outbox (so it can never be reclaimed or retried) and appends a DeadLetteredMessage snapshot carrying the final failure context: payload, headers, correlation id, enqueue time, total attempt count, final error, and the dead-letter instant.
Routing is idempotent on the row id. A redelivered or crash-replayed terminal-path call for an already-routed row is a no-op, so a message lands in the store exactly once and produces no duplicate metrics or alerts. Storage that does not implement IDeadLetterStore keeps the prior in-place status flip, so this is backward compatible.
This is distinct from the v0.2.18 IDeadLetterSink observer. The sink is a fire-and-forget triage notification (Slack, PagerDuty); the store is the durable destination the message is moved into.
// InMemoryOutboxStorage (and any storage that implements IDeadLetterStore) is detected by the
// dispatcher automatically. To inspect or replay abandoned messages, query the store directly:
if (storage is IDeadLetterStore deadLetterStore)
{
IReadOnlyList<DeadLetteredMessage> abandoned =
await deadLetterStore.GetDeadLetteredAsync(ct);
foreach (var message in abandoned)
{
// message.Id, message.MessageType, message.Payload, message.FinalError,
// message.AttemptCount, message.DeadLetteredAtUtc ...
}
}
Archival (IOutboxArchivalStore)
Successfully dispatched (Processed) rows accumulate in the hot outbox; an ever-growing table degrades claim-query planning and storage cost. ArchiveProcessedAsync reaps Processed rows whose ProcessedAtUtc is at or before nowUtc - retention out of the active outbox and returns the count moved. Pending, Claimed, and DeadLettered rows are never touched, and a processed row still inside the retention window is never touched. The reap is idempotent and incremental, so it is safe to call on a schedule.
OrionPatchOptions.ArchiveRetention (default 7 days, validated non-negative) expresses the retention horizon. ArchiveProcessedAsync is operator-invoked maintenance: OrionPatch does not start a background reaper, so call it from your own scheduled job (a hosted BackgroundService, Quartz.NET, Hangfire, or a cron-triggered endpoint).
// Run from a scheduled maintenance job, e.g. nightly.
if (storage is IOutboxArchivalStore archivalStore)
{
int reaped = await archivalStore.ArchiveProcessedAsync(
options.Value.ArchiveRetention, DateTime.UtcNow, ct);
}
The bundled InMemoryOutboxStorage supports an archive mode (default; reaped rows are observable via GetArchivedAsync) and a purge mode (new InMemoryOutboxStorage(purgeOnArchive: true); reaped rows are discarded).
Telemetry
ActivitySourceandMeternamedMoongazing.OrionPatch.- Spans:
OrionPatch.Dispatchper envelope, tagged withorionpatch.message.typeandorionpatch.attempt. - Counters:
orionpatch.outbox.enqueued,.dispatched,.failed,.deadlettered,.attempts. - Histogram:
orionpatch.outbox.dispatch.duration(milliseconds).
Wire them up with the standard OpenTelemetry .NET helpers.
Benchmarks
See benchmarks.md for the scenarios we plan to measure and the current status of the BenchmarkDotNet harness. A formal bench/Moongazing.OrionPatch.Bench project is on the v0.2 roadmap; the dispatcher has only been profiled informally during development so far.
Roadmap
The current release is 0.4.2, which shipped the outbox dead-letter store (IDeadLetterStore) and outbox archival (IOutboxArchivalStore) described above. See the CHANGELOG for the full per-version history.
12-month forward plan in ROADMAP.md. The next milestones:
- Push-based dispatch (LISTEN/NOTIFY, Service Broker).
- Operator dashboard, schema-evolution helpers.
- v1.0.0: API freeze, LTS window.
If something on the list matters to you, open an issue with the roadmap label.
More from the Orion family
OrionPatch is one of several standalone .NET libraries:
- OrionGuard — input validation, guard clauses, DDD primitives.
- OrionAudit — EF Core audit trail with JSON Patch diffs and time-travel reconstruction.
- OrionKey — source-generated strongly-typed IDs.
- OrionLock — distributed lock primitive with auto-renewing leases.
Each ships separately; none depends on another at runtime.
See it in a real app
Moongazing.OrionShowcase is a production-shaped banking sample integrating all six Orion packages end-to-end. OrionPatch outbox interceptor captures Account/Customer domain events into the same transaction as SaveChanges. DomainEventOutboxAdapter walks AggregateRoot.DomainEvents and enqueues via reflection on IOutbox.Enqueue
- src/Moongazing.OrionShowcase.Infrastructure/Outbox/DomainEventOutboxAdapter.cs
- src/Moongazing.OrionShowcase.Infrastructure/DependencyInjection/InfrastructureServiceCollectionExtensions.cs
Contributing
Issues and pull requests welcome. Please read CONTRIBUTING.md and the Code of Conduct before opening one.
License
MIT. See LICENSE.txt.
Packages
| Package | Version | Downloads |
|---|---|---|
| OrionPatch | 0.4.2 | 7,108 |
| OrionPatch.EntityFrameworkCore | 0.4.2 | 4,617 |
| OrionPatch.Testing | 0.4.2 | 4,478 |
| OrionPatch.Kafka | 0.4.2 | 2,358 |
| OrionPatch.RabbitMQ | 0.4.2 | 1,835 |
| OrionPatch.AzureServiceBus | 0.4.2 | 1,788 |