Stratara.Outbox.RabbitMQ 4.1.0

Prefix Reserved
There is a newer version of this package available.
See the version list below for details.
dotnet add package Stratara.Outbox.RabbitMQ --version 4.1.0
                    
NuGet\Install-Package Stratara.Outbox.RabbitMQ -Version 4.1.0
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="Stratara.Outbox.RabbitMQ" Version="4.1.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Stratara.Outbox.RabbitMQ" Version="4.1.0" />
                    
Directory.Packages.props
<PackageReference Include="Stratara.Outbox.RabbitMQ" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add Stratara.Outbox.RabbitMQ --version 4.1.0
                    
#r "nuget: Stratara.Outbox.RabbitMQ, 4.1.0"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package Stratara.Outbox.RabbitMQ@4.1.0
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=Stratara.Outbox.RabbitMQ&version=4.1.0
                    
Install as a Cake Addin
#tool nuget:?package=Stratara.Outbox.RabbitMQ&version=4.1.0
                    
Install as a Cake Tool

Stratara.Outbox.RabbitMQ

Derived. The behaviour described here is specified under openspec/specs/. Those specifications are the source; this page explains and illustrates them.

License: MIT.

Outbox-pattern command + event dispatch for the Stratara event-sourced stack with a RabbitMQ message-bus implementation. Contains the write-side dispatchers, the outbox-retry worker, the read-side mediator command worker, the RabbitMQ bus, and the ProjectionReplayState that coordinates dispatch skip during projection replay — shared over Redis where a connection is registered, held in process otherwise. Azure Service Bus ships separately as Stratara.Outbox.AzureServiceBus.

What's in the box

Folder Contents
Outbox/ OutboxOptions, CommandOutboxDispatcher (write-side ICommand fan-out via IMessageBus, falls back to outbox table on bus failure), EventBundleOutboxDispatcher (same for EventBundle), OutboxWorker (hosted service that retries unpublished outbox rows on a polling interval), NullOutboxLock + RedisOutboxLock (IOutboxLock implementations — default no-op for single-instance deployments, Redis-leased distributed lock for multi-replica setups)
Messaging/ RabbitMqBusIMessageBus over RabbitMQ. A worker subscription is a quorum queue <subscription>.v2 with <subscription>.dead-letter beside it; a message a handler cannot take is redelivered under MessageRetryOptions and then dead-lettered. Azure Service Bus ships as the sibling Stratara.Outbox.AzureServiceBus package.
Mediator/ MediatorCommandWorker (hosted service that subscribes to the command topic and dispatches into the in-process IMediator)
Projections/ ProjectionReplayState (concrete IProjectionReplayState, Redis-backed where a connection is registered, in process otherwise; dispatchers skip publishing while replay is active), ProjectionReplayOptions (how long the replay marking survives without renewal)
DependencyInjection/ AddOutboxDispatcher(), AddOutboxWorker(IConfiguration), AddRedisOutboxLock() (opt-in distributed lock), AddProjectionReplayState(), AddMediatorWorker(), AddMessaging()
Diagnostics/Extensions/ LoggerOutboxExtensions, LoggerMessagingExtensions (source-generated logger surfaces)

Quick start

// In your API host:
builder.AddMessaging();                          // IMessageBus + MessagingOptions binding
builder.Services
    .AddOutboxDispatcher()                       // CommandOutboxDispatcher + EventBundleOutboxDispatcher + ProjectionReplayState
    .AddOutboxWorker(builder.Configuration);     // OutboxWorker hosted service (only if this host owns retries)

// In your command worker:
builder.Services
    .AddMediatorWorker();                        // MediatorCommandWorker hosted service

The dispatchers consult IProjectionReplayState.IsReplayActive before each publish and skip dispatch (writing to the outbox table only) while a replay is in progress. They also record the outbox.published counter for what the broker accepted — the worker counts nothing, because only the dispatcher knows what went out.

The replay marking and its progress counters are held on the lease configured by ProjectionReplayOptions.LeaseSeconds (default 300), which the replay renews every time it reports progress. A replay whose host is killed stops renewing it and the marking lapses on its own, instead of suppressing publication indefinitely. AddProjectionReplayState() registers the options with their defaults; override with services.Configure<ProjectionReplayOptions>(...). Set it longer than the slowest stretch between two progress reports — too short lets the marking lapse while the replay is still running.

A drain cycle takes one batch of each kind and ends. Rows the broker would not accept stay in the table for the next interval; a cycle never re-reads what it just failed to publish, so an unreachable broker or a suppressed drain cannot turn a cycle into a loop over the same rows.

Multi-instance outbox workers

AddOutboxWorker registers NullOutboxLock as the default IOutboxLock — a no-op that preserves the single-instance assumption. For multi-replica deployments call AddRedisOutboxLock() afterwards; it replaces the no-op with a Redis-leased lock (SET stratara:outbox:lock NX EX) so only one replica drains at a time:

builder.AddCaching();                              // registers IConnectionMultiplexer
builder.Services
    .AddOutboxDispatcher()
    .AddOutboxWorker(builder.Configuration)
    .AddRedisOutboxLock();                         // multi-replica safe

The lease defaults to 60 s (OutboxOptions.LockLeaseSeconds). Tune it so it exceeds the worst-case drain duration; otherwise the lock can expire mid-cycle and a peer may start a concurrent drain. Outbox semantics are still at-least-once, so a duplicate publish is recoverable provided handlers stay idempotent.

Dependencies

  • Stratara.Abstractions — for ICommand, IEvent, IMessageBus, ICommandOutboxDispatcher, IEventBundleOutboxDispatcher, IProjectionReplayState, IMessagingIdentifier, IWriteUnitOfWork (used at runtime via the outbox repository).
  • Stratara.Contracts — for EventBundle + CommandEnvelope messages.
  • Stratara.MediatorMediatorCommandWorker dispatches into the in-process IMediator.
  • Stratara.Sessions — dispatcher hydrates CommandEnvelope from the current session context.
  • Stratara.Shared — for messaging primitives, resilience pipeline names, mapping helpers, and the diagnostics base.
  • RabbitMQ.Client, StackExchange.Redis (only used when a connection is registered: shared replay state, optional outbox lock).
  • Microsoft.Extensions.Hosting.Abstractions + Microsoft.Extensions.Options.ConfigurationExtensions — for hosted services + options binding.

The outbox dispatcher persists rows through IWriteUnitOfWork.CreateOutboxRepository — that interface lives in Stratara.Abstractions, but the concrete implementation comes from Stratara.EventSourcing.EntityFrameworkCore. Reference that package alongside this one to get a working stack.

Product 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. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

NuGet packages (2)

Showing the top 2 NuGet packages that depend on Stratara.Outbox.RabbitMQ:

Package Downloads
Stratara.Infrastructure

Infrastructure glue for the Stratara framework — authorization decorators, configuration providers, and DI composition helpers that wire Mediator, Outbox, Identity, and EF Core into a hosted app.

Stratara.EventSourcing.WorkerDefaults

Worker-host wiring composites for the Stratara event-sourced stack. IHostApplicationBuilder extensions (AddBackendServices, AddCommandWorkerServices, AddHeavyCommandWorkerServices, AddEventProjectionWorkerServices, AddEventStreamHashWorkerServices, AddSagaWorkerServices, AddOutboxWorkerServices) bundle the per-concern DI calls so each worker host opts in with one line.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
4.1.1 66 9/16/2026
4.1.0 76 9/16/2026
4.0.4 511 9/14/2026
4.0.3 177 9/3/2026
4.0.2 607 9/3/2026
4.0.1 210 9/2/2026
4.0.0 165 8/31/2026
4.0.0-preview.1 74 8/31/2026
3.4.0 159 8/28/2026
3.3.0 125 8/25/2026
3.2.3 140 8/22/2026
3.2.2 145 8/14/2026
3.2.1 151 8/2/2026
3.2.0 142 7/18/2026
3.1.7 155 7/1/2026
3.1.6 570 6/22/2026
3.1.5 152 6/22/2026
3.1.4 153 6/15/2026
3.1.3 154 6/10/2026
3.1.2 171 6/5/2026
Loading failed

The Orleans execution model ships as two new packages, `Stratara.Orleans` and
`Stratara.Orleans.EntityFrameworkCore`: one writer per aggregate across a cluster, accepted commands that
survive a crash, and projections and sagas that never miss a committed fact. Every host on the Entity
Framework store generates one migration when it upgrades, whether or not it adopts the execution model.
Telemetry type tags now carry the simple type name on every instrument, so a dashboard that filters the
projection, saga or conflict series on the qualified name needs its filter updated. Eight fixes land on
the retry, conflict, redaction and encryption work of 4.0.4, with one new warning to know about.

### Added

- **The Orleans execution model: `Stratara.Orleans` and `Stratara.Orleans.EntityFrameworkCore`.**
 Commands, projections, sagas, durable timers and singleton work run as virtual actors on an Orleans
 10.3 cluster. One aggregate has one writer across the deployment; a command the execution model's
 dispatcher accepts is recorded before the call returns and resumed after a crash, a bounded number of
 times, through the same mediator pipeline as on the bus; projections and sagas read the store in
 commit order from a checkpoint, so a crash costs latency and never a committed fact; singleton work
 runs once per cluster without a lock; a failing entry stops its partition and is retried instead of
 being dropped. Handlers, projections and sagas are unchanged, each role is adopted with one
 registration after its composite, and both models can run side by side during a rollout. Every
 setting is validated at start. See *Choose an Execution Model*, *Migrate to the Orleans Execution
 Model* and *Operate the Orleans Execution Model* on the documentation site.
- **Schema additions — every host on the Entity Framework store generates a migration.** The framework's
 write context declares `event_stream_entry.partition_position`, the `partition_position` table, on
 PostgreSQL `event_stream_entry.commit_transaction_id`, and the `outbox_entry` columns `aggregate_id`,
 `heavy`, `attempt_count`, `last_handed_over_at`, `kept_at` and `last_failure`; the read context declares
 `projection_checkpoint`. They are part of the model whether or not a host adopts the execution model, so
 every host on `Stratara.EventSourcing.EntityFrameworkCore` generates and applies a migration after
 upgrading — without it, appends and outbox writes fail on the missing columns. A store that never runs
 the execution model carries them unfilled; on PostgreSQL an append also reads back the transaction id.
- **Composites without the bus-fed worker.** `AddEventProjectionServices()` and `AddSagaServices()`
 register the projection and saga runtimes without their bus workers; `AddProjectionHandling` and
 `AddSagaHandling` do the same on `IServiceCollection`.
- **Every instrument name is a published constant.** `ApplicationDiagnostics.Metrics` carries a
 `…Name` constant beside each instrument — `EventsAppendedName`, `ProjectionEventsProcessedName`,
 `SagasInFlightName` and the others — so a query or listener references the name instead of a literal.
 `ApplicationDiagnostics.MetricTags.TypeNameValue` gives the value a type tag carries for a type name.
- `LogEvents.Messaging.WorkerQueueDeclaredWithOtherArguments` (`108_112`), the warning a RabbitMQ
 worker logs when it uses an existing queue declared with other arguments.

### Changed

- **`event.type` and `aggregate.type` carry the simple type name on every instrument.** On
 `projection.events.processed`, `saga.events.processed` and `event_source.append.conflicts` the tag
 value changes from the assembly-qualified name (`Shop.Orders.OrderPlaced, Shop.Orders`) to the simple
 name (`OrderPlaced`), the form `event_source.events.appended` already used, so write and read series
 join on the same value. A dashboard or alert that filters those three series on the qualified name
 needs its filter updated; no instrument or tag name changes.
- **`AddAuthorizingCommandOutboxDispatcher()` decorates whichever command dispatcher is registered**, not
 only the RabbitMQ one, and composes with a dispatcher registered after it.

### Fixed

- **RabbitMQ: changing the retry bounds no longer stops a worker from subscribing.** A worker queue
 carries `x-delivery-limit` from the bounds it was first declared with, and RabbitMQ refuses a
 redeclaration with a different value — so a deployment that changed `MessageRetry` failed every
 subscription with `PRECONDITION_FAILED`. An existing queue is now used as it is, and the bus logs
 `108_112` naming it. From RabbitMQ 4.3 the old limit does not cut the new bounds short; before 4.3
 a raised bound needs the drained queue deleted once.
- **Azure Service Bus: a subscription whose `MaxDeliveryCount` is below the bounds no longer hides
 its dead-lettering.** Where the host can read the subscription, the bounds for it are lowered to
 fit under the broker's limit, so the framework makes the move and it reaches `108_110` and
 `messaging.dead_lettered`. The warning `108_111` now says so. A failure to read the subscription
 of any kind other than cancellation ends the check instead of the subscription.
- **A unique violation is a concurrency conflict whatever exception type carries it.** The event
 source consulted the registered `IStoreConflictDetector`s only for an Entity Framework
 `DbUpdateException`, so a unit of work that surfaced the provider's exception unwrapped got a
 persistence failure instead of a `ConcurrencyException`. Every detector now sees the exception
 as the save threw it, as its contract says.
- **Proxy credentials and response cookies are redacted from traces in the form the semantic
 conventions name them.** The HTTP client and ASP.NET Core enrichment callbacks replaced
 `http.request.header.proxy_authorization` and `http.response.header.set_cookie`, but the
 OpenTelemetry semantic conventions keep the dash — `proxy-authorization`, `set-cookie` — so a host
 that captured headers under those names exported the values. Both forms are redacted now.
- **The concurrency-conflict policy retries a conflict the event source reports.** An append that
 lost the race surfaces from `IEventSource.SaveChangesAsync` as `ConcurrencyException`, and
 `ResilienceNames.ConcurrencyConflict` retried only `ConcurrencyConflictException` — so an
 `IResilientRequest` whose handler appends events ran once and failed on the first conflict. The
 policy now retries both, and still nothing else.
- **A revoked encrypted field of a value type reads as its default instead of failing the object.**
 Deserialization wrote an unreadable field back as JSON `null`, which a `decimal`, `int` or other
 non-nullable value type cannot hold — so erasing a subject's key made every record with such a
 field throw instead of degrading. The field now reads as the type's default, and the object's
 other fields are recovered.
- **`InMemoryKeyStore` keeps a scope usable after its current version is revoked.** Revoking the
 current key left the scope pointing at a key that no longer existed, so the next
 `GetOrCreateCurrentKeyAsync` threw. It now falls back to the highest remaining version, or creates
 a new one, as `EnvelopeFileKeyStore` does.
- **`IAggregationService.AggregateAsync` no longer documents `fromVersion` as a start version.** The
 parameter has never been honoured: a rebuild starts from the stream's beginning, or from the latest
 snapshot at or below `toVersion`. Its documentation now says so; the signature is unchanged.