Skip to main content

OrionPatch

OrionPatch

OrionPatch

Transactional outbox primitive for .NET. Enqueue inside SaveChanges, dispatch at-least-once through a pluggable sink.

NuGet Downloads License Target


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?​

FeatureOrionPatchDIY interceptorMassTransitWolverine
Transactional enqueueYesYesYesYes
At-least-once dispatchYesMaybeYesYes
EF Core SaveChangesInterceptorYesYesOptional-
Multi-provider claim (SQL Server/Postgres/MySQL/SQLite)Yes (native SKIP LOCKED on Postgres/MySQL)MaybeOptionalYes
Pluggable sink (no broker bundled)Yes-BundledBundled
Built-in retry + dead-letterYesMaybeYesYes
Dead-letter store (route exhausted rows out of the hot outbox)Yes (v0.3)MaybeYesYes
Outbox archival / retention (reap processed rows)Yes (v0.3)MaybeOptionalOptional
OpenTelemetryYesMaybeYesYes
In-process test sinkYes-YesYes
Saga / process managerNo (out of scope)-YesYes
Standalone primitive (no framework)YesYesNoNo

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​

PackageDescription
OrionPatchCore: IOutbox, IOutboxSink, IOutboxStorage, dispatcher hosted service, telemetry, options. Includes ChannelOutboxSink.
OrionPatch.EntityFrameworkCoreEF 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.TestingTest helpers: in-memory storage, deterministic dispatcher, capturing sink, test clock, fluent assertions. Zero EF Core dependency.
OrionPatch.RabbitMQRabbitMQ broker sink (AddOrionPatchRabbitMqSink).
OrionPatch.AzureServiceBusAzure Service Bus broker sink (AddOrionPatchAzureServiceBusSink).
OrionPatch.KafkaKafka 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:

  1. The sink succeeds but the subsequent CompleteAsync write fails or the process crashes before it runs. The row stays Claimed, the lease expires, another dispatcher re-delivers.
  2. 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​

  • ActivitySource and Meter named Moongazing.OrionPatch.
  • Spans: OrionPatch.Dispatch per envelope, tagged with orionpatch.message.type and orionpatch.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. Concrete usage:

Contributing​

Issues and pull requests welcome. Please read CONTRIBUTING.md and the Code of Conduct before opening one.

License​

MIT. See LICENSE.txt.

Packages​

PackageVersionDownloads
OrionPatch0.4.27,108
OrionPatch.EntityFrameworkCore0.4.24,617
OrionPatch.Testing0.4.24,478
OrionPatch.Kafka0.4.22,358
OrionPatch.RabbitMQ0.4.21,835
OrionPatch.AzureServiceBus0.4.21,788