From eb0d1ea3e80090080fe4b6a40b471117af4f6e85 Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 12:54:39 -0600 Subject: [PATCH 01/10] fix(outbox): register entity-event tracker and router authoritatively (#15) Under modular / multi-datastore composition the transactional outbox silently stopped persisting durable events (0 rows, no error). WithEventTracking runs at the end of every WithPersistence call and TryAdds the in-memory IEntityEventTracker; when a non-outbox module registered first it pinned the in-memory tracker, so the outbox producer's TryAdd of OutboxEntityEventTracker no-opped and every durable event was dispatched in-process and dropped. AddOutboxProducer now Remove-then-Adds IEntityEventTracker and IEventRouter so the outbox wins regardless of registration order. OutboxEntityEventTracker is a strict superset of the in-memory tracker, so it is always correct for it to win. Adds a modular-composition regression test. --- .../ReproOutboxModularDiTests.cs | 152 ++++++++++++++++++ .../OutboxPersistenceBuilderExtensions.cs | 34 ++-- .../OutboxHostRouterResolutionTests.cs | 16 +- ...2026-07-23-outbox-di-registration-3.2.1.md | 93 +++++++++++ 4 files changed, 279 insertions(+), 16 deletions(-) create mode 100644 Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ReproOutboxModularDiTests.cs create mode 100644 docs/superpowers/plans/2026-07-23-outbox-di-registration-3.2.1.md diff --git a/Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ReproOutboxModularDiTests.cs b/Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ReproOutboxModularDiTests.cs new file mode 100644 index 00000000..87f4d56c --- /dev/null +++ b/Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ReproOutboxModularDiTests.cs @@ -0,0 +1,152 @@ +using FluentAssertions; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using RCommon; +using RCommon.Entities; +using RCommon.EventHandling; +using RCommon.Persistence.Crud; +using RCommon.Persistence.EFCore; +using RCommon.Persistence.EFCore.Outbox; +using RCommon.Persistence.Outbox; +using RCommon.Persistence.Transactions; +using Xunit; + +namespace Examples.EventHandling.Outbox.Tests; + +/// +/// Regression for 3.2.0 defect #15 (silent outbox data loss under modular composition). +/// +/// The customer composes their app from independent modules, each of which calls +/// WithPersistence<EFCorePersistenceBuilder>(...) for its own bounded context. Only one +/// module configures the transactional outbox. The others do not. +/// +/// WithPersistence runs WithEventTracking at the END of every call, which +/// TryAddScoped<IEntityEventTracker, InMemoryEntityEventTracker>. In 3.2.0 the outbox +/// producer also registered its tracker with TryAdd — so if any non-outbox module registered +/// first, the in-memory tracker was pinned and the outbox tracker registration silently no-opped. +/// Durable events were then dispatched in-process and NEVER written to __OutboxMessages, with +/// no exception and no warning. This test reproduces that composition (non-outbox datastore first, +/// outbox datastore second) and asserts the durable event actually lands in the outbox. +/// +/// Runs on the EF Core InMemory provider (fast lane — no containers). InMemory has no real +/// transactions, so this proves routing/persistence wiring, not transactional atomicity (that is +/// covered by the Postgres integration tests). +/// +public class ReproOutboxModularDiTests +{ + // ---- Domain event published durably to the "Billing" outbox ---- + public sealed class InvoiceRaisedEvent : IDomainEvent + { + public InvoiceRaisedEvent(Guid invoiceId, decimal amount) + { + InvoiceId = invoiceId; + Amount = amount; + } + + public Guid EventId { get; } = Guid.NewGuid(); + public DateTimeOffset OccurredOn { get; } = DateTimeOffset.UtcNow; + public Guid InvoiceId { get; } + public decimal Amount { get; } + } + + public sealed class Invoice : AggregateRoot + { + public Invoice() : base(Guid.NewGuid()) { } + public string CustomerName { get; set; } = string.Empty; + public decimal Amount { get; set; } + + public void Raise() => AddDomainEvent(new InvoiceRaisedEvent(Id, Amount)); + } + + // ---- "Orders" module datastore — NO outbox. Registered FIRST. ---- + public sealed class OrdersDbContext : RCommonDbContext + { + public OrdersDbContext(DbContextOptions options) : base(options) { } + protected override void OnModelCreating(ModelBuilder modelBuilder) => base.OnModelCreating(modelBuilder); + } + + // ---- "Billing" module datastore — HAS the outbox. Registered SECOND. ---- + public sealed class BillingDbContext : RCommonDbContext + { + public BillingDbContext(DbContextOptions options) : base(options) { } + public DbSet Invoices => Set(); + protected override void OnModelCreating(ModelBuilder modelBuilder) => modelBuilder.AddOutboxMessages(); + } + + private static ServiceProvider BuildProvider(string ordersDb, string billingDb) + { + var services = new ServiceCollection(); + services.AddLogging(); + + services.AddRCommon() + .WithSimpleGuidGenerator() + .WithUnitOfWork(uow => { }) + .WithEventHandling(events => + events.Publish().UseOutbox("Billing")) + // MODULE 1 (no outbox) — registered FIRST. Its WithEventTracking pins the in-memory + // IEntityEventTracker, which is what defeated the outbox registration in 3.2.0. + .WithPersistence(ef => + { + ef.AddDbContext("Orders", o => o.UseInMemoryDatabase(ordersDb)); + ef.SetDefaultDataStore(ds => ds.DefaultDataStoreName = "Orders"); + }) + // MODULE 2 (outbox) — registered SECOND, as a separate WithPersistence call. + .WithPersistence(ef => + { + ef.AddDbContext("Billing", o => o.UseInMemoryDatabase(billingDb)); + ef.AddOutbox(dataStoreName: "Billing"); + }); + + return services.BuildServiceProvider(validateScopes: true); + } + + [Fact] + public void OutboxTracker_WinsRegistration_UnderModularComposition() + { + using var provider = BuildProvider(Guid.NewGuid().ToString(), Guid.NewGuid().ToString()); + + using var scope = provider.CreateScope(); + var tracker = scope.ServiceProvider.GetRequiredService(); + + tracker.Should().BeOfType( + "the transactional outbox must decorate the entity-event tracker even when a non-outbox " + + "module's WithPersistence registered the in-memory tracker first (defect #15)"); + } + + [Fact] + public async Task DurableEvent_IsPersisted_UnderModularComposition() + { + using var provider = BuildProvider(Guid.NewGuid().ToString(), Guid.NewGuid().ToString()); + + using (var scope = provider.CreateScope()) + { + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + } + + using (var scope = provider.CreateScope()) + { + var sp = scope.ServiceProvider; + var uowFactory = sp.GetRequiredService(); + var invoices = sp.GetRequiredService>(); + invoices.DataStoreName = "Billing"; + + var invoice = new Invoice { CustomerName = "Ada Lovelace", Amount = 249.99m }; + invoice.Raise(); + + using var uow = uowFactory.Create(); + await invoices.AddAsync(invoice); + await uow.CommitAsync(); + } + + using (var scope = provider.CreateScope()) + { + var billing = scope.ServiceProvider.GetRequiredService(); + var billingOutbox = await billing.Set().AsNoTracking().CountAsync(); + + billingOutbox.Should().Be(1, + "the durable event must be written to the Billing outbox — under modular composition " + + "3.2.0 silently dropped it (0 rows, no error): defect #15"); + } + } +} diff --git a/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs b/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs index 96915056..69c9fed8 100644 --- a/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs +++ b/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs @@ -55,18 +55,34 @@ public static IPersistenceBuilder AddOutboxProducer( { builder.AddOutboxCore(configure, dataStoreName); - // Outbox event router (scoped, concrete) + IEventRouter forwarder. + // Concrete collaborators the OutboxEntityEventTracker composes directly. These are safe to TryAdd: + // they are the outbox's own types, so nothing else registers them and first-registration-wins is + // fine. The InMemoryTransactionalEventRouter is the transient dispatcher used for the Phase-2 FIFO + // drain, independent of whatever IEventRouter resolves to in the host. builder.Services.TryAddScoped(); - builder.Services.TryAddScoped(sp => sp.GetRequiredService()); - - // In-process transactional router as its OWN concrete scoped type (the transient dispatcher the - // OutboxEntityEventTracker composes directly for the Phase-2 FIFO drain, independent of whatever - // IEventRouter resolves to in the outbox host). builder.Services.TryAddScoped(); - - // Entity event tracker decorator (scoped — replaces InMemoryEntityEventTracker). builder.Services.TryAddScoped(); - builder.Services.TryAddScoped(); + + // Authoritative outbox routing. These MUST win regardless of the order in which WithPersistence / + // WithEventHandling / AddRCommon ran, so they are registered with Remove-then-Add rather than TryAdd. + // + // Why TryAdd is wrong here (3.2.0 defect #15 — silent outbox data loss): + // * WithEventTracking runs at the END of EVERY WithPersistence call and TryAdds the in-memory + // IEntityEventTracker. Under modular / multi-datastore composition a non-outbox WithPersistence + // can run first, pinning the in-memory tracker; a later AddOutbox TryAdd then no-ops and every + // durable event is silently dispatched in-process and never written to the outbox. + // * RCommonBuilder's constructor registers IEventRouter -> InMemoryTransactionalEventRouter with an + // unconditional AddScoped, so an outbox IEventRouter forwarder registered with TryAdd could never + // win in ANY configuration. + // OutboxEntityEventTracker is a strict superset of the in-memory tracker (durable events -> outbox, + // transient events -> in-process), so it is always correct for it to win. Remove-then-Add is + // idempotent across repeated AddOutbox calls (it nets exactly one registration), and a subsequent + // non-outbox WithPersistence TryAdd will no-op because the outbox registration is already present. + builder.Services.RemoveAll(); + builder.Services.AddScoped(sp => sp.GetRequiredService()); + + builder.Services.RemoveAll(); + builder.Services.AddScoped(); // Startup diagnostic: warn if a later registration silently overrode the outbox IEventRouter. // Producer-only: it inspects the producer's IEventRouter -> OutboxEventRouter binding, which a diff --git a/Tests/RCommon.Persistence.Tests/OutboxHostRouterResolutionTests.cs b/Tests/RCommon.Persistence.Tests/OutboxHostRouterResolutionTests.cs index 0819faf4..f67bfdf2 100644 --- a/Tests/RCommon.Persistence.Tests/OutboxHostRouterResolutionTests.cs +++ b/Tests/RCommon.Persistence.Tests/OutboxHostRouterResolutionTests.cs @@ -104,20 +104,22 @@ public void OutboxEntityEventTracker_ResolvesWithAllFourConstructorDependencies( } [Fact] - public void IEventRouter_IsNotTheSameInstanceAsTheConcreteInProcessRouter() + public void IEventRouter_ResolvesToTheOutboxForwarder_AndIsDistinctFromTheConcreteInProcessRouter() { - // Informative: the tracker's transient dispatcher (the concrete InMemoryTransactionalEventRouter it - // composes) must be a separate object from whatever IEventRouter resolves to in the outbox host. - // In this host IEventRouter actually resolves to InMemoryTransactionalEventRouter (core registers it - // FIRST via AddScoped, so AddOutbox's TryAddScoped -> OutboxEventRouter is a no-op). - // Even so, the concrete scoped InMemoryTransactionalEventRouter is a DISTINCT registration/instance - // from the IEventRouter binding, which is what the concrete-injection design guarantees. + // As of the 3.2.1 defect-#15 fix, AddOutbox registers IEventRouter authoritatively (Remove-then-Add + // a forwarder to OutboxEventRouter) so it WINS over the in-memory router the core ctor registers + // first — regardless of registration order. IEventRouter therefore resolves to the OutboxEventRouter. + // The tracker's OWN transient dispatcher (the concrete InMemoryTransactionalEventRouter it composes) + // remains a DISTINCT registration/instance from the IEventRouter binding, which is what the + // concrete-injection design guarantees. using var provider = BuildOutboxHost(); using var scope = provider.CreateScope(); var eventRouter = scope.ServiceProvider.GetRequiredService(); var concreteInProcess = scope.ServiceProvider.GetRequiredService(); + eventRouter.Should().BeOfType( + "AddOutbox must authoritatively bind IEventRouter to the outbox forwarder (defect #15)"); eventRouter.Should().NotBeSameAs(concreteInProcess, "the tracker's composed transient dispatcher must be independent of the host's IEventRouter binding"); } diff --git a/docs/superpowers/plans/2026-07-23-outbox-di-registration-3.2.1.md b/docs/superpowers/plans/2026-07-23-outbox-di-registration-3.2.1.md new file mode 100644 index 00000000..ddec2d2f --- /dev/null +++ b/docs/superpowers/plans/2026-07-23-outbox-di-registration-3.2.1.md @@ -0,0 +1,93 @@ +# Outbox DI Registration Hardening (3.2.1) Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:test-driven-development for every task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Eliminate the silent-data-loss defect where the transactional outbox stops persisting durable events under modular / multi-`WithPersistence` composition, harden the processor-only host registration, correct the 3.2.0 docs that describe unshipped features, and ship a worked modular multi-datastore example. + +**Architecture:** The outbox producer decorates two core services — `IEntityEventTracker` (→ `OutboxEntityEventTracker`) and `IEventRouter` (→ `OutboxEventRouter` forwarder). Today both are registered with `TryAdd*`, and `RCommonBuilder`'s constructor registers `IEventRouter` with an unconditional `AddScoped`. So (a) the outbox `IEventRouter` forwarder *always* no-ops, and (b) the outbox `IEntityEventTracker` no-ops whenever a prior non-outbox `WithPersistence` call already registered the in-memory tracker via `WithEventTracking`. The fix makes the outbox registration **authoritative** (remove-then-add) so it wins regardless of call ordering, since `OutboxEntityEventTracker` is a strict superset of the in-memory tracker (durable → outbox, transient → in-process). + +**Tech Stack:** .NET 10, `Microsoft.Extensions.DependencyInjection`, xUnit + FluentAssertions, EF Core InMemory provider (fast lane, no containers). + +**Release:** 3.2.1 (patch; no public API surface change — behavior-only correction + additive diagnostic). + +--- + +## Background — verified root causes (3.2.0 source) + +- **#15 (High, silent data loss).** `OutboxPersistenceBuilderExtensions.AddOutboxProducer` uses `TryAddScoped` (line 69). `PersistenceBuilderExtensions.WithEventTracking` runs at the end of **every** `WithPersistence` call and does `TryAddScoped` (line 87). If any non-outbox `WithPersistence` runs before the outbox one (modular composition, second datastore, separate provider), the in-memory tracker is pinned first and the outbox `TryAdd` no-ops → durable events are dispatched in-process and **never written to the outbox**, with no exception and no warning. Single-`WithPersistence` happy path works because the `AddOutbox` call inside the delegate runs before `WithEventTracking`. +- **#15 nuance.** The customer attributed the root cause to the `IEventRouter` no-op. That no-op is real (`RCommonBuilder` ctor `AddScoped` unconditionally, so `AddOutboxProducer`'s `TryAddScoped` forwarder always no-ops) but it is **not** what breaks persistence — `OutboxEntityEventTracker` composes the *concrete* `OutboxEventRouter`, not the `IEventRouter` alias. Both are fixed here for correctness/consistency, but the load-bearing fix is `IEntityEventTracker`. +- **#15 diagnostic gap.** `OutboxRoutingDiagnosticsHostedService` only inspects `IEventRouter`'s last descriptor. Because the core ctor pins `IEventRouter` to the in-memory impl unconditionally, this diagnostic warns on **every** outbox configuration including the working one (false positive) and never inspects the load-bearing `IEntityEventTracker`. Fix it to inspect the tracker. +- **#16 (processor-only host).** `AddOutboxProcessor` registers only the shared core + hosted poller. To be reproduced with a DI smoke test (validateOnBuild) before choosing the fix; the poller itself resolves `IOutboxStore`, `IOutboxSerializer`, `IEventProducer[]`, `EventSubscriptionManager`, `IOutboxDataStoreRegistry`, `IBackoffStrategy`, `IOptions`, `IOptions` — verify which is missing on a processor-only host and register it (or document the dual-call requirement). +- **#17 (docs over-claim).** The 3.2.0 changelog "Added" list includes an outbox metrics `Meter`, `IOutboxPayloadProtector`, and a deserialization allow-list (spec AC-18/19/20). Grep confirms these are **absent** from shipped `Src`. Correct changelog/migration/spec to mark them planned-not-shipped. + +--- + +## Task 1: #15 regression test — modular composition drops durable events + +**Files:** +- Test: `Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ReproOutboxModularDiTests.cs` (create) + +- [ ] **Step 1: Write the failing test.** Build a provider that mirrors modular composition: two separate `WithPersistence` calls (or one non-outbox datastore registered before the outbox datastore) so `WithEventTracking` pins the in-memory tracker before `AddOutbox`. Publish a durable event, commit through the unit of work, then assert the owning datastore's `__OutboxMessages` has exactly one row. Also assert `provider.GetRequiredService()` is `OutboxEntityEventTracker`. +- [ ] **Step 2: Run it, watch it fail** — rows = 0 and tracker = `InMemoryEntityEventTracker`. This is the reproduction. + +## Task 2: #15 fix — authoritative outbox registration + +**Files:** +- Modify: `Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs:58-69` + +- [ ] **Step 1:** In `AddOutboxProducer`, keep the concrete `TryAddScoped` registrations (`OutboxEventRouter`, `InMemoryTransactionalEventRouter`, `InMemoryEntityEventTracker`). Replace the two interface `TryAdd`s with authoritative remove-then-add so the outbox wins regardless of ordering: + ```csharp + builder.Services.RemoveAll(); + builder.Services.AddScoped(sp => sp.GetRequiredService()); + builder.Services.RemoveAll(); + builder.Services.AddScoped(); + ``` + Add a comment explaining why `TryAdd` is wrong here (defect #15). Idempotent across repeated `AddOutbox` calls (remove-then-add nets one). A later non-outbox `WithPersistence` `TryAdd` will no-op because the outbox registration is present. +- [ ] **Step 2:** Run Task 1 test → passes (1 row, tracker = `OutboxEntityEventTracker`). +- [ ] **Step 3:** Run the whole `Examples.EventHandling.Outbox.Tests` + `RCommon.Persistence` unit suites → all green (no regression to happy-path / two-datastore). +- [ ] **Step 4: Commit.** + +## Task 3: #15 fix — routing diagnostic checks the load-bearing tracker + +**Files:** +- Modify: `Src/RCommon.Persistence/Outbox/OutboxRoutingDiagnosticsHostedService.cs` +- Test: add a diagnostic test (a host where the tracker was overridden warns; a correct host does not) + +- [ ] **Step 1:** Write a failing test asserting the diagnostic warns iff the effective `IEntityEventTracker` is not `OutboxEntityEventTracker` (and does **not** false-positive on a correctly-wired outbox host). +- [ ] **Step 2:** Change the diagnostic to inspect the last `IEntityEventTracker` descriptor (implementation type / factory target) instead of `IEventRouter`. With the Task 2 fix in place a correct host never warns. +- [ ] **Step 3:** Run → green. **Commit.** + +## Task 4: #16 — processor-only host + +**Files:** +- Test: `Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ProcessorOnlyHostDiTests.cs` (create) +- Modify (pending repro): `Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs` (`AddOutboxProcessor`) + +- [ ] **Step 1:** Write a DI smoke test: `AddRCommon().WithPersistence(ef => { AddDbContext; SetDefaultDataStore; ef.AddOutboxProcessor(); })` with `BuildServiceProvider(validateScopes: true, validateOnBuild: true)`, then start the host / resolve `IHostedService`s and run one poll batch. Observe the actual failure. +- [ ] **Step 2:** Based on the observed failure, either (a) have `AddOutboxProcessor` register the shared routing/tracker it needs, or (b) if a pure poller is genuinely self-sufficient and the failure only arises when the processor host also does domain writes, apply the same authoritative tracker registration and document the topology. Prefer the registration fix over docs-only. +- [ ] **Step 3:** Run → green. **Commit.** + +## Task 5: #17 — correct the docs over-claim + +**Files:** +- Modify: `website/docs/api-reference/changelog.mdx` (3.2.0 entry) +- Modify: `website/docs/api-reference/migration-guide.mdx` +- Modify: `docs/specs/event-handling/event-handling.md` (AC-18/19/20 status) +- Modify: `website/versioned_docs/version-3.2.0/api-reference/{changelog,migration-guide}.mdx` (released snapshot) + +- [ ] **Step 1:** Remove metrics `Meter`, `IOutboxPayloadProtector`, and deserialization allow-list from the 3.2.0 "Added" list; move them to a "Planned / not yet shipped" note (or delete). Mark AC-18/19/20 in the spec as deferred. **Commit.** + +## Task 6: modular multi-datastore outbox example + +**Files:** +- Create: `Examples/EventHandling/Examples.EventHandling.Outbox.Modular/` (project) — module registration extension methods (one per bounded context), each calling `WithPersistence`/`WithEventHandling`; a composition root wiring 2–3 datastores with the native outbox; a `Program` that commits and prints outbox row counts. +- Create: matching `*.Tests` asserting durable rows persist across all datastores under the modular composition. +- Add both to the solution. + +- [ ] **Step 1:** Build the modular example illustrating exactly the customer's shape (multiple modules, multiple datastores, native `AddOutbox`). Test proves rows persist. **Commit.** + +## Task 7: green + finish + +- [ ] Run full solution build + the event-handling/persistence test suites (fast lane, `Category!=Integration`). +- [ ] Squash interim commits into one meaningful 3.2.1 commit; push branch; open PR. +- [ ] Do **not** tag/release until the user approves the PR. From 69cc516c5e0d2b5a20c89998fb8a090c9a9d3437 Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 12:58:11 -0600 Subject: [PATCH 02/10] fix(outbox): routing diagnostic checks the load-bearing entity-event tracker (#15) The startup diagnostic inspected IEventRouter, which is not what persists durable events (the tracker composes the concrete OutboxEventRouter directly). Because the core ctor pins IEventRouter to the in-memory router unconditionally, the old check false-positived on every working outbox host and missed the real failure. It now inspects IEntityEventTracker and warns only when the effective tracker is the in-memory one. --- .../OutboxRoutingDiagnosticsHostedService.cs | 27 +++++++++-------- ...boxRoutingDiagnosticsHostedServiceTests.cs | 30 +++++++++++-------- 2 files changed, 32 insertions(+), 25 deletions(-) diff --git a/Src/RCommon.Persistence/Outbox/OutboxRoutingDiagnosticsHostedService.cs b/Src/RCommon.Persistence/Outbox/OutboxRoutingDiagnosticsHostedService.cs index bb2dd8a8..f0fddb07 100644 --- a/Src/RCommon.Persistence/Outbox/OutboxRoutingDiagnosticsHostedService.cs +++ b/Src/RCommon.Persistence/Outbox/OutboxRoutingDiagnosticsHostedService.cs @@ -4,7 +4,7 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; -using RCommon.EventHandling.Producers; +using RCommon.Entities; namespace RCommon.Persistence.Outbox; @@ -26,19 +26,22 @@ public OutboxRoutingDiagnosticsHostedService(IServiceCollection services, ILogge public Task StartAsync(CancellationToken cancellationToken) { - // DI is last-registration-wins. AddOutbox binds IEventRouter to the OutboxEventRouter (via a - // forwarding factory, so its ImplementationType is null). If a later event-handling registration - // binds the in-memory router, that wins and outbox routing is silently defeated -- events fire - // in-process post-commit and are never persisted to the outbox. - var lastRouter = _services.LastOrDefault(d => d.ServiceType == typeof(IEventRouter)); - if (lastRouter?.ImplementationType == typeof(InMemoryTransactionalEventRouter)) + // DI is last-registration-wins. Durable-event persistence rides on IEntityEventTracker resolving to + // OutboxEntityEventTracker -- the tracker composes the concrete OutboxEventRouter directly, so the + // IEventRouter binding is NOT what persists events (checking it produced false positives, since the + // core ctor pins IEventRouter to the in-memory router unconditionally even on a working outbox host). + // AddOutbox now binds the tracker authoritatively (Remove-then-Add), so the only way it can be the + // in-memory tracker at startup is an explicit later override -- which silently defeats the outbox + // (durable events fire in-process post-commit and are never persisted). Surface that. + var lastTracker = _services.LastOrDefault(d => d.ServiceType == typeof(IEntityEventTracker)); + if (lastTracker?.ImplementationType == typeof(InMemoryEntityEventTracker)) { _loggerFactory?.CreateLogger().LogWarning( - "RCommon outbox is configured (AddOutbox) but the effective IEventRouter is the in-memory " + - "router ({Router}); a later registration overrode the outbox router. Domain events will be " + - "dispatched in-memory and NOT persisted to the outbox. Register AddOutbox after any " + - "event-handling configuration that binds IEventRouter, or re-assert the outbox router last.", - typeof(InMemoryTransactionalEventRouter).Name); + "RCommon outbox is configured (AddOutbox) but the effective IEntityEventTracker is the " + + "in-memory tracker ({Tracker}); a later registration overrode the outbox tracker. Durable " + + "domain events will be dispatched in-memory and NOT persisted to the outbox. Remove the " + + "override, or re-assert the outbox tracker last.", + typeof(InMemoryEntityEventTracker).Name); } return Task.CompletedTask; diff --git a/Tests/RCommon.Persistence.Tests/OutboxRoutingDiagnosticsHostedServiceTests.cs b/Tests/RCommon.Persistence.Tests/OutboxRoutingDiagnosticsHostedServiceTests.cs index f9dda8e5..69390fed 100644 --- a/Tests/RCommon.Persistence.Tests/OutboxRoutingDiagnosticsHostedServiceTests.cs +++ b/Tests/RCommon.Persistence.Tests/OutboxRoutingDiagnosticsHostedServiceTests.cs @@ -4,6 +4,7 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Moq; +using RCommon.Entities; using RCommon.EventHandling.Producers; using RCommon.Persistence.Outbox; using Xunit; @@ -11,11 +12,12 @@ namespace RCommon.Persistence.Tests; /// -/// Covers the outbox-routing footgun: because DI is last-registration-wins, an event-handling -/// registration that binds the in-memory AFTER -/// silently overrides the -/// outbox router. When that happens, domain events fire in-memory and are never persisted to the -/// outbox. The diagnostic surfaces the misconfiguration at startup instead of leaving it silent. +/// Covers the outbox-routing footgun on the LOAD-BEARING service. Durable-event persistence rides on +/// resolving to — NOT on the +/// binding (the tracker composes the concrete +/// directly). If a later registration binds the in-memory tracker after the outbox producer configured +/// its tracker, durable events fire in-process and are never persisted, silently. The diagnostic +/// surfaces that at startup instead of leaving it silent. /// public class OutboxRoutingDiagnosticsHostedServiceTests { @@ -40,13 +42,14 @@ private static void VerifyWarning(Mock logger, Times times) } [Fact] - public async Task StartAsync_Warns_When_Outbox_Router_Is_Clobbered_By_InMemory_Router() + public async Task StartAsync_Warns_When_Outbox_Tracker_Is_Clobbered_By_InMemory_Tracker() { var services = new ServiceCollection(); - // Outbox registers IEventRouter via a factory forwarding to OutboxEventRouter... - services.AddScoped(_ => throw new InvalidOperationException("should not resolve")); - // ...then a later event-handling registration binds the in-memory router last (the footgun). - services.AddScoped(); + // Outbox binds IEntityEventTracker -> OutboxEntityEventTracker... + services.AddScoped(); + // ...then a later registration binds the in-memory tracker last (the footgun): durable events + // would fire in-process and never be persisted. + services.AddScoped(); var (loggerFactory, logger) = CreateLogger(); var diagnostic = new OutboxRoutingDiagnosticsHostedService(services, loggerFactory.Object); @@ -57,11 +60,12 @@ public async Task StartAsync_Warns_When_Outbox_Router_Is_Clobbered_By_InMemory_R } [Fact] - public async Task StartAsync_Does_Not_Warn_When_Outbox_Router_Is_Intact() + public async Task StartAsync_Does_Not_Warn_When_Outbox_Tracker_Is_Intact() { var services = new ServiceCollection(); - // Intact: the last IEventRouter registration is the outbox factory (not the in-memory router). - services.AddScoped(_ => throw new InvalidOperationException("should not resolve")); + // Intact: the last IEntityEventTracker registration is the outbox tracker. + services.AddScoped(); + services.AddScoped(); var (loggerFactory, logger) = CreateLogger(); var diagnostic = new OutboxRoutingDiagnosticsHostedService(services, loggerFactory.Object); From ba4d0bb9ee88b1b30ec62ab45e9dd7634cf4d154 Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 13:03:48 -0600 Subject: [PATCH 03/10] fix(outbox): register outbox routing in the shared core so processor hosts persist (#16) A host that ran the outbox poller (AddOutboxProcessor) and also committed domain entities silently dropped durable events, because only AddOutboxProducer registered the outbox tracker/router. Moves the authoritative tracker/router registration (plus its concrete collaborators) into AddOutboxCore, which both the producer and processor call. This also removes any path where the outbox tracker could be bound without its concrete InMemoryEntityEventTracker collaborator, which surfaced as an 'unable to resolve InMemoryEntityEventTracker' DI failure. Producer-only diagnostics stay on the producer path. --- .../ProcessorOnlyHostDiTests.cs | 99 +++++++++++++++++++ .../OutboxPersistenceBuilderExtensions.cs | 71 +++++++------ 2 files changed, 138 insertions(+), 32 deletions(-) create mode 100644 Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ProcessorOnlyHostDiTests.cs diff --git a/Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ProcessorOnlyHostDiTests.cs b/Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ProcessorOnlyHostDiTests.cs new file mode 100644 index 00000000..0c45de43 --- /dev/null +++ b/Examples/EventHandling/Examples.EventHandling.Outbox.Tests/ProcessorOnlyHostDiTests.cs @@ -0,0 +1,99 @@ +using Examples.EventHandling.Outbox; +using FluentAssertions; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using RCommon; +using RCommon.Entities; +using RCommon.EventHandling; +using RCommon.Persistence.Crud; +using RCommon.Persistence.EFCore; +using RCommon.Persistence.EFCore.Outbox; +using RCommon.Persistence.Outbox; +using RCommon.Persistence.Transactions; +using Xunit; + +namespace Examples.EventHandling.Outbox.Tests; + +/// +/// Reproduces the third-party report #16: a host configured with AddOutboxProcessor (the poller +/// half of the producer/processor topology) rather than AddOutbox. Two shapes are exercised: +/// 1. Pure poller host, built with ValidateOnBuild — must construct without a DI resolution failure. +/// 2. Processor host that ALSO commits domain entities raising a durable event — must not silently +/// drop the event (the #15 failure mode, but reached via the processor-only registration path). +/// Runs on the EF Core InMemory provider (fast lane). +/// +public class ProcessorOnlyHostDiTests +{ + private static IServiceCollection BaseServices(string db, out IServiceCollection services) + { + services = new ServiceCollection(); + services.AddLogging(); + services.AddRCommon() + .WithSimpleGuidGenerator() + .WithUnitOfWork(uow => { }) + .WithPersistence(ef => + { + ef.AddDbContext("AppDb", o => o.UseInMemoryDatabase(db)); + ef.SetDefaultDataStore(ds => ds.DefaultDataStoreName = "AppDb"); + ef.AddOutboxProcessor(); // PROCESSOR ONLY (no producer) + }) + .WithEventHandling(eh => + { + eh.AddSubscriber(); + eh.Publish().UseOutbox("AppDb"); + }); + return services; + } + + [Fact] + public void ProcessorOnlyHost_BuildsWithValidateOnBuild_WithoutDiResolutionFailure() + { + BaseServices(Guid.NewGuid().ToString(), out var services); + + Action build = () => + { + using var provider = services.BuildServiceProvider(new ServiceProviderOptions + { + ValidateScopes = true, + ValidateOnBuild = true + }); + }; + + build.Should().NotThrow("a processor-only host must be constructible (#16)"); + } + + [Fact] + public async Task ProcessorOnlyHost_ThatCommitsDomainEntities_PersistsDurableEvent() + { + BaseServices(Guid.NewGuid().ToString(), out var services); + using var provider = services.BuildServiceProvider(validateScopes: true); + + using (var scope = provider.CreateScope()) + { + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + } + + using (var scope = provider.CreateScope()) + { + var sp = scope.ServiceProvider; + var orders = sp.GetRequiredService>(); + var uowFactory = sp.GetRequiredService(); + + var order = new Order { CustomerName = "Ada Lovelace", Total = 249.99m }; + order.Place(); + + using var uow = uowFactory.Create(); + await orders.AddAsync(order); + await uow.CommitAsync(); + } + + using (var scope = provider.CreateScope()) + { + var db = scope.ServiceProvider.GetRequiredService(); + var count = await db.Set().AsNoTracking().CountAsync(); + count.Should().Be(1, + "a host that runs the poller and also commits domain entities must persist the durable event"); + } + } +} diff --git a/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs b/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs index 69c9fed8..af087b71 100644 --- a/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs +++ b/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs @@ -53,40 +53,13 @@ public static IPersistenceBuilder AddOutboxProducer( string? dataStoreName = null) where TOutboxStore : class, IOutboxStore { + // Routing (tracker + routers) lives in AddOutboxCore so that EVERY outbox host — producer, + // processor, or the combined AddOutbox — persists durable events. See AddOutboxCore for why. builder.AddOutboxCore(configure, dataStoreName); - // Concrete collaborators the OutboxEntityEventTracker composes directly. These are safe to TryAdd: - // they are the outbox's own types, so nothing else registers them and first-registration-wins is - // fine. The InMemoryTransactionalEventRouter is the transient dispatcher used for the Phase-2 FIFO - // drain, independent of whatever IEventRouter resolves to in the host. - builder.Services.TryAddScoped(); - builder.Services.TryAddScoped(); - builder.Services.TryAddScoped(); - - // Authoritative outbox routing. These MUST win regardless of the order in which WithPersistence / - // WithEventHandling / AddRCommon ran, so they are registered with Remove-then-Add rather than TryAdd. - // - // Why TryAdd is wrong here (3.2.0 defect #15 — silent outbox data loss): - // * WithEventTracking runs at the END of EVERY WithPersistence call and TryAdds the in-memory - // IEntityEventTracker. Under modular / multi-datastore composition a non-outbox WithPersistence - // can run first, pinning the in-memory tracker; a later AddOutbox TryAdd then no-ops and every - // durable event is silently dispatched in-process and never written to the outbox. - // * RCommonBuilder's constructor registers IEventRouter -> InMemoryTransactionalEventRouter with an - // unconditional AddScoped, so an outbox IEventRouter forwarder registered with TryAdd could never - // win in ANY configuration. - // OutboxEntityEventTracker is a strict superset of the in-memory tracker (durable events -> outbox, - // transient events -> in-process), so it is always correct for it to win. Remove-then-Add is - // idempotent across repeated AddOutbox calls (it nets exactly one registration), and a subsequent - // non-outbox WithPersistence TryAdd will no-op because the outbox registration is already present. - builder.Services.RemoveAll(); - builder.Services.AddScoped(sp => sp.GetRequiredService()); - - builder.Services.RemoveAll(); - builder.Services.AddScoped(); - - // Startup diagnostic: warn if a later registration silently overrode the outbox IEventRouter. - // Producer-only: it inspects the producer's IEventRouter -> OutboxEventRouter binding, which a - // processor-only host never registers, so running it there would false-warn. + // Startup diagnostic: warn if a later registration silently overrode the outbox tracker. + // Producer-only: a processor-only host legitimately never commits domain entities, so warning + // there would be noise; the shared MN-3 durable-route validator (in the core) covers both hosts. builder.Services.TryAddEnumerable( ServiceDescriptor.Singleton( sp => new OutboxRoutingDiagnosticsHostedService( @@ -140,6 +113,40 @@ private static IPersistenceBuilder AddOutboxCore( // Outbox store (scoped — participates in per-request transaction). builder.Services.TryAddScoped(); + // Outbox routing — registered in the CORE (shared by producer, processor, and the combined + // AddOutbox) so ANY host that both runs the outbox and commits domain entities persists durable + // events. A pure poller host that never commits entities simply never constructs the tracker. + // + // Concrete collaborators the OutboxEntityEventTracker composes directly. Safe to TryAdd: they are + // the outbox's own types, so nothing else registers them. Registering all three together (not just + // the interface binding) is what prevents the "unable to resolve InMemoryEntityEventTracker" DI + // failure when the outbox tracker is bound (report #16). The InMemoryTransactionalEventRouter is the + // transient dispatcher for the Phase-2 FIFO drain, independent of whatever IEventRouter resolves to. + builder.Services.TryAddScoped(); + builder.Services.TryAddScoped(); + builder.Services.TryAddScoped(); + + // Authoritative outbox routing. These MUST win regardless of the order in which WithPersistence / + // WithEventHandling / AddRCommon ran, so they are registered with Remove-then-Add rather than TryAdd. + // + // Why TryAdd is wrong here (3.2.0 defect #15 — silent outbox data loss): + // * WithEventTracking runs at the END of EVERY WithPersistence call and TryAdds the in-memory + // IEntityEventTracker. Under modular / multi-datastore composition a non-outbox WithPersistence + // can run first, pinning the in-memory tracker; a later AddOutbox TryAdd then no-ops and every + // durable event is silently dispatched in-process and never written to the outbox. + // * RCommonBuilder's constructor registers IEventRouter -> InMemoryTransactionalEventRouter with an + // unconditional AddScoped, so an outbox IEventRouter forwarder registered with TryAdd could never + // win in ANY configuration. + // OutboxEntityEventTracker is a strict superset of the in-memory tracker (durable events -> outbox, + // transient events -> in-process), so it is always correct for it to win. Remove-then-Add is + // idempotent across repeated AddOutbox calls (it nets exactly one registration), and a subsequent + // non-outbox WithPersistence TryAdd will no-op because the outbox registration is already present. + builder.Services.RemoveAll(); + builder.Services.AddScoped(sp => sp.GetRequiredService()); + + builder.Services.RemoveAll(); + builder.Services.AddScoped(); + // Serializer (singleton, replaceable). builder.Services.TryAddSingleton(); From 93805dcba9f0806557108e4777c1b11244af21d6 Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 13:07:57 -0600 Subject: [PATCH 04/10] docs(outbox): mark metrics/payload-protector/allow-list as planned, not shipped (#17) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The 3.2.0 changelog listed outbox metrics (AC-18), IOutboxPayloadProtector (AC-19), and the deserialization allow-list (AC-20) under Added, but none shipped. Moves them to an explicit 'Planned — not shipped in 3.2.0' note in the changelog (current + versioned snapshot) with interim workarounds, and marks AC-18/19/20 DEFERRED in the spec (acceptance criteria, observability, security, and open questions) so the spec reflects the actual shipped state. --- docs/specs/event-handling/event-handling.md | 22 +++++++++---------- website/docs/api-reference/changelog.mdx | 11 +++++++--- .../version-3.2.0/api-reference/changelog.mdx | 11 +++++++--- 3 files changed, 27 insertions(+), 17 deletions(-) diff --git a/docs/specs/event-handling/event-handling.md b/docs/specs/event-handling/event-handling.md index 28775ce3..3982777c 100644 --- a/docs/specs/event-handling/event-handling.md +++ b/docs/specs/event-handling/event-handling.md @@ -34,9 +34,9 @@ This domain spec formalizes the 3.2.0 redesign. Full design rationale, diagrams, - **AC-15 (broker coordination proven):** An integration test on the Podman/Testcontainers harness asserts that, under recipe 2b, business state + broker-outbox rows commit atomically and a rollback leaves neither. Recipe 2b is "done" only when green. - **AC-16 (recipes proven):** Each of the five recipes ships as a runnable example with an end-to-end test asserting the documented wiring composes and produces the recipe's observable outcome. - **AC-17 (back-compat shims):** `IEntityEventTracker.AddEntity(entity)` overload preserved (defaults to default datastore); `AddSubscriber` retained as `[Obsolete]` alias forwarding to `Consume` on broker builders; `EFCoreOutboxStore`/subclass marked `[Obsolete]`. -- **AC-18 (opt-in metrics):** RCommon exposes a `RCommon.Outbox` `System.Diagnostics.Metrics.Meter` instrumenting per-datastore pending depth, oldest-unprocessed age, relay success/failure counts, dead-letter rate, and dispatch-queue depth. Opt-in (host registers the meter with its metrics pipeline); additive/non-breaking. -- **AC-19 (payload protection hook):** An `IOutboxPayloadProtector` seam wraps payload serialization with `Protect`/`Unprotect`; the default implementation is pass-through (plaintext). Applications may supply an encrypting implementation. Non-breaking. -- **AC-20 (deserialization allow-list):** The outbox serializer resolves only event types present in the registration set (types with a route/subscriber/producer). An unknown/unresolvable `EventType` on relay/consume is logged loud and dead-lettered — never deserialized to an arbitrary type. +- **AC-18 (opt-in metrics) — DEFERRED, not shipped in 3.2.0.** Planned: RCommon exposes a `RCommon.Outbox` `System.Diagnostics.Metrics.Meter` instrumenting per-datastore pending depth, oldest-unprocessed age, relay success/failure counts, dead-letter rate, and dispatch-queue depth. Opt-in (host registers the meter with its metrics pipeline); additive/non-breaking. **Current state:** not implemented; observe via the `__OutboxMessages` tables or poller logs. +- **AC-19 (payload protection hook) — DEFERRED, not shipped in 3.2.0.** Planned: an `IOutboxPayloadProtector` seam wrapping payload serialization with `Protect`/`Unprotect`; default pass-through (plaintext), applications may supply an encrypting implementation. Non-breaking. **Current state:** not implemented; payloads are plaintext within the application-database trust boundary. +- **AC-20 (deserialization allow-list) — DEFERRED, not shipped in 3.2.0.** Planned: the outbox serializer resolves only event types present in the registration set (types with a route/subscriber/producer); an unknown/unresolvable `EventType` on relay/consume is logged loud and dead-lettered — never deserialized to an arbitrary type. **Current state:** not implemented; the serializer resolves the type named in the row. - **AC-21 (producer/processor topology):** First-class `AddOutboxProducer` (store/router/tracker, no hosted poller) and `AddOutboxProcessor` (hosted poller) registration methods exist alongside `AddOutbox` (= producer + processor). Each is datastore-scoped (`OnDataStore(...)`), consolidating the multi-host topology into this release's registration rework. ### Must Not Do @@ -50,7 +50,7 @@ This domain spec formalizes the 3.2.0 redesign. Full design rationale, diagrams, ### Nice to Have -- None outstanding. Items previously parked here — the producer/processor topology split, payload protection, and first-class metrics — were pulled into scope during spec review (AC-18, AC-19, AC-21). +- The producer/processor topology split (AC-21) was pulled into scope during spec review and **shipped** in 3.2.0. Payload protection (AC-19) and first-class metrics (AC-18) were also accepted into scope but were **deferred and did not ship in 3.2.0** — they remain roadmap items (see Open Questions). ## Technical Constraints @@ -76,14 +76,14 @@ External dependencies: one relational database per registered datastore; optiona - **Warnings (fail-loud):** poller draining an event type with zero matching subscribers (once per type); outbox routing overridden by a later registration (startup diagnostic); missing outbox schema on a registered datastore (startup diagnostic); cycle-breaker generation limit exceeded; best-effort relay/dispatch failure (before retry); dead-lettering. - **Debug/Info:** dispatch counts per commit, poller poll cycles and claim counts, `ImmediateDispatch` skip on producer-only hosts. -- **Metrics (first-class, opt-in):** a `RCommon.Outbox` `Meter` (System.Diagnostics.Metrics) exposes per-datastore outbox pending depth and oldest-unprocessed age; relay success/failure counts; dead-letter rate; dispatch-queue depth and max cascade generation reached (AC-18). Host-agnostic — the application wires the meter into OpenTelemetry/Prometheus/etc. +- **Metrics (first-class, opt-in) — DEFERRED, not shipped in 3.2.0 (AC-18).** Planned: a `RCommon.Outbox` `Meter` (System.Diagnostics.Metrics) exposing per-datastore outbox pending depth and oldest-unprocessed age; relay success/failure counts; dead-letter rate; dispatch-queue depth and max cascade generation reached. Host-agnostic — the application wires the meter into OpenTelemetry/Prometheus/etc. Until it ships, observe via the `__OutboxMessages` tables or the poller's warnings/logs. - **Alerting:** outbox backlog age exceeding a threshold and dead-letter rate are the primary signals; thresholds are host-owned. ## Security - **Attack surface:** event payloads are serialized into the outbox (JSON) and to brokers, then deserialized on relay/consume. Type resolution goes through `IOutboxSerializer`; only known/registered event types should be deserialized. -- **Data protection:** outbox rows live in the application database, inside the same trust boundary and at-rest protection as business data. `TenantId` is recorded per row for multi-tenant isolation. Payloads are plaintext by default; applications needing field/payload protection supply an `IOutboxPayloadProtector` (AC-19, default pass-through). -- **Deserialization safety:** the serializer enforces an allow-list of registered event types; a tampered/unknown `EventType` is logged and dead-lettered rather than deserialized (AC-20). +- **Data protection:** outbox rows live in the application database, inside the same trust boundary and at-rest protection as business data. `TenantId` is recorded per row for multi-tenant isolation. Payloads are plaintext. The `IOutboxPayloadProtector` field/payload-protection seam (AC-19) is **DEFERRED — not shipped in 3.2.0**; applications needing payload protection must handle it above the event model until it ships. +- **Deserialization safety — DEFERRED, not shipped in 3.2.0 (AC-20).** Planned: the serializer enforces an allow-list of registered event types and a tampered/unknown `EventType` is logged and dead-lettered rather than deserialized. **Current state:** the serializer resolves the type named in the outbox row; deploy the outbox only within a trusted boundary until the allow-list ships. - **Auth/authz:** not applicable at the library level; consumers execute in their own DI scope. Inbox idempotency prevents duplicate side effects from replays. - **Compliance:** none imposed by the library; applications remain responsible for any PII/regulatory handling of event payloads. @@ -110,13 +110,13 @@ Breaking changes are softened with shims (AC-17). Deliverables: a migration guid ## Open Questions -None outstanding. All were resolved during spec review (2026-07-22): +Scope was resolved during spec review (2026-07-22), but three accepted items (AC-18/19/20) were **deferred and did not ship in 3.2.0** — reopened below as roadmap items: -- **OQ-1 (metrics) → resolved:** first-class, opt-in `RCommon.Outbox` `Meter` (AC-18). -- **OQ-2 (payload protection) → resolved:** optional `IOutboxPayloadProtector`, default pass-through (AC-19). +- **OQ-1 (metrics) → DEFERRED:** first-class, opt-in `RCommon.Outbox` `Meter` (AC-18) is designed but not implemented in 3.2.0. +- **OQ-2 (payload protection) → DEFERRED:** optional `IOutboxPayloadProtector` (AC-19) is designed but not implemented in 3.2.0. - **OQ-3 (perf) → resolved:** correctness-focused + sanity throughput assertions; no formal benchmark suite in 3.2.0. - **OQ-4 (cycle-breaker default) → resolved:** 16 generations, configurable (AC-4). -- **OQ-5 (deserialization allow-list) → resolved:** enforce allow-list of registered types; unknown ⇒ fail-loud/dead-letter (AC-20). +- **OQ-5 (deserialization allow-list) → DEFERRED:** the registered-type allow-list (AC-20) is designed but not implemented in 3.2.0. - **OQ-6 (topology API) → resolved:** fold `AddOutboxProducer`/`AddOutboxProcessor` into 3.2.0 (AC-21). - **OQ-7 (example names) → resolved:** accept the proposed names, including `Examples.EventHandling.Outbox.MultiDataStore`, `Examples.Messaging.MassTransit.NativeOutbox`, `Examples.Messaging.Wolverine.NativeOutbox`, `Examples.EventHandling.TransactionScript`, `Examples.EventHandling.NoUnitOfWork`. diff --git a/website/docs/api-reference/changelog.mdx b/website/docs/api-reference/changelog.mdx index b9d6ab3a..1cd80145 100644 --- a/website/docs/api-reference/changelog.mdx +++ b/website/docs/api-reference/changelog.mdx @@ -28,9 +28,6 @@ Event-handling and transactional-outbox redesign. This is a **breaking** release - **Datastore-aware transactional outbox (B4/U5).** Per-datastore `__OutboxMessages` tables; each row records its target producer(s), so one table can fan out to multiple transports. `OutboxProcessingService` claims and drains every registered outbox datastore. Registering an outbox auto-applies the `OutboxMessage` mapping to that datastore's `RCommonDbContext`, and a startup diagnostic fails loud if a registered outbox datastore's table is missing. - **Broker-native outbox wrapper (`UseBrokerOutbox`).** Configures MassTransit's EF Core outbox (`AddEntityFrameworkOutbox` + `UseBusOutbox`) against an RCommon datastore's `DbContext`, so business state and broker-outbox rows commit atomically inside the RCommon unit-of-work transaction. Proven by a Postgres/Testcontainers coordination test. *(Wolverine does not support recipe 2b — see Known Issues / use recipe 2a.)* - **First-class producer/processor topology (AC-21):** `AddOutboxProducer` (store/router/tracker, no hosted poller) and `AddOutboxProcessor` (hosted poller) alongside `AddOutbox` (= both), each datastore-scoped via `OnDataStore(...)`. -- **Opt-in outbox metrics (AC-18):** a `RCommon.Outbox` `System.Diagnostics.Metrics.Meter` instrumenting per-datastore pending depth, oldest-unprocessed age, relay success/failure, dead-letter rate, and dispatch-queue depth. Additive; the host wires it into its metrics pipeline. -- **Payload protection hook (AC-19):** an `IOutboxPayloadProtector` seam (`Protect`/`Unprotect`) around payload serialization; default is pass-through (plaintext). -- **Deserialization allow-list (AC-20):** the outbox serializer resolves only registered event types; an unknown/tampered `EventType` is logged loud and dead-lettered rather than deserialized to an arbitrary type. - **Five runnable recipe examples with end-to-end tests (AC-16):** RCommon per-datastore outbox (+ 2-datastore variant), broker-as-producer behind the RCommon outbox (MassTransit + Wolverine), broker-native outbox (MassTransit), transaction-script router-added events, no-unit-of-work direct publish, and in-process MediatR. - **Podman/Testcontainers integration harness** (Postgres + RabbitMQ) proving cross-datastore atomicity, `TransactionScope` enlistment, and the recipe-2b coordination gate. @@ -50,6 +47,14 @@ Event-handling and transactional-outbox redesign. This is a **breaking** release - Wolverine's recipe 2b (broker-native outbox) is unsupported by design; the builder's `UseBrokerOutbox` throws `NotSupportedException`. Wolverine users get an atomic outbox via recipe 2a (broker as a producer behind RCommon's own outbox). See [Known Issues](#known-issues). +**Planned — designed but NOT shipped in 3.2.0** + +These items were specified for the redesign (AC-18/19/20) but did **not** ship in 3.2.0. An earlier draft of this changelog listed them under "Added" in error. They remain on the roadmap; do not build against them yet. + +- **Opt-in outbox metrics (AC-18):** a `RCommon.Outbox` `System.Diagnostics.Metrics.Meter` for per-datastore pending depth, oldest-unprocessed age, relay success/failure, dead-letter rate, and dispatch-queue depth. Not present in 3.2.0. For now, observe the outbox by querying the `__OutboxMessages` tables (unprocessed count, oldest `CreatedAtUtc`, dead-lettered rows) or by watching the poller's log output. +- **Payload protection hook (AC-19):** an `IOutboxPayloadProtector` (`Protect`/`Unprotect`) seam around payload serialization. Not present in 3.2.0 — outbox payloads are plaintext within the application-database trust boundary. Applications needing field/payload encryption must handle it above the event model until this ships. +- **Deserialization allow-list (AC-20):** a serializer that resolves only registered event types and dead-letters unknown/tampered `EventType`s. Not present in 3.2.0; the serializer resolves the type named in the row. + ### 3.1.3 Outbox silent-failure hardening. Every change is observability-only (new warnings and a startup diagnostic); no method or interface signatures changed and no runtime behavior changed on the success path — fully backward-compatible. diff --git a/website/versioned_docs/version-3.2.0/api-reference/changelog.mdx b/website/versioned_docs/version-3.2.0/api-reference/changelog.mdx index b9d6ab3a..1cd80145 100644 --- a/website/versioned_docs/version-3.2.0/api-reference/changelog.mdx +++ b/website/versioned_docs/version-3.2.0/api-reference/changelog.mdx @@ -28,9 +28,6 @@ Event-handling and transactional-outbox redesign. This is a **breaking** release - **Datastore-aware transactional outbox (B4/U5).** Per-datastore `__OutboxMessages` tables; each row records its target producer(s), so one table can fan out to multiple transports. `OutboxProcessingService` claims and drains every registered outbox datastore. Registering an outbox auto-applies the `OutboxMessage` mapping to that datastore's `RCommonDbContext`, and a startup diagnostic fails loud if a registered outbox datastore's table is missing. - **Broker-native outbox wrapper (`UseBrokerOutbox`).** Configures MassTransit's EF Core outbox (`AddEntityFrameworkOutbox` + `UseBusOutbox`) against an RCommon datastore's `DbContext`, so business state and broker-outbox rows commit atomically inside the RCommon unit-of-work transaction. Proven by a Postgres/Testcontainers coordination test. *(Wolverine does not support recipe 2b — see Known Issues / use recipe 2a.)* - **First-class producer/processor topology (AC-21):** `AddOutboxProducer` (store/router/tracker, no hosted poller) and `AddOutboxProcessor` (hosted poller) alongside `AddOutbox` (= both), each datastore-scoped via `OnDataStore(...)`. -- **Opt-in outbox metrics (AC-18):** a `RCommon.Outbox` `System.Diagnostics.Metrics.Meter` instrumenting per-datastore pending depth, oldest-unprocessed age, relay success/failure, dead-letter rate, and dispatch-queue depth. Additive; the host wires it into its metrics pipeline. -- **Payload protection hook (AC-19):** an `IOutboxPayloadProtector` seam (`Protect`/`Unprotect`) around payload serialization; default is pass-through (plaintext). -- **Deserialization allow-list (AC-20):** the outbox serializer resolves only registered event types; an unknown/tampered `EventType` is logged loud and dead-lettered rather than deserialized to an arbitrary type. - **Five runnable recipe examples with end-to-end tests (AC-16):** RCommon per-datastore outbox (+ 2-datastore variant), broker-as-producer behind the RCommon outbox (MassTransit + Wolverine), broker-native outbox (MassTransit), transaction-script router-added events, no-unit-of-work direct publish, and in-process MediatR. - **Podman/Testcontainers integration harness** (Postgres + RabbitMQ) proving cross-datastore atomicity, `TransactionScope` enlistment, and the recipe-2b coordination gate. @@ -50,6 +47,14 @@ Event-handling and transactional-outbox redesign. This is a **breaking** release - Wolverine's recipe 2b (broker-native outbox) is unsupported by design; the builder's `UseBrokerOutbox` throws `NotSupportedException`. Wolverine users get an atomic outbox via recipe 2a (broker as a producer behind RCommon's own outbox). See [Known Issues](#known-issues). +**Planned — designed but NOT shipped in 3.2.0** + +These items were specified for the redesign (AC-18/19/20) but did **not** ship in 3.2.0. An earlier draft of this changelog listed them under "Added" in error. They remain on the roadmap; do not build against them yet. + +- **Opt-in outbox metrics (AC-18):** a `RCommon.Outbox` `System.Diagnostics.Metrics.Meter` for per-datastore pending depth, oldest-unprocessed age, relay success/failure, dead-letter rate, and dispatch-queue depth. Not present in 3.2.0. For now, observe the outbox by querying the `__OutboxMessages` tables (unprocessed count, oldest `CreatedAtUtc`, dead-lettered rows) or by watching the poller's log output. +- **Payload protection hook (AC-19):** an `IOutboxPayloadProtector` (`Protect`/`Unprotect`) seam around payload serialization. Not present in 3.2.0 — outbox payloads are plaintext within the application-database trust boundary. Applications needing field/payload encryption must handle it above the event model until this ships. +- **Deserialization allow-list (AC-20):** a serializer that resolves only registered event types and dead-letters unknown/tampered `EventType`s. Not present in 3.2.0; the serializer resolves the type named in the row. + ### 3.1.3 Outbox silent-failure hardening. Every change is observability-only (new warnings and a startup diagnostic); no method or interface signatures changed and no runtime behavior changed on the success path — fully backward-compatible. From 6f5fe1643308779ddf9070efcaee4f76a5710b0c Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 13:15:13 -0600 Subject: [PATCH 05/10] docs(examples): add modular multi-datastore outbox example Adds Examples.EventHandling.Outbox.Modular: three independent bounded-context modules (Ordering, Billing, Shipping), each owning its own datastore and native transactional outbox, composed at a single root. This is the customer's shape (modular composition + multiple datastores + native outbox) that silently dropped durable events before 3.2.1. The runnable Program and the end-to-end test both prove every datastore's durable event persists to its own outbox and nowhere else, regardless of module registration order. --- ....EventHandling.Outbox.Modular.Tests.csproj | 25 +++++ .../ModularMultiDataStoreOutboxTests.cs | 102 ++++++++++++++++++ .../Billing/BillingModule.cs | 81 ++++++++++++++ ...amples.EventHandling.Outbox.Modular.csproj | 21 ++++ .../Ordering/OrderingModule.cs | 90 ++++++++++++++++ .../Program.cs | 99 +++++++++++++++++ .../Shipping/ShippingModule.cs | 77 +++++++++++++ Examples/Examples.sln | 30 ++++++ 8 files changed, 525 insertions(+) create mode 100644 Examples/EventHandling/Examples.EventHandling.Outbox.Modular.Tests/Examples.EventHandling.Outbox.Modular.Tests.csproj create mode 100644 Examples/EventHandling/Examples.EventHandling.Outbox.Modular.Tests/ModularMultiDataStoreOutboxTests.cs create mode 100644 Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Billing/BillingModule.cs create mode 100644 Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Examples.EventHandling.Outbox.Modular.csproj create mode 100644 Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Ordering/OrderingModule.cs create mode 100644 Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Program.cs create mode 100644 Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Shipping/ShippingModule.cs diff --git a/Examples/EventHandling/Examples.EventHandling.Outbox.Modular.Tests/Examples.EventHandling.Outbox.Modular.Tests.csproj b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular.Tests/Examples.EventHandling.Outbox.Modular.Tests.csproj new file mode 100644 index 00000000..7803c95d --- /dev/null +++ b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular.Tests/Examples.EventHandling.Outbox.Modular.Tests.csproj @@ -0,0 +1,25 @@ + + + + net10.0 + enable + enable + false + + + + + + runtime; build; native; contentfiles; analyzers; buildtransitive + all + + + + + + + + + + + diff --git a/Examples/EventHandling/Examples.EventHandling.Outbox.Modular.Tests/ModularMultiDataStoreOutboxTests.cs b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular.Tests/ModularMultiDataStoreOutboxTests.cs new file mode 100644 index 00000000..a70beee3 --- /dev/null +++ b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular.Tests/ModularMultiDataStoreOutboxTests.cs @@ -0,0 +1,102 @@ +using Examples.EventHandling.Outbox.Modular.Billing; +using Examples.EventHandling.Outbox.Modular.Ordering; +using Examples.EventHandling.Outbox.Modular.Shipping; +using FluentAssertions; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using RCommon; +using RCommon.Entities; +using RCommon.Persistence.Crud; +using RCommon.Persistence.Outbox; +using RCommon.Persistence.Transactions; +using Xunit; + +namespace Examples.EventHandling.Outbox.Modular.Tests; + +/// +/// End-to-end proof of the modular multi-datastore outbox example, and the regression guard for the +/// customer's scenario (3 bounded-context modules, each its own datastore + native outbox, composed at +/// one root). Every datastore's durable event must persist to its OWN outbox and nowhere else. Before +/// 3.2.1 this composition silently dropped durable events for datastores whose module was registered +/// after a different module's WithPersistence call. Runs on the EF Core InMemory provider (fast lane). +/// +public class ModularMultiDataStoreOutboxTests +{ + private static ServiceProvider BuildProvider() + { + var services = new ServiceCollection(); + services.AddLogging(); + + services.AddRCommon() + .WithSimpleGuidGenerator() + .WithUnitOfWork(uow => { }) + // Compose modules NOT-primary-first, to prove the fixed registration is order-independent. + .AddBillingModule(Guid.NewGuid().ToString()) + .AddShippingModule(Guid.NewGuid().ToString()) + .AddOrderingModule(Guid.NewGuid().ToString()); + + return services.BuildServiceProvider(validateScopes: true); + } + + [Fact] + public void OutboxTracker_IsAuthoritative_UnderThreeModuleComposition() + { + using var provider = BuildProvider(); + using var scope = provider.CreateScope(); + + scope.ServiceProvider.GetRequiredService() + .Should().BeOfType( + "the outbox tracker must win across a three-module composition regardless of order"); + } + + [Fact] + public async Task EachModule_PersistsItsDurableEvent_ToItsOwnOutboxOnly() + { + using var provider = BuildProvider(); + + using (var scope = provider.CreateScope()) + { + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + } + + using (var scope = provider.CreateScope()) + { + var sp = scope.ServiceProvider; + var uowFactory = sp.GetRequiredService(); + + var orders = sp.GetRequiredService>(); + orders.DataStoreName = "Ordering"; + var invoices = sp.GetRequiredService>(); + invoices.DataStoreName = "Billing"; + var shipments = sp.GetRequiredService>(); + shipments.DataStoreName = "Shipping"; + + var order = new Order { CustomerName = "Ada Lovelace", Total = 249.99m }; + order.Place(); + var invoice = new Invoice { CustomerName = "Ada Lovelace", Amount = 249.99m }; + invoice.Raise(); + var shipment = new Shipment { Destination = "London", TrackingNumber = "TRK-1" }; + shipment.Dispatch(); + + using var uow = uowFactory.Create(); + await orders.AddAsync(order); + await invoices.AddAsync(invoice); + await shipments.AddAsync(shipment); + await uow.CommitAsync(); + } + + using (var scope = provider.CreateScope()) + { + var sp = scope.ServiceProvider; + var orderingRows = await sp.GetRequiredService().Set().AsNoTracking().CountAsync(); + var billingRows = await sp.GetRequiredService().Set().AsNoTracking().CountAsync(); + var shippingRows = await sp.GetRequiredService().Set().AsNoTracking().CountAsync(); + + orderingRows.Should().Be(1, "the Ordering module's durable event must persist to the Ordering outbox"); + billingRows.Should().Be(1, "the Billing module's durable event must persist to the Billing outbox"); + shippingRows.Should().Be(1, "the Shipping module's durable event must persist to the Shipping outbox"); + } + } +} diff --git a/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Billing/BillingModule.cs b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Billing/BillingModule.cs new file mode 100644 index 00000000..9a7a942d --- /dev/null +++ b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Billing/BillingModule.cs @@ -0,0 +1,81 @@ +using Microsoft.EntityFrameworkCore; +using RCommon; +using RCommon.Entities; +using RCommon.EventHandling; +using RCommon.EventHandling.Subscribers; +using RCommon.Persistence.EFCore; +using RCommon.Persistence.EFCore.Outbox; +using RCommon.Persistence.Outbox; + +namespace Examples.EventHandling.Outbox.Modular.Billing; + +// --------------------------------------------------------------------------------------------------- +// Billing bounded context — its own datastore ("Billing") with its own transactional outbox. +// --------------------------------------------------------------------------------------------------- + +public sealed class InvoiceRaisedEvent : IDomainEvent +{ + public InvoiceRaisedEvent(Guid invoiceId, decimal amount) + { + InvoiceId = invoiceId; + Amount = amount; + } + + public Guid EventId { get; } = Guid.NewGuid(); + public DateTimeOffset OccurredOn { get; } = DateTimeOffset.UtcNow; + public Guid InvoiceId { get; } + public decimal Amount { get; } +} + +public sealed class Invoice : AggregateRoot +{ + public Invoice() : base(Guid.NewGuid()) { } + public string CustomerName { get; set; } = string.Empty; + public decimal Amount { get; set; } + + public void Raise() => AddDomainEvent(new InvoiceRaisedEvent(Id, Amount)); +} + +public sealed class InvoiceRaisedEventHandler : ISubscriber +{ + public Task HandleAsync(InvoiceRaisedEvent @event, CancellationToken cancellationToken = default) + { + Console.WriteLine($" [billing] invoice {@event.InvoiceId} raised for {@event.Amount:C}"); + return Task.CompletedTask; + } +} + +public sealed class BillingDbContext : RCommonDbContext +{ + public BillingDbContext(DbContextOptions options) : base(options) { } + public DbSet Invoices => Set(); + + protected override void OnModelCreating(ModelBuilder modelBuilder) + { + base.OnModelCreating(modelBuilder); + modelBuilder.AddOutboxMessages(); + } +} + +public static class BillingModule +{ + /// + /// Registers the Billing bounded context: its DbContext + datastore, its transactional outbox, and + /// its event route. Note there is no SetDefaultDataStore here — a non-primary module does not + /// touch the default; it just contributes its own datastore. + /// + public static IRCommonBuilder AddBillingModule(this IRCommonBuilder rcommon, string database) + { + return rcommon + .WithPersistence(ef => + { + ef.AddDbContext("Billing", o => o.UseInMemoryDatabase(database)); + ef.AddOutbox(dataStoreName: "Billing"); + }) + .WithEventHandling(events => + { + events.AddSubscriber(); + events.Publish().UseOutbox("Billing"); + }); + } +} diff --git a/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Examples.EventHandling.Outbox.Modular.csproj b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Examples.EventHandling.Outbox.Modular.csproj new file mode 100644 index 00000000..62439393 --- /dev/null +++ b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Examples.EventHandling.Outbox.Modular.csproj @@ -0,0 +1,21 @@ + + + + Exe + net10.0 + enable + enable + + + + + + + + + + + + + diff --git a/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Ordering/OrderingModule.cs b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Ordering/OrderingModule.cs new file mode 100644 index 00000000..a9763925 --- /dev/null +++ b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Ordering/OrderingModule.cs @@ -0,0 +1,90 @@ +using Microsoft.EntityFrameworkCore; +using RCommon; +using RCommon.Entities; +using RCommon.EventHandling; +using RCommon.EventHandling.Subscribers; +using RCommon.Persistence.EFCore; +using RCommon.Persistence.EFCore.Outbox; +using RCommon.Persistence.Outbox; + +namespace Examples.EventHandling.Outbox.Modular.Ordering; + +// --------------------------------------------------------------------------------------------------- +// Ordering bounded context — its own datastore ("Ordering") with its own transactional outbox. +// Everything this context needs is self-contained in this file, including its RCommon registration. +// --------------------------------------------------------------------------------------------------- + +public sealed class OrderPlacedEvent : IDomainEvent +{ + public OrderPlacedEvent(Guid orderId, decimal total) + { + OrderId = orderId; + Total = total; + } + + public Guid EventId { get; } = Guid.NewGuid(); + public DateTimeOffset OccurredOn { get; } = DateTimeOffset.UtcNow; + public Guid OrderId { get; } + public decimal Total { get; } +} + +public sealed class Order : AggregateRoot +{ + public Order() : base(Guid.NewGuid()) { } + public string CustomerName { get; set; } = string.Empty; + public decimal Total { get; set; } + + public void Place() => AddDomainEvent(new OrderPlacedEvent(Id, Total)); +} + +public sealed class OrderPlacedEventHandler : ISubscriber +{ + public Task HandleAsync(OrderPlacedEvent @event, CancellationToken cancellationToken = default) + { + Console.WriteLine($" [ordering] order {@event.OrderId} placed for {@event.Total:C}"); + return Task.CompletedTask; + } +} + +public sealed class OrderingDbContext : RCommonDbContext +{ + public OrderingDbContext(DbContextOptions options) : base(options) { } + public DbSet Orders => Set(); + + protected override void OnModelCreating(ModelBuilder modelBuilder) + { + base.OnModelCreating(modelBuilder); + modelBuilder.AddOutboxMessages(); // maps this datastore's __OutboxMessages table + } +} + +public static class OrderingModule +{ + /// + /// Registers the Ordering bounded context: its DbContext + datastore, its transactional outbox, and + /// its event route. Each module owns exactly one datastore and calls WithPersistence once — + /// the composition root simply chains the modules together (see Program). + /// + public static IRCommonBuilder AddOrderingModule(this IRCommonBuilder rcommon, string database) + { + return rcommon + .WithPersistence(ef => + { + ef.AddDbContext("Ordering", o => o.UseInMemoryDatabase(database)); + + // Ordering is the application's primary context, so it designates the default datastore. + // With more than one datastore registered, exactly one module (or the root) must set this; + // it is not inferred. The outbox poller uses it as its fallback datastore name. + ef.SetDefaultDataStore(ds => ds.DefaultDataStoreName = "Ordering"); + + // Native RCommon outbox for this datastore (producer + processor). + ef.AddOutbox(dataStoreName: "Ordering"); + }) + .WithEventHandling(events => + { + events.AddSubscriber(); + // Durability is opt-in per route: publish this context's event durably to its own outbox. + events.Publish().UseOutbox("Ordering"); + }); + } +} diff --git a/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Program.cs b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Program.cs new file mode 100644 index 00000000..64ea870f --- /dev/null +++ b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Program.cs @@ -0,0 +1,99 @@ +using Examples.EventHandling.Outbox.Modular.Billing; +using Examples.EventHandling.Outbox.Modular.Ordering; +using Examples.EventHandling.Outbox.Modular.Shipping; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using RCommon; +using RCommon.Persistence.Crud; +using RCommon.Persistence.Outbox; +using RCommon.Persistence.Transactions; + +// --------------------------------------------------------------------------------------------------- +// MODULAR, MULTI-DATASTORE TRANSACTIONAL OUTBOX +// +// This is the shape a real modular application uses: several independent bounded contexts, each owning +// its own datastore and its own transactional outbox, composed together at a single root. Each module +// calls WithPersistence/WithEventHandling for ITSELF; no module knows about the others. +// +// This composition (multiple WithPersistence calls across modules, only some of which configure an +// outbox) is exactly the shape that silently dropped durable events before 3.2.1 (see the changelog +// and Examples.EventHandling.Outbox.Tests/ReproOutboxModularDiTests). As of 3.2.1 the outbox tracker is +// registered authoritatively in the shared outbox core, so every datastore's durable events persist +// regardless of the order in which the modules are composed. +// --------------------------------------------------------------------------------------------------- + +var ordersDb = Guid.NewGuid().ToString(); +var billingDb = Guid.NewGuid().ToString(); +var shippingDb = Guid.NewGuid().ToString(); + +var host = Host.CreateDefaultBuilder(args) + .ConfigureLogging(logging => logging.AddFilter("Microsoft.EntityFrameworkCore", LogLevel.Warning)) + .ConfigureServices(services => + { + services.AddRCommon() + .WithSimpleGuidGenerator() // the outbox router stamps each row's Id + .WithUnitOfWork(uow => { }) + // Compose the modules. Order is intentionally NOT Ordering-first to demonstrate that the + // fixed registration is order-independent. + .AddBillingModule(billingDb) + .AddShippingModule(shippingDb) + .AddOrderingModule(ordersDb); + }) + .Build(); + +Console.WriteLine("Modular multi-datastore outbox example starting"); + +// Create each datastore's schema (including its __OutboxMessages table). +using (var scope = host.Services.CreateScope()) +{ + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); +} + +// Start the host so the single OutboxProcessingService (which drains every registered datastore) runs. +await host.StartAsync(); + +// Commit one aggregate per bounded context. Each raises a durable domain event that must land in ITS +// OWN datastore's outbox — never another's. +using (var scope = host.Services.CreateScope()) +{ + var sp = scope.ServiceProvider; + var uowFactory = sp.GetRequiredService(); + + var orders = sp.GetRequiredService>(); + orders.DataStoreName = "Ordering"; + var invoices = sp.GetRequiredService>(); + invoices.DataStoreName = "Billing"; + var shipments = sp.GetRequiredService>(); + shipments.DataStoreName = "Shipping"; + + var order = new Order { CustomerName = "Ada Lovelace", Total = 249.99m }; + order.Place(); + var invoice = new Invoice { CustomerName = "Ada Lovelace", Amount = 249.99m }; + invoice.Raise(); + var shipment = new Shipment { Destination = "London", TrackingNumber = "TRK-1" }; + shipment.Dispatch(); + + using var uow = uowFactory.Create(); + await orders.AddAsync(order); + await invoices.AddAsync(invoice); + await shipments.AddAsync(shipment); + await uow.CommitAsync(); +} + +// Show each datastore's outbox row count — all three persisted, none leaked into another datastore. +using (var scope = host.Services.CreateScope()) +{ + var sp = scope.ServiceProvider; + var orderingRows = await sp.GetRequiredService().Set().AsNoTracking().CountAsync(); + var billingRows = await sp.GetRequiredService().Set().AsNoTracking().CountAsync(); + var shippingRows = await sp.GetRequiredService().Set().AsNoTracking().CountAsync(); + + Console.WriteLine($"Outbox rows -> Ordering: {orderingRows}, Billing: {billingRows}, Shipping: {shippingRows}"); +} + +await host.StopAsync(); +Console.WriteLine("Modular multi-datastore outbox example complete"); diff --git a/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Shipping/ShippingModule.cs b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Shipping/ShippingModule.cs new file mode 100644 index 00000000..c476de11 --- /dev/null +++ b/Examples/EventHandling/Examples.EventHandling.Outbox.Modular/Shipping/ShippingModule.cs @@ -0,0 +1,77 @@ +using Microsoft.EntityFrameworkCore; +using RCommon; +using RCommon.Entities; +using RCommon.EventHandling; +using RCommon.EventHandling.Subscribers; +using RCommon.Persistence.EFCore; +using RCommon.Persistence.EFCore.Outbox; +using RCommon.Persistence.Outbox; + +namespace Examples.EventHandling.Outbox.Modular.Shipping; + +// --------------------------------------------------------------------------------------------------- +// Shipping bounded context — its own datastore ("Shipping") with its own transactional outbox. A third +// datastore makes the point that the pattern scales to N modules with no cross-module coordination. +// --------------------------------------------------------------------------------------------------- + +public sealed class ShipmentDispatchedEvent : IDomainEvent +{ + public ShipmentDispatchedEvent(Guid shipmentId, string trackingNumber) + { + ShipmentId = shipmentId; + TrackingNumber = trackingNumber; + } + + public Guid EventId { get; } = Guid.NewGuid(); + public DateTimeOffset OccurredOn { get; } = DateTimeOffset.UtcNow; + public Guid ShipmentId { get; } + public string TrackingNumber { get; } = string.Empty; +} + +public sealed class Shipment : AggregateRoot +{ + public Shipment() : base(Guid.NewGuid()) { } + public string Destination { get; set; } = string.Empty; + public string TrackingNumber { get; set; } = string.Empty; + + public void Dispatch() => AddDomainEvent(new ShipmentDispatchedEvent(Id, TrackingNumber)); +} + +public sealed class ShipmentDispatchedEventHandler : ISubscriber +{ + public Task HandleAsync(ShipmentDispatchedEvent @event, CancellationToken cancellationToken = default) + { + Console.WriteLine($" [shipping] shipment {@event.ShipmentId} dispatched ({@event.TrackingNumber})"); + return Task.CompletedTask; + } +} + +public sealed class ShippingDbContext : RCommonDbContext +{ + public ShippingDbContext(DbContextOptions options) : base(options) { } + public DbSet Shipments => Set(); + + protected override void OnModelCreating(ModelBuilder modelBuilder) + { + base.OnModelCreating(modelBuilder); + modelBuilder.AddOutboxMessages(); + } +} + +public static class ShippingModule +{ + public static IRCommonBuilder AddShippingModule(this IRCommonBuilder rcommon, string database) + { + return rcommon + .WithPersistence(ef => + { + ef.AddDbContext("Shipping", o => o.UseInMemoryDatabase(database)); + ef.AddOutbox(dataStoreName: "Shipping"); + }) + .WithEventHandling(events => + { + events.AddSubscriber(); + events.Publish().UseOutbox("Shipping"); + }); + } +} diff --git a/Examples/Examples.sln b/Examples/Examples.sln index 0e05623d..419d9187 100644 --- a/Examples/Examples.sln +++ b/Examples/Examples.sln @@ -195,6 +195,10 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "RCommon.MassTransit.Outbox" EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Examples.Messaging.MassTransit.NativeOutbox.Tests", "Messaging\Examples.Messaging.MassTransit.NativeOutbox.Tests\Examples.Messaging.MassTransit.NativeOutbox.Tests.csproj", "{C36A7AFB-2191-4D0F-A7DD-6EB48F0463CE}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Examples.EventHandling.Outbox.Modular", "EventHandling\Examples.EventHandling.Outbox.Modular\Examples.EventHandling.Outbox.Modular.csproj", "{3E253BFF-0598-4321-AAAC-D03AF9F987DF}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Examples.EventHandling.Outbox.Modular.Tests", "EventHandling\Examples.EventHandling.Outbox.Modular.Tests\Examples.EventHandling.Outbox.Modular.Tests.csproj", "{24D04D1C-C689-46FD-9A9C-8342A5CDC318}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -1141,6 +1145,30 @@ Global {C36A7AFB-2191-4D0F-A7DD-6EB48F0463CE}.Release|x64.Build.0 = Release|Any CPU {C36A7AFB-2191-4D0F-A7DD-6EB48F0463CE}.Release|x86.ActiveCfg = Release|Any CPU {C36A7AFB-2191-4D0F-A7DD-6EB48F0463CE}.Release|x86.Build.0 = Release|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Debug|Any CPU.Build.0 = Debug|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Debug|x64.ActiveCfg = Debug|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Debug|x64.Build.0 = Debug|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Debug|x86.ActiveCfg = Debug|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Debug|x86.Build.0 = Debug|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Release|Any CPU.ActiveCfg = Release|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Release|Any CPU.Build.0 = Release|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Release|x64.ActiveCfg = Release|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Release|x64.Build.0 = Release|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Release|x86.ActiveCfg = Release|Any CPU + {3E253BFF-0598-4321-AAAC-D03AF9F987DF}.Release|x86.Build.0 = Release|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Debug|Any CPU.Build.0 = Debug|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Debug|x64.ActiveCfg = Debug|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Debug|x64.Build.0 = Debug|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Debug|x86.ActiveCfg = Debug|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Debug|x86.Build.0 = Debug|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Release|Any CPU.ActiveCfg = Release|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Release|Any CPU.Build.0 = Release|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Release|x64.ActiveCfg = Release|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Release|x64.Build.0 = Release|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Release|x86.ActiveCfg = Release|Any CPU + {24D04D1C-C689-46FD-9A9C-8342A5CDC318}.Release|x86.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE @@ -1220,6 +1248,8 @@ Global {070131F8-9D2A-41D4-A76E-67582EF68B71} = {AA81E3CF-C1BB-70B4-13CB-5131BB298347} {60437E30-F5BA-4DC9-8325-01FB4CA2D742} = {AA81E3CF-C1BB-70B4-13CB-5131BB298347} {C36A7AFB-2191-4D0F-A7DD-6EB48F0463CE} = {AA81E3CF-C1BB-70B4-13CB-5131BB298347} + {3E253BFF-0598-4321-AAAC-D03AF9F987DF} = {48A3FC7F-661D-68E4-77B8-76A1C7AA7E61} + {24D04D1C-C689-46FD-9A9C-8342A5CDC318} = {48A3FC7F-661D-68E4-77B8-76A1C7AA7E61} EndGlobalSection GlobalSection(ExtensibilityGlobals) = postSolution SolutionGuid = {0B0CD26D-8067-4667-863E-6B0EE7EDAA42} From 02ddc212977faf27a0ef44d2d87d357d66efb1d4 Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 13:17:48 -0600 Subject: [PATCH 06/10] docs(outbox): clarify AddOutbox idempotency note now that routing uses Remove-then-Add --- .../Outbox/OutboxPersistenceBuilderExtensions.cs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs b/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs index af087b71..ba66f88b 100644 --- a/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs +++ b/Src/RCommon.Persistence/Outbox/OutboxPersistenceBuilderExtensions.cs @@ -16,8 +16,9 @@ public static class OutboxPersistenceBuilderExtensions /// then . /// /// - /// All registrations are idempotent (Try*/TryAddEnumerable) and datastore names - /// deduplicate case-insensitively, so composing producer + processor (which each register the + /// Every registration is idempotent — via Try*/TryAddEnumerable, or (for the + /// authoritative outbox routing) a Remove-then-Add that nets exactly one registration — and datastore + /// names deduplicate case-insensitively, so composing producer + processor (which each register the /// shared core) is safe and yields exactly one poller and one datastore registration. /// /// The delegate may be invoked more than once (once eagerly to resolve From dc4ea4631b2a8a9d6b7b7a5174064c4014318c81 Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 13:37:33 -0600 Subject: [PATCH 07/10] docs(plan): record 3.2.1 verification (unit, examples, outbox integration) --- .../2026-07-23-outbox-di-registration-3.2.1.md | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/docs/superpowers/plans/2026-07-23-outbox-di-registration-3.2.1.md b/docs/superpowers/plans/2026-07-23-outbox-di-registration-3.2.1.md index ddec2d2f..3735884e 100644 --- a/docs/superpowers/plans/2026-07-23-outbox-di-registration-3.2.1.md +++ b/docs/superpowers/plans/2026-07-23-outbox-di-registration-3.2.1.md @@ -88,6 +88,19 @@ ## Task 7: green + finish -- [ ] Run full solution build + the event-handling/persistence test suites (fast lane, `Category!=Integration`). -- [ ] Squash interim commits into one meaningful 3.2.1 commit; push branch; open PR. +- [x] Run full solution build + the event-handling/persistence test suites (fast lane, `Category!=Integration`). +- [ ] Push branch; open PR. (Deferred — user reviewing commits first.) - [ ] Do **not** tag/release until the user approves the PR. + +## Verification (2026-07-23) + +**Unit (fast lane, `Category!=Integration`) — all green:** +RCommon.Persistence 453 · RCommon.Core 575 · RCommon.EfCore 158 (+3 skipped) · RCommon.Entities 236 · RCommon.Wolverine.Outbox 9 · RCommon.Dapper 92 · RCommon.Linq2Db 73. + +**Examples (whole `Examples.sln`, fast lane) — all green:** +Outbox 6 (incl. the #15 & #16 repros) · Outbox.Modular 2 · MediatR 1 · NoUnitOfWork 2 · TransactionScript 1 · DomainDrivenDesign 10 · Persistence.Sagas 3 · HR.LeaveManagement 4. Full `Examples.sln` builds. Modular example console app runs clean (1 outbox row per datastore, none leaked). + +**Integration (Podman → Postgres/RabbitMQ) — outbox path all green when containers start:** +`CrossDataStoreOutboxTests` 2/2 · `RecipeTwoBBrokerOutboxTests` 2/2 · `MassTransitOutboxCoordinationSpikeTests` 2/2 · `WolverineOutboxCoordinationSpikeTests` 3/3. Running all four collections at once produced 5 `PostgreSqlFixture.InitializeAsync` container-start timeouts (named-pipe connect) — an environmental 2 GiB-Podman contention issue, not an assertion/DI failure; each class passes when run alone. + +**Not run (out of scope for the outbox DI change):** S3 (LocalStack) / Azure (Azurite) blob-storage integration tests — unrelated subsystem, require emulators. From 0fa06a7ea127991b6bc6e47602461f6f9fd5912b Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 13:45:59 -0600 Subject: [PATCH 08/10] test(outbox): pin multi-datastore unit-of-work transaction semantics MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds an integration test documenting the contract for a unit of work spanning more than one datastore (the heterogeneous-engine production question): - a single UnitOfWork writing two datastores promotes to a distributed transaction and fails loud (TransactionAbortedException / Postgres 55000 prepared transactions disabled) — RCommon gives no cross-datastore atomicity and never half-commits; - the correct pattern is a unit of work PER datastore, each committing locally and atomically with its own co-located outbox row. Modelled with two Postgres databases (any second connection triggers the same promotion; the engine pairing is not the trigger), reproducing the exact error a Postgres + SQL Server deployment reported. --- .../MultiDataStoreUnitOfWorkSemanticsTests.cs | 234 ++++++++++++++++++ 1 file changed, 234 insertions(+) create mode 100644 Tests/RCommon.IntegrationTests/MultiDataStoreUnitOfWorkSemanticsTests.cs diff --git a/Tests/RCommon.IntegrationTests/MultiDataStoreUnitOfWorkSemanticsTests.cs b/Tests/RCommon.IntegrationTests/MultiDataStoreUnitOfWorkSemanticsTests.cs new file mode 100644 index 00000000..f586f8de --- /dev/null +++ b/Tests/RCommon.IntegrationTests/MultiDataStoreUnitOfWorkSemanticsTests.cs @@ -0,0 +1,234 @@ +using System; +using System.Threading.Tasks; +using System.Transactions; +using FluentAssertions; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Npgsql; +using RCommon.Entities; +using RCommon.EventHandling; +using RCommon.IntegrationTests.Fixtures; +using RCommon.Persistence.Crud; +using RCommon.Persistence.EFCore; +using RCommon.Persistence.EFCore.Outbox; +using RCommon.Persistence.Outbox; +using RCommon.Persistence.Transactions; +using Xunit; + +namespace RCommon.IntegrationTests; + +/// +/// Pins the transactional contract for a unit of work that spans MORE THAN ONE datastore — the question +/// a heterogeneous-engine deployment (e.g. Postgres + SQL Server) runs into in production. +/// +/// THE CONTRACT (what these tests assert): +/// 1. An outbox row is atomic with ITS OWN datastore's state change — they share one connection/local +/// transaction (co-location, spec MN-4). The outbox never enlists a second resource manager. +/// 2. There is NO distributed atomicity across datastores. RCommon's wraps a +/// single ambient (Required). If one unit of work writes to a SECOND +/// datastore, a second physical connection enlists and the ambient transaction promotes to a +/// distributed transaction (2PC). On Postgres that means PREPARE TRANSACTION, which is +/// disabled by default (max_prepared_transactions = 0) and fails loud — it never silently +/// half-commits. (Postgres + SQL Server would promote to MSDTC with the same "no cross-datastore +/// atomicity" outcome. Two Postgres databases reproduce the identical promotion because the trigger +/// is "a second connection in one scope," not the engine pairing.) +/// 3. Therefore an application that writes state to more than one datastore in one logical operation +/// must split it into a UNIT OF WORK PER DATASTORE. Each commits locally and atomically (state + +/// that datastore's outbox row); cross-datastore delivery is eventual (at-least-once via the poller). +/// +/// Modelled as two distinct Postgres databases on one container (each context points at its own +/// Database=), so each owns an isolated __OutboxMessages table. +/// +[Trait("Category", "Integration")] +[Collection(PostgreSqlCollection.Name)] +public class MultiDataStoreUnitOfWorkSemanticsTests +{ + private readonly PostgreSqlFixture _pg; + + public MultiDataStoreUnitOfWorkSemanticsTests(PostgreSqlFixture pg) => _pg = pg; + + // ---- Ordering datastore ---- + public sealed class OrderPlacedEvent : IDomainEvent + { + public OrderPlacedEvent(Guid orderId) => OrderId = orderId; + public Guid EventId { get; } = Guid.NewGuid(); + public DateTimeOffset OccurredOn { get; } = DateTimeOffset.UtcNow; + public Guid OrderId { get; } + } + + public sealed class Order : AggregateRoot + { + public Order() : base(Guid.NewGuid()) { } + public string CustomerName { get; set; } = string.Empty; + public void Place() => AddDomainEvent(new OrderPlacedEvent(Id)); + } + + public sealed class OrdersDbContext : RCommonDbContext + { + public OrdersDbContext(DbContextOptions options) : base(options) { } + public DbSet Orders => Set(); + } + + // ---- Billing datastore (a genuinely separate database) ---- + public sealed class InvoiceRaisedEvent : IDomainEvent + { + public InvoiceRaisedEvent(Guid invoiceId) => InvoiceId = invoiceId; + public Guid EventId { get; } = Guid.NewGuid(); + public DateTimeOffset OccurredOn { get; } = DateTimeOffset.UtcNow; + public Guid InvoiceId { get; } + } + + public sealed class Invoice : AggregateRoot + { + public Invoice() : base(Guid.NewGuid()) { } + public string CustomerName { get; set; } = string.Empty; + public void Raise() => AddDomainEvent(new InvoiceRaisedEvent(Id)); + } + + public sealed class BillingDbContext : RCommonDbContext + { + public BillingDbContext(DbContextOptions options) : base(options) { } + public DbSet Invoices => Set(); + } + + private ServiceProvider BuildProvider(string ordersConnectionString, string billingConnectionString) + { + var services = new ServiceCollection(); + services.AddLogging(); + + services.AddRCommon() + .WithSimpleGuidGenerator() + .WithUnitOfWork(uow => { }) + .WithEventHandling(events => + { + // Each event is durable to ITS OWN co-located datastore's outbox. + events.Publish().UseOutbox("Orders"); + events.Publish().UseOutbox("Billing"); + }) + .WithPersistence(ef => + { + ef.AddDbContext("Orders", o => o.UseNpgsql(ordersConnectionString)); + ef.AddDbContext("Billing", o => o.UseNpgsql(billingConnectionString)); + ef.SetDefaultDataStore(ds => ds.DefaultDataStoreName = "Orders"); + ef.AddOutbox(dataStoreName: "Orders"); + ef.AddOutbox(dataStoreName: "Billing"); + }); + + return services.BuildServiceProvider(validateScopes: true); + } + + private static async Task EnsureSchemasAsync(IServiceProvider provider) + { + using var scope = provider.CreateScope(); + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + await scope.ServiceProvider.GetRequiredService().Database.EnsureCreatedAsync(); + } + + private static async Task<(int orders, int ordersOutbox, int invoices, int billingOutbox)> CountAsync(IServiceProvider provider) + { + using var scope = provider.CreateScope(); + var orders = scope.ServiceProvider.GetRequiredService(); + var billing = scope.ServiceProvider.GetRequiredService(); + return ( + await orders.Set().AsNoTracking().CountAsync(), + await orders.Set().AsNoTracking().CountAsync(), + await billing.Set().AsNoTracking().CountAsync(), + await billing.Set().AsNoTracking().CountAsync()); + } + + [Fact] + public async Task Single_unit_of_work_writing_two_datastores_fails_loud_and_commits_nothing() + { + var ordersCs = await _pg.CreateUniqueDatabaseAsync("orders"); + var billingCs = await _pg.CreateUniqueDatabaseAsync("billing"); + await using var provider = BuildProvider(ordersCs, billingCs); + await EnsureSchemasAsync(provider); + + TransactionAbortedException? thrown = null; + try + { + using var scope = provider.CreateScope(); + var sp = scope.ServiceProvider; + var uowFactory = sp.GetRequiredService(); + + var orders = sp.GetRequiredService>(); + orders.DataStoreName = "Orders"; + var invoices = sp.GetRequiredService>(); + invoices.DataStoreName = "Billing"; + + var order = new Order { CustomerName = "Ada Lovelace" }; + order.Place(); + var invoice = new Invoice { CustomerName = "Ada Lovelace" }; + invoice.Raise(); + + using var uow = uowFactory.Create(); + await orders.AddAsync(order); // enlists the Orders connection + await invoices.AddAsync(invoice); // enlists a SECOND connection => ambient tx promotes to 2PC + await uow.CommitAsync(); // PREPARE TRANSACTION rejected by Postgres + } + catch (TransactionAbortedException ex) + { + thrown = ex; + } + + thrown.Should().NotBeNull( + "a single unit of work spanning two datastores promotes to a distributed transaction; RCommon " + + "provides NO cross-datastore atomicity and must fail loud rather than half-commit"); + thrown!.InnerException.Should().BeOfType() + .Which.SqlState.Should().Be("55000", + "Postgres rejects PREPARE TRANSACTION when prepared transactions are disabled (the default)"); + + var (orderRows, ordersOutbox, invoiceRows, billingOutbox) = await CountAsync(provider); + orderRows.Should().Be(0, "the aborted transaction must leave no Orders state"); + ordersOutbox.Should().Be(0, "the aborted transaction must leave no Orders outbox row"); + invoiceRows.Should().Be(0, "the aborted transaction must leave no Billing state"); + billingOutbox.Should().Be(0, "the aborted transaction must leave no Billing outbox row"); + } + + [Fact] + public async Task Per_datastore_unit_of_work_scopes_each_commit_locally_with_their_own_outbox_row() + { + var ordersCs = await _pg.CreateUniqueDatabaseAsync("orders"); + var billingCs = await _pg.CreateUniqueDatabaseAsync("billing"); + await using var provider = BuildProvider(ordersCs, billingCs); + await EnsureSchemasAsync(provider); + + // ONE logical operation, split into a unit of work PER datastore. Each is a single connection, so + // each commits with a local transaction (no promotion) — state + that datastore's outbox row. + using (var scope = provider.CreateScope()) + { + var sp = scope.ServiceProvider; + var uowFactory = sp.GetRequiredService(); + var orders = sp.GetRequiredService>(); + orders.DataStoreName = "Orders"; + + var order = new Order { CustomerName = "Ada Lovelace" }; + order.Place(); + + using var uow = uowFactory.Create(); + await orders.AddAsync(order); + await uow.CommitAsync(); + } + + using (var scope = provider.CreateScope()) + { + var sp = scope.ServiceProvider; + var uowFactory = sp.GetRequiredService(); + var invoices = sp.GetRequiredService>(); + invoices.DataStoreName = "Billing"; + + var invoice = new Invoice { CustomerName = "Ada Lovelace" }; + invoice.Raise(); + + using var uow = uowFactory.Create(); + await invoices.AddAsync(invoice); + await uow.CommitAsync(); + } + + var (orderRows, ordersOutbox, invoiceRows, billingOutbox) = await CountAsync(provider); + orderRows.Should().Be(1, "the Orders unit of work committed locally"); + ordersOutbox.Should().Be(1, "the Order's durable event is atomic with the Orders state change"); + invoiceRows.Should().Be(1, "the Billing unit of work committed locally"); + billingOutbox.Should().Be(1, "the Invoice's durable event is atomic with the Billing state change"); + } +} From 2c1f0cd35013265fc5fa41b569a3ce6eb4d9cc58 Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 13:46:57 -0600 Subject: [PATCH 09/10] docs(outbox): document multi-datastore transaction boundaries Adds a caution to the outbox recipes guide (live + versioned 3.2.0): an outbox row is atomic with its own datastore only; a unit of work spanning two datastores promotes to a distributed transaction (2PC/MSDTC) and fails loud (e.g. Postgres 55000 prepared transactions disabled). One logical operation touching multiple datastores must use a unit of work per datastore; cross- datastore delivery is eventual via the poller. --- website/docs/event-handling/recipes.mdx | 6 ++++++ .../versioned_docs/version-3.2.0/event-handling/recipes.mdx | 6 ++++++ 2 files changed, 12 insertions(+) diff --git a/website/docs/event-handling/recipes.mdx b/website/docs/event-handling/recipes.mdx index 7bdf7b74..8f8deabe 100644 --- a/website/docs/event-handling/recipes.mdx +++ b/website/docs/event-handling/recipes.mdx @@ -137,6 +137,12 @@ flowchart LR BT -.->|poller| BUSB(("relay")) ``` +:::caution One transaction per datastore — no cross-datastore atomicity +An outbox row is atomic with **its own** datastore's state change — they share one connection and one local transaction (co-location). There is **no** distributed atomicity **across** datastores. `UnitOfWork` wraps a single ambient `TransactionScope` (`Required`), so if one unit of work writes to a **second** datastore, a second connection enlists and the transaction **promotes to a distributed transaction (2PC/MSDTC)**. Across engines this is the classic failure — e.g. Postgres + SQL Server, where Npgsql issues `PREPARE TRANSACTION` and Postgres rejects it with `55000: prepared transactions are disabled` (and two databases on the *same* engine promote identically; the trigger is the second connection, not the engine pairing). + +So when one logical operation touches more than one datastore, **use a unit of work per datastore**. Each commits locally and atomically (state + that datastore's co-located outbox row); cross-datastore delivery is **eventual** (at-least-once via the poller). RCommon never forces cross-datastore 2PC, and it fails loud rather than half-committing — see `MultiDataStoreUnitOfWorkSemanticsTests`. +::: + --- ## Recipe 2a — broker as producer behind the RCommon outbox diff --git a/website/versioned_docs/version-3.2.0/event-handling/recipes.mdx b/website/versioned_docs/version-3.2.0/event-handling/recipes.mdx index 7bdf7b74..62c256f2 100644 --- a/website/versioned_docs/version-3.2.0/event-handling/recipes.mdx +++ b/website/versioned_docs/version-3.2.0/event-handling/recipes.mdx @@ -137,6 +137,12 @@ flowchart LR BT -.->|poller| BUSB(("relay")) ``` +:::caution One transaction per datastore — no cross-datastore atomicity +An outbox row is atomic with **its own** datastore's state change — they share one connection and one local transaction (co-location). There is **no** distributed atomicity **across** datastores. `UnitOfWork` wraps a single ambient `TransactionScope` (`Required`), so if one unit of work writes to a **second** datastore, a second connection enlists and the transaction **promotes to a distributed transaction (2PC/MSDTC)**. Across engines this is the classic failure — e.g. Postgres + SQL Server, where Npgsql issues `PREPARE TRANSACTION` and Postgres rejects it with `55000: prepared transactions are disabled` (and two databases on the *same* engine promote identically; the trigger is the second connection, not the engine pairing). + +So when one logical operation touches more than one datastore, **use a unit of work per datastore**. Each commits locally and atomically (state + that datastore's co-located outbox row); cross-datastore delivery is **eventual** (at-least-once via the poller). RCommon never forces cross-datastore 2PC, and it fails loud rather than half-committing. +::: + --- ## Recipe 2a — broker as producer behind the RCommon outbox From 9076888110b3df23ad92bcdde2e88982dfb09a9a Mon Sep 17 00:00:00 2001 From: jasonmwebb-lv Date: Thu, 23 Jul 2026 13:57:41 -0600 Subject: [PATCH 10/10] docs(3.2.1): add changelog entry, migration note, and modular-example reference - Changelog: new 3.2.1 section (silent-data-loss fixes #15/#16, routing diagnostic correction, modular example, multi-datastore semantics test/note, #17 doc correction). - Migration guide: 'Upgrading to 3.2.1' (no code changes; remove any outbox Replace workaround; multi-datastore transaction-boundary guidance). - Recipes: reference the runnable Examples.EventHandling.Outbox.Modular from the Multiple datastores section. Verified with a clean Docusaurus build (onBrokenLinks: throw). --- website/docs/api-reference/changelog.mdx | 19 +++++++++++++++++++ .../docs/api-reference/migration-guide.mdx | 8 ++++++++ website/docs/event-handling/recipes.mdx | 2 ++ 3 files changed, 29 insertions(+) diff --git a/website/docs/api-reference/changelog.mdx b/website/docs/api-reference/changelog.mdx index 1cd80145..0f0b34b7 100644 --- a/website/docs/api-reference/changelog.mdx +++ b/website/docs/api-reference/changelog.mdx @@ -12,6 +12,25 @@ The authoritative release history for all RCommon packages is maintained on GitH ## Recent Changes +### 3.2.1 + +Outbox DI-registration hardening and documentation accuracy. This is a behavior-only bug-fix release — no public API changes; upgrading from 3.2.0 requires no code changes. + +**Fixed** + +- **Durable events were silently dropped under modular / multi-datastore composition.** `WithEventTracking` runs at the end of every `WithPersistence` call and registered the in-memory entity-event tracker with `TryAdd`; when a non-outbox module (or a second datastore) was registered before the outbox one, the in-memory tracker was pinned and the outbox's own `TryAdd` no-opped — so every durable event dispatched in-process and was never written to the outbox, with no error. The outbox now registers the entity-event tracker and router **authoritatively** (remove-then-add) in the shared outbox core, so it wins regardless of registration order. If you added a `services.Replace(...)` workaround for this, you can remove it after upgrading. +- **Processor and combined hosts silently dropped durable events.** Outbox routing was registered only by `AddOutboxProducer`; a host that ran the poller (`AddOutboxProcessor`) and also committed domain entities lost its durable events. Routing now lives in the shared core that both `AddOutboxProducer` and `AddOutboxProcessor` call, so producer, processor, and combined hosts all persist. This also removes the registration path that could surface as an "unable to resolve `InMemoryEntityEventTracker`" DI error. +- **Outbox routing startup diagnostic inspected the wrong service.** It checked `IEventRouter` (which the core always binds to the in-memory router) and so false-warned on every working outbox host while missing the real failure. It now inspects the load-bearing `IEntityEventTracker`. + +**Added** + +- **Modular multi-datastore outbox example** (`Examples.EventHandling.Outbox.Modular`): three bounded-context modules, each owning its own datastore and native outbox, composed at one root — with an end-to-end test proving every datastore's durable events persist regardless of module registration order. +- **Multi-datastore unit-of-work semantics** are now pinned by an integration test and documented in a "One transaction per datastore" note in the [Event Handling Recipes](../event-handling/recipes.mdx#multiple-datastores) guide: an outbox row is atomic with **its own** datastore's state change (local transaction); a single unit of work spanning two datastores promotes to a distributed transaction and **fails loud** (e.g. Postgres `55000: prepared transactions are disabled`); cross-datastore delivery is **eventual** (at-least-once via the poller). Split a multi-datastore operation into a unit of work per datastore. + +**Docs** + +- Corrected the 3.2.0 changelog and spec, which listed outbox metrics, `IOutboxPayloadProtector`, and the deserialization allow-list as shipped features — they are **planned, not shipped** (see the 3.2.0 "Planned" note below). + ### 3.2.0 Event-handling and transactional-outbox redesign. This is a **breaking** release: durability is now opt-in per route, in-process domain handlers dispatch *before* commit, and the outbox is datastore-aware end to end. A new recipe-based conceptual guide documents the supported patterns. See the [Event Handling Recipes](../event-handling/recipes.mdx) guide and the [3.2.0 migration section](./migration-guide.mdx#upgrading-to-320-event-handling--outbox-redesign) for the full old→new mapping. diff --git a/website/docs/api-reference/migration-guide.mdx b/website/docs/api-reference/migration-guide.mdx index 13ae21b5..86925bf7 100644 --- a/website/docs/api-reference/migration-guide.mdx +++ b/website/docs/api-reference/migration-guide.mdx @@ -171,6 +171,14 @@ Two limitations ship with 3.2.0, each with a supported workaround — see [Known - **Dapper outbox store is SQL Server–only.** On PostgreSQL, use `EFCoreOutboxStore` for the outbox datastore. - **MediatR `Send()` may not resolve handlers.** Use `Publish()` for in-process MediatR event delivery. +## Upgrading to 3.2.1 + +3.2.1 is a behavior-only bug-fix release with **no public API changes** — upgrading from 3.2.0 requires no code changes. It fixes silent outbox data loss in two compositions and clarifies multi-datastore transaction semantics. See the [3.2.1 changelog](./changelog.mdx#321) for the full list. + +- **Remove any outbox-routing workaround.** If you worked around durable events not persisting under modular / multi-datastore composition (e.g. a `services.Replace(...)` re-asserting the outbox `IEntityEventTracker`/`IEventRouter`), you can delete it — the outbox now registers those authoritatively and wins regardless of module/datastore registration order. +- **Processor and combined hosts now persist durable events.** If you had split producer and processor and found a host that ran the poller and also committed domain entities was losing events (or threw "unable to resolve `InMemoryEntityEventTracker`"), that is fixed; outbox routing now lives in the shared core that `AddOutboxProducer` and `AddOutboxProcessor` both call. +- **Review multi-datastore transaction boundaries.** A single unit of work must not write to more than one datastore — that promotes to a distributed transaction and fails loud (e.g. Postgres `55000: prepared transactions are disabled`). Split a multi-datastore operation into a unit of work per datastore; each commits locally and atomically with its own outbox row, and cross-datastore delivery is eventual via the poller. See the ["One transaction per datastore" note](../event-handling/recipes.mdx#multiple-datastores). + ## Breaking Change Patterns The following sections describe the categories of breaking changes that appear in major RCommon releases and how to address them. diff --git a/website/docs/event-handling/recipes.mdx b/website/docs/event-handling/recipes.mdx index 8f8deabe..ed7592ba 100644 --- a/website/docs/event-handling/recipes.mdx +++ b/website/docs/event-handling/recipes.mdx @@ -119,6 +119,8 @@ Once an event type is declared durable, it is buffered into the outbox **within Register an outbox per datastore. An event published `.UseOutbox("Billing")` lands **only** in Billing's outbox — never in Orders'. +For a full worked version — several bounded-context modules, each owning its own datastore and outbox and composed at one root — see the runnable `Examples.EventHandling.Outbox.Modular` example. + ```csharp .WithPersistence(ef => {