diff --git a/Directory.Packages.props b/Directory.Packages.props index 258127b76..0e36fd163 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -70,7 +70,7 @@ - + diff --git a/docs/guide/messaging/transports/nats.md b/docs/guide/messaging/transports/nats.md index cddf56e28..64b4a85bc 100644 --- a/docs/guide/messaging/transports/nats.md +++ b/docs/guide/messaging/transports/nats.md @@ -299,6 +299,24 @@ opts.ListenToNatsSubject("orders.received") .Named("orders-listener"); ``` +### Load Balancing with Queue Groups + +Multiple listeners sharing a NATS queue group have each message delivered to only one member, spreading load +across instances. Set a transport-wide default so every listener joins the same group: + +```csharp +opts.UseNats(nats => +{ + nats.ConnectionString = "nats://localhost:4222"; + nats.DefaultQueueGroup = "orders-workers"; +}); +``` + +::: tip Subject normalization +By default the transport normalizes `/` separators in subjects to NATS `.` tokens (`NormalizeSubjects`, on by +default). Set `nats.NormalizeSubjects = false` if you need to use literal subjects that contain `/`. +::: + ## Publishing Messages ### To a Specific Subject @@ -325,6 +343,88 @@ opts.PublishAllMessages() .SendInline(); ``` +### Static Outgoing Headers + +Attach a constant header to every message published to a subject with `AddOutgoingHeader`: + +```csharp +opts.PublishMessage() + .ToNatsSubject("orders.created") + .AddOutgoingHeader("x-source", "orders-service"); +``` + +### Per-Message (Dynamic) Subjects + +`ToNatsSubject("...")` publishes to a single static subject. To compute the subject *per message* — e.g. an +aggregate-scoped subject like `orders.events.{id}` — use `PublishMessagesToNatsSubject`. This is built on +Wolverine's generic topic routing (`RoutingMode.ByTopic` / `Envelope.TopicName`), the same mechanism the +RabbitMQ, Kafka, and MQTT transports use, so it also participates in `IMessageBus.BroadcastToTopicAsync`. + +```csharp +opts.UseNats("nats://localhost:4222").AutoProvision(); + +// The subject is derived from each message instance. +opts.PublishMessagesToNatsSubject(e => $"orders.events.{e.OrderId}"); +``` + +The same endpoint is automatically enrolled for explicit topic broadcasts, where the caller supplies the +subject directly (overriding the function): + +```csharp +await bus.BroadcastToTopicAsync("orders.events.12345", new OrderShipped(...)); +``` + +::: tip Consuming dynamic subjects +Because the publish subject varies, a consumer must subscribe to the whole space with a NATS wildcard. +For Core NATS, listen on `orders.events.>`. For JetStream, provision the stream over a wildcard subject +(`orders.events.>`) so it captures every computed subject, then listen with a matching consumer filter. A +too-narrow stream subject silently fails to capture the dynamic subjects. +::: + +For subject shaping that a strongly-typed `Func` can't express — for example deriving the subject +from an envelope header or tenant id — configure an `ISubjectResolver`. It runs after the base/topic subject +is determined and can rewrite it from any envelope state: + +```csharp +opts.UseNats(nats => +{ + nats.ConnectionString = "nats://localhost:4222"; + nats.SubjectResolver = new MyAggregateSubjectResolver(); +}); +``` + +### Deduplication (JetStream `Nats-Msg-Id`) + +Wolverine stamps a `Nats-Msg-Id` on every JetStream publish, so the stream's duplicate window discards +duplicates server-side — the idempotency key external (non-Wolverine) consumers can rely on, independent of +Wolverine's own durable-inbox dedup on `Envelope.Id`. + +By default the id is the Wolverine `Envelope.Id`. Project a domain identity instead with `DeduplicateUsing` +so a logical event dedups even across separate sends (e.g. `{stream}/{version}`): + +```csharp +opts.UseNats("nats://localhost:4222") + .AutoProvision() + // Any two publishes resolving to the same key within the stream's duplicate window collapse to one. + .DeduplicateUsing(envelope => $"{envelope.GroupId}/{envelope.Id}"); +``` + +Precedence for the dedup key: + +1. An explicit `Nats-Msg-Id` header already on the outgoing envelope always wins. +2. Otherwise the configured `DeduplicateUsing` function is used. +3. Otherwise the Wolverine `Envelope.Id`. + +The duplicate window itself is configured per stream (`WithDeduplicationWindow`) or transport-wide via +`JetStreamDefaults.DuplicateWindow` (default two minutes): + +```csharp +opts.UseNats("nats://localhost:4222") + .DefineStream("ORDERS", s => s + .WithSubjects("orders.>") + .WithDeduplicationWindow(TimeSpan.FromMinutes(5))); +``` + ## Scheduled Message Delivery NATS Server 2.12+ supports native scheduled message delivery. When enabled, Wolverine uses NATS headers for scheduling instead of database persistence. @@ -380,9 +480,26 @@ When native scheduled send is not available (server < 2.12 or stream not configu ## Multi-Tenancy -NATS transport supports subject-based tenant isolation. +::: tip +For a holistic overview of multi-tenancy across all of Wolverine, see the [Multi-Tenancy Tutorial](/tutorials/multi-tenancy) +and [Multi-Tenancy with Wolverine](/guide/handlers/multi-tenancy) for how Wolverine tracks the tenant id across messages. +::: + +The NATS transport supports two flavors of tenant isolation: + +- **Subject-based** — all tenants share one connection and are separated by a tenant subject prefix + (`{tenantId}.{subject}`). This is soft partitioning within a single NATS account. +- **Connection-based** — a tenant gets its own dedicated NATS connection to a different server or **account**. -### Basic Multi-Tenancy +::: info NATS accounts are the native tenancy boundary +In NATS, true multi-tenancy is [Accounts](https://docs.nats.io/running-a-nats-service/configuration/securing_nats/accounts): +each account is a fully isolated subject namespace, and a single connection authenticates into exactly **one** +account. So a genuinely isolated tenant means a **dedicated connection with its own credentials** (see +[Per-Tenant Connections](#per-tenant-connections)). A subject prefix on a shared connection is only +partitioning within one account, not account-level isolation. +::: + +### Basic Multi-Tenancy (Subject Isolation) ```csharp opts.UseNats("nats://localhost:4222") @@ -396,6 +513,33 @@ opts.UseNats("nats://localhost:4222") - `TenantIdRequired`: Throws if tenant ID is missing - `FallbackToDefault`: Uses base subject if tenant ID is missing +### Per-Tenant Connections + +To route a tenant to its own NATS server or account, add it with a configuration action. The action receives +a copy of the transport's own connection settings, so you only override what differs for this tenant — a +different URL, or any of the NATS auth mechanisms (token, JWT/NKey, credentials file, client certificate): + +```csharp +opts.UseNats("nats://shared:4222") + .ConfigureMultiTenancy(TenantedIdBehavior.FallbackToDefault) + .AddTenant("tenant-a", cfg => cfg.ConnectionString = "nats://tenant-a-host:4222") + .AddTenant("tenant-b", cfg => + { + cfg.ConnectionString = "nats://tenant-b-host:4222"; + cfg.CredentialsFile = "/etc/nats/tenant-b.creds"; + }); +``` + +Each tenant with its own configuration gets a dedicated connection, owned by the transport for its lifetime. +Tenants added without a configuration action keep sharing the transport connection (subject-prefix isolation +only). + +Both **sending and listening** are tenant-aware. A listener consumes on the shared connection *and* on each +tenant's dedicated connection: when a message arrives on a tenant connection it is stamped with that tenant's +id, and its ack/nak/dead-letter is routed back over the same connection. Sending a message tagged with a +`TenantId` publishes it over that tenant's connection. If any tenant streams need JetStream, the configured +streams are auto-provisioned on each tenant server as well (when `AutoProvision()` is on). + ### Custom Subject Mapper ```csharp @@ -430,8 +574,13 @@ The response endpoint always uses Core NATS for low-latency replies, even when t ### JetStream -- **Retry**: Message is requeued via `NakAsync()` with optional delay -- **Dead Letter**: Message is terminated via `AckTerminateAsync()` +- **Retry**: Message is requeued via `NakAsync()` with optional delay, up to the consumer's maximum delivery + attempts (`JetStreamDefaults.MaxDeliver`, default 5, or a per-endpoint `MaxDeliveryAttempts` override). +- **Dead Letter**: Once delivery attempts are exhausted, the poison message is first forwarded to the + configured dead-letter subject (so a terminate failure can't lose it), then terminated on the consumer via + `AckTerminateAsync(reason)` so the server stops redelivering and records why. If **no** dead-letter subject + is configured, Wolverine logs a warning and the message is terminated without being retained — configure a + dead-letter subject to keep poison messages. ### Core NATS diff --git a/src/Transports/NATS/Wolverine.Nats.Tests/Helpers/NatsTestHelpers.cs b/src/Transports/NATS/Wolverine.Nats.Tests/Helpers/NatsTestHelpers.cs new file mode 100644 index 000000000..7a99ee739 --- /dev/null +++ b/src/Transports/NATS/Wolverine.Nats.Tests/Helpers/NatsTestHelpers.cs @@ -0,0 +1,106 @@ +using JasperFx.Core; +using Microsoft.Extensions.Hosting; +using NATS.Client.Core; + +namespace Wolverine.Nats.Tests; + +/// +/// Shared helpers for the NATS integration tests. The whole NATS suite resolves the broker from the +/// NATS_URL environment variable (set by when Testcontainers is used), +/// falling back to the docker-compose broker on localhost:4222. +/// +internal static class NatsTestHelpers +{ + public static string ResolveUrl() + { + return Environment.GetEnvironmentVariable("NATS_URL") ?? "nats://localhost:4222"; + } + + /// + /// Probe the broker by spinning up a throwaway Wolverine host. Returns false (so the caller can skip) + /// when no NATS server is reachable, mirroring the guard used across the existing NATS integration tests. + /// + public static async Task IsNatsAvailable(string natsUrl) + { + try + { + using var testHost = await Host.CreateDefaultBuilder() + .UseWolverine(opts => opts.UseNats(natsUrl)) + .StartAsync(); + + await testHost.StopAsync(); + return true; + } + catch + { + return false; + } + } + + /// + /// Subscribe a raw NATS client to an exact subject. SubscribeCoreAsync registers the subscription + /// synchronously (the SUB is written before it returns) and the follow-up ping flushes it to the server, + /// so there is no subscribe-before-publish race. Used to prove the exact concrete subject a message was + /// published to — something a wildcard Wolverine listener can't (the receiving pipeline overwrites + /// Envelope.Destination with the listener's own address). + /// + public static async Task SubscribeRawAsync(string url, string subject) + { + var connection = new NatsConnection(new NatsOpts { Url = url }); + await connection.ConnectAsync(); + var sub = await connection.SubscribeCoreAsync(subject); + await connection.PingAsync(); + return new RawSubscription(connection, sub); + } +} + +/// +/// A raw NATS subscription plus its dedicated connection, disposed together. +/// +internal sealed class RawSubscription : IAsyncDisposable +{ + private readonly NatsConnection _connection; + private readonly INatsSub _sub; + + public RawSubscription(NatsConnection connection, INatsSub sub) + { + _connection = connection; + _sub = sub; + } + + /// + /// Read the next message, returning null if none arrives within . + /// + public async Task?> ReadAsync(TimeSpan timeout) + { + using var cts = new CancellationTokenSource(timeout); + try + { + return await _sub.Msgs.ReadAsync(cts.Token); + } + catch (OperationCanceledException) + { + return null; + } + } + + public async ValueTask DisposeAsync() + { + await _sub.DisposeAsync(); + await _connection.DisposeAsync(); + } +} + +/// +/// Message used by the dynamic-subject and dedup tests. The doubles as a stable +/// domain identity for JetStream deduplication (see DeduplicateUsing). +/// +public record OrderPlaced(string OrderId); + +public class OrderPlacedHandler +{ + // No-op: receivers only need this so Wolverine tracking counts the message as "received". + public void Handle(OrderPlaced message) + { + } +} diff --git a/src/Transports/NATS/Wolverine.Nats.Tests/MultiTenancyTests.cs b/src/Transports/NATS/Wolverine.Nats.Tests/MultiTenancyTests.cs index 2fc838972..24df11b9e 100644 --- a/src/Transports/NATS/Wolverine.Nats.Tests/MultiTenancyTests.cs +++ b/src/Transports/NATS/Wolverine.Nats.Tests/MultiTenancyTests.cs @@ -260,10 +260,11 @@ public class NatsTenantTests public void creates_tenant_with_id() { var tenant = new NatsTenant("tenant1"); - + tenant.TenantId.ShouldBe("tenant1"); tenant.SubjectMapper.ShouldBeNull(); - tenant.ConnectionString.ShouldBeNull(); + tenant.ConnectionConfiguration.ShouldBeNull(); + tenant.HasOwnConnection.ShouldBeFalse(); } [Fact] diff --git a/src/Transports/NATS/Wolverine.Nats.Tests/NatsDynamicSubjectTenancyTests.cs b/src/Transports/NATS/Wolverine.Nats.Tests/NatsDynamicSubjectTenancyTests.cs new file mode 100644 index 000000000..8a599bd2c --- /dev/null +++ b/src/Transports/NATS/Wolverine.Nats.Tests/NatsDynamicSubjectTenancyTests.cs @@ -0,0 +1,76 @@ +using IntegrationTests; +using JasperFx.Core; +using Microsoft.Extensions.Hosting; +using Shouldly; +using Wolverine.Transports.Sending; +using Xunit; +using Xunit.Abstractions; + +namespace Wolverine.Nats.Tests; + +/// +/// Regression coverage for the intersection of the two features this branch adds: per-message dynamic subjects +/// (RoutingMode.ByTopic / Envelope.TopicName) and subject-isolation multi-tenancy. +/// +/// A static endpoint subject is tenant-qualified once, at sender construction (NatsEndpoint.CreateSender). +/// A subject computed per message is not — so without the fix a subject-isolation tenant's dynamic-subject send +/// would publish to the raw, un-prefixed subject on the shared connection, silently defeating isolation. This +/// asserts the computed subject is tenant-qualified ({tenantId}.{computed}) and does not leak to the other +/// tenant's subject or the bare (un-prefixed) subject. +/// +/// Uses raw NATS subscribers on the exact expected subjects (a single shared broker), so it is a local sanity +/// check and is not required to run in CI. +/// +[Collection("NATS Integration")] +[Trait("Category", "Integration")] +public class NatsDynamicSubjectTenancyTests +{ + private readonly ITestOutputHelper _output; + + public NatsDynamicSubjectTenancyTests(ITestOutputHelper output) => _output = output; + + [Fact] + public async Task dynamic_subject_is_tenant_qualified_for_subject_isolation_tenants() + { + var url = NatsTestHelpers.ResolveUrl(); + if (!await NatsTestHelpers.IsNatsAvailable(url)) + { + return; + } + + var root = $"orders.events.{Guid.NewGuid():N}"; + var orderId = Guid.NewGuid().ToString("N"); + + using var host = await Host.CreateDefaultBuilder() + .ConfigureLogging(l => l.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "DynamicSubjectTenancy"; + opts.UseNats(url) + .ConfigureMultiTenancy(TenantedIdBehavior.FallbackToDefault) + // Both tenants share the connection; isolation is purely by subject prefix. + .AddTenant("tenantA") + .AddTenant("tenantB"); + + opts.Policies.DisableConventionalLocalRouting(); + opts.PublishMessagesToNatsSubject(m => $"{root}.{m.OrderId}").SendInline(); + }) + .StartAsync(); + + // Each tenant's dynamic subject must be tenant-qualified: {tenantId}.{computed subject}. + await using var subA = await NatsTestHelpers.SubscribeRawAsync(url, $"tenantA.{root}.{orderId}"); + await using var subB = await NatsTestHelpers.SubscribeRawAsync(url, $"tenantB.{root}.{orderId}"); + // The un-prefixed subject proves nothing escapes isolation (this is where the pre-fix bug landed). + await using var subBare = await NatsTestHelpers.SubscribeRawAsync(url, $"{root}.{orderId}"); + + await host.MessageBus().SendAsync(new OrderPlaced(orderId), new DeliveryOptions { TenantId = "tenantA" }); + + var onA = await subA.ReadAsync(15.Seconds()); + onA.ShouldNotBeNull(); + onA!.Value.Subject.ShouldBe($"tenantA.{root}.{orderId}"); + + // Not on the other tenant's subject, and not on the bare un-prefixed subject. + (await subB.ReadAsync(2.Seconds())).ShouldBeNull(); + (await subBare.ReadAsync(2.Seconds())).ShouldBeNull(); + } +} diff --git a/src/Transports/NATS/Wolverine.Nats.Tests/NatsDynamicSubjectTests.cs b/src/Transports/NATS/Wolverine.Nats.Tests/NatsDynamicSubjectTests.cs new file mode 100644 index 000000000..81f4052ee --- /dev/null +++ b/src/Transports/NATS/Wolverine.Nats.Tests/NatsDynamicSubjectTests.cs @@ -0,0 +1,327 @@ +using IntegrationTests; +using JasperFx.Core; +using Microsoft.Extensions.Hosting; +using Shouldly; +using Wolverine.Nats.Configuration; +using Wolverine.Tracking; +using Xunit; +using Xunit.Abstractions; + +namespace Wolverine.Nats.Tests; + +/// +/// Integration coverage for per-message dynamic NATS subjects. Two mechanisms are exercised, both built on +/// Wolverine's generic topic routing (RoutingMode.ByTopic / Envelope.TopicName): +/// +/// the strongly-typed PublishMessagesToNatsSubject<T>(Func<T,string>) registration, and +/// IMessageBus.BroadcastToTopicAsync, which the same ByTopic endpoint auto-enrolls in. +/// +/// A third test covers the advanced escape hatch that rewrites the subject from +/// envelope-level state a strongly-typed function can't reach. +/// +/// Each test proves two things: a {root}.> wildcard Wolverine listener consumes the message +/// end-to-end (tracking), and a raw NATS subscriber bound to the exact expected subject receives it — +/// the raw subscriber is what pins down the concrete subject, since the receiving pipeline overwrites +/// Envelope.Destination with the listener's own (wildcard) address. +/// +[Collection("NATS Integration")] +[Trait("Category", "Integration")] +public class NatsDynamicSubjectTests : IAsyncLifetime +{ + private readonly ITestOutputHelper _output; + private IHost? _sender; + private IHost? _receiver; + private string _root = null!; + private string _natsUrl = null!; + + public NatsDynamicSubjectTests(ITestOutputHelper output) => _output = output; + + public async Task InitializeAsync() + { + _natsUrl = NatsTestHelpers.ResolveUrl(); + _root = $"orders.events.{Guid.NewGuid():N}"; + + if (!await NatsTestHelpers.IsNatsAvailable(_natsUrl)) + { + _output.WriteLine("NATS not available, skipping test"); + return; + } + + _receiver = await Host.CreateDefaultBuilder() + .ConfigureLogging(logging => logging.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "DynamicSubjectReceiver"; + opts.UseNats(_natsUrl).AutoProvision(); + + // One wildcard Core NATS listener captures every per-message subject under the root. + opts.ListenToNatsSubject($"{_root}.>").Named("wildcard"); + }) + .StartAsync(); + + _sender = await Host.CreateDefaultBuilder() + .ConfigureLogging(logging => logging.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "DynamicSubjectSender"; + opts.UseNats(_natsUrl).AutoProvision(); + + // The sender also discovers OrderPlacedHandler (same assembly); without this the additive + // TopicRouting source AND the local handler would both fire, double-counting in tracking. + opts.Policies.DisableConventionalLocalRouting(); + + // Per-message subject computed from the message body. + opts.PublishMessagesToNatsSubject(m => $"{_root}.{m.OrderId}").SendInline(); + }) + .StartAsync(); + } + + public async Task DisposeAsync() + { + if (_sender != null) + { + await _sender.StopAsync(); + _sender.Dispose(); + } + + if (_receiver != null) + { + await _receiver.StopAsync(); + _receiver.Dispose(); + } + } + + [Fact] + public async Task publishes_to_a_per_message_subject_computed_from_the_body() + { + if (_sender == null || _receiver == null) + { + return; + } + + var orderId = Guid.NewGuid().ToString("N"); + var expectedSubject = $"{_root}.{orderId}"; + + await using var raw = await NatsTestHelpers.SubscribeRawAsync(_natsUrl, expectedSubject); + + var session = await _sender + .TrackActivity() + .AlsoTrack(_receiver) + .Timeout(30.Seconds()) + .SendMessageAndWaitAsync(new OrderPlaced(orderId)); + + // The wildcard Wolverine listener consumed it end-to-end. + session.Received.SingleMessage().OrderId.ShouldBe(orderId); + + // ...and it landed on exactly the computed subject. + var delivered = await raw.ReadAsync(5.Seconds()); + delivered.ShouldNotBeNull(); + delivered!.Value.Subject.ShouldBe(expectedSubject); + } + + [Fact] + public async Task two_messages_land_on_two_distinct_computed_subjects() + { + if (_sender == null || _receiver == null) + { + return; + } + + var firstId = Guid.NewGuid().ToString("N"); + var secondId = Guid.NewGuid().ToString("N"); + + await using var rawFirst = await NatsTestHelpers.SubscribeRawAsync(_natsUrl, $"{_root}.{firstId}"); + await using var rawSecond = await NatsTestHelpers.SubscribeRawAsync(_natsUrl, $"{_root}.{secondId}"); + + var session = await _sender + .TrackActivity() + .AlsoTrack(_receiver) + .Timeout(30.Seconds()) + .ExecuteAndWaitAsync((Func)(async context => + { + await context.SendAsync(new OrderPlaced(firstId)); + await context.SendAsync(new OrderPlaced(secondId)); + })); + + session.Received.MessagesOf().Count().ShouldBe(2); + + var first = await rawFirst.ReadAsync(5.Seconds()); + first.ShouldNotBeNull(); + first!.Value.Subject.ShouldBe($"{_root}.{firstId}"); + + var second = await rawSecond.ReadAsync(5.Seconds()); + second.ShouldNotBeNull(); + second!.Value.Subject.ShouldBe($"{_root}.{secondId}"); + } + + [Fact] + public async Task broadcast_to_topic_async_publishes_to_the_explicit_subject() + { + if (_sender == null || _receiver == null) + { + return; + } + + // BroadcastToTopicAsync routes to every ByTopic endpoint (the one created by + // PublishMessagesToNatsSubject above) with an explicit topic that overrides the Func. + var subject = $"{_root}.broadcast"; + + await using var raw = await NatsTestHelpers.SubscribeRawAsync(_natsUrl, subject); + + var session = await _sender + .TrackActivity() + .AlsoTrack(_receiver) + .Timeout(30.Seconds()) + .ExecuteAndWaitAsync(context => + context.BroadcastToTopicAsync(subject, new OrderPlaced("broadcast"))); + + session.Received.SingleMessage().OrderId.ShouldBe("broadcast"); + + var delivered = await raw.ReadAsync(5.Seconds()); + delivered.ShouldNotBeNull(); + delivered!.Value.Subject.ShouldBe(subject); + } + + [Fact] + public async Task subject_resolver_escape_hatch_rewrites_the_subject_from_envelope_state() + { + if (!await NatsTestHelpers.IsNatsAvailable(_natsUrl)) + { + return; + } + + var root = $"aggregates.{Guid.NewGuid():N}"; + var baseSubject = $"{root}.base"; + var expectedSubject = $"{baseSubject}.A-42"; + + using var receiver = await Host.CreateDefaultBuilder() + .ConfigureLogging(logging => logging.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "ResolverReceiver"; + opts.UseNats(_natsUrl).AutoProvision(); + opts.ListenToNatsSubject($"{root}.>").Named("resolver-wildcard"); + }) + .StartAsync(); + + using var sender = await Host.CreateDefaultBuilder() + .ConfigureLogging(logging => logging.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "ResolverSender"; + opts.UseNats(natsCfg => + { + natsCfg.ConnectionString = _natsUrl; + // Advanced hook: rewrite the outgoing subject from an envelope header the typed + // Func can't see. + natsCfg.SubjectResolver = new AggregateSubjectResolver(); + }).AutoProvision(); + + opts.Policies.DisableConventionalLocalRouting(); + opts.PublishMessage().ToNatsSubject(baseSubject).SendInline(); + }) + .StartAsync(); + + await using var raw = await NatsTestHelpers.SubscribeRawAsync(_natsUrl, expectedSubject); + + var session = await sender + .TrackActivity() + .AlsoTrack(receiver) + .Timeout(30.Seconds()) + .SendMessageAndWaitAsync(new OrderPlaced("A-42"), + new DeliveryOptions { Headers = { ["aggregate-id"] = "A-42" } }); + + session.Received.SingleMessage().OrderId.ShouldBe("A-42"); + + var delivered = await raw.ReadAsync(5.Seconds()); + delivered.ShouldNotBeNull(); + delivered!.Value.Subject.ShouldBe(expectedSubject); + } + + [Fact] + public async Task subject_resolver_output_is_normalized() + { + if (!await NatsTestHelpers.IsNatsAvailable(_natsUrl)) + { + return; + } + + var root = $"normalized.{Guid.NewGuid():N}"; + var baseSubject = $"{root}.base"; + // The resolver deliberately returns a '/'-separated subject. With NormalizeSubjects on (the default) + // it must be published with '.' tokens, consistent with static subjects and TopicName routing — so + // the raw subscriber bound to the normalized subject is what receives it. + var expectedSubject = $"{baseSubject}.A-42"; + + using var receiver = await Host.CreateDefaultBuilder() + .ConfigureLogging(logging => logging.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "NormalizeResolverReceiver"; + opts.UseNats(_natsUrl).AutoProvision(); + opts.ListenToNatsSubject($"{root}.>").Named("normalize-wildcard"); + }) + .StartAsync(); + + using var sender = await Host.CreateDefaultBuilder() + .ConfigureLogging(logging => logging.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "NormalizeResolverSender"; + opts.UseNats(natsCfg => + { + natsCfg.ConnectionString = _natsUrl; + natsCfg.SubjectResolver = new SlashSubjectResolver(); + }).AutoProvision(); + + opts.Policies.DisableConventionalLocalRouting(); + opts.PublishMessage().ToNatsSubject(baseSubject).SendInline(); + }) + .StartAsync(); + + await using var raw = await NatsTestHelpers.SubscribeRawAsync(_natsUrl, expectedSubject); + + await sender + .TrackActivity() + .AlsoTrack(receiver) + .Timeout(30.Seconds()) + .SendMessageAndWaitAsync(new OrderPlaced("A-42"), + new DeliveryOptions { Headers = { ["aggregate-id"] = "A-42" } }); + + var delivered = await raw.ReadAsync(5.Seconds()); + delivered.ShouldNotBeNull(); + delivered!.Value.Subject.ShouldBe(expectedSubject); + } + + /// + /// Rewrites the base subject to {base}.{aggregate-id} using the aggregate-id header — + /// the kind of envelope-level shaping the strongly-typed subject function can't express. + /// + private sealed class AggregateSubjectResolver : ISubjectResolver + { + public string ResolveSubject(string baseSubject, Envelope envelope) + { + return envelope.Headers.TryGetValue("aggregate-id", out var id) && id.IsNotEmpty() + ? $"{baseSubject}.{id}" + : baseSubject; + } + + public string? ExtractTenantId(string subject) => null; + } + + /// + /// Returns a '/'-separated subject on purpose, to prove the transport normalizes a resolver's output the + /// same way it normalizes static subjects and Envelope.TopicName (when NormalizeSubjects is on). + /// + private sealed class SlashSubjectResolver : ISubjectResolver + { + public string ResolveSubject(string baseSubject, Envelope envelope) + { + return envelope.Headers.TryGetValue("aggregate-id", out var id) && id.IsNotEmpty() + ? $"{baseSubject}/{id}" + : baseSubject; + } + + public string? ExtractTenantId(string subject) => null; + } +} diff --git a/src/Transports/NATS/Wolverine.Nats.Tests/NatsJetStreamDedupTests.cs b/src/Transports/NATS/Wolverine.Nats.Tests/NatsJetStreamDedupTests.cs new file mode 100644 index 000000000..15d95c023 --- /dev/null +++ b/src/Transports/NATS/Wolverine.Nats.Tests/NatsJetStreamDedupTests.cs @@ -0,0 +1,156 @@ +using IntegrationTests; +using JasperFx.Core; +using Microsoft.Extensions.Hosting; +using NATS.Client.Core; +using NATS.Client.JetStream; +using NATS.Net; +using Shouldly; +using Xunit; +using Xunit.Abstractions; + +namespace Wolverine.Nats.Tests; + +/// +/// Integration coverage for server-side JetStream deduplication. Wolverine now stamps a +/// Nats-Msg-Id on every JetStream publish (default: the envelope Id, overridable via +/// DeduplicateUsing or an explicit Nats-Msg-Id header), so the stream's duplicate +/// window actually discards duplicates — the idempotency guarantee external (non-Wolverine) +/// consumers rely on. +/// +/// Assertions are deterministic: each publish goes through an inline JetStream sender that awaits the +/// publish-ack, so by the time the send returns the message is either persisted or discarded, and the +/// stream's message count is authoritative. +/// +[Collection("NATS Integration")] +[Trait("Category", "Integration")] +public class NatsJetStreamDedupTests +{ + // The NATS JetStream server-side dedup header (mirrors JetStreamPublisher.NatsMsgIdHeader, which is internal). + private const string NatsMsgIdHeader = "Nats-Msg-Id"; + + private readonly ITestOutputHelper _output; + + public NatsJetStreamDedupTests(ITestOutputHelper output) => _output = output; + + [Fact] + public async Task same_domain_msg_id_is_deduplicated_by_the_stream() + { + var natsUrl = NatsTestHelpers.ResolveUrl(); + if (!await NatsTestHelpers.IsNatsAvailable(natsUrl)) return; + + var stream = $"DEDUP_{Guid.NewGuid():N}"; + var subject = $"dedup.domain.{Guid.NewGuid():N}"; + + using var host = await Host.CreateDefaultBuilder() + .ConfigureLogging(l => l.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "DedupDomain"; + opts.UseNats(natsUrl) + .AutoProvision() + // Project a stable domain identity into the dedup key instead of the per-send envelope Id. + .DeduplicateUsing(e => ((OrderPlaced)e.Message!).OrderId) + .DefineStream(stream, s => s + .WithSubjects(subject) + .WithDeduplicationWindow(5.Minutes())); + + opts.Policies.DisableConventionalLocalRouting(); + opts.PublishMessage().ToNatsSubject(subject).UseJetStream(stream).SendInline(); + }) + .StartAsync(); + + var bus = host.MessageBus(); + var orderId = Guid.NewGuid().ToString("N"); + + await bus.SendAsync(new OrderPlaced(orderId)); // persisted + await bus.SendAsync(new OrderPlaced(orderId)); // same key -> discarded + await bus.SendAsync(new OrderPlaced(Guid.NewGuid().ToString("N"))); // different key -> persisted + + (await CountStreamMessagesAsync(natsUrl, stream)).ShouldBe(2); + } + + [Fact] + public async Task explicit_nats_msg_id_header_is_honored_for_dedup() + { + var natsUrl = NatsTestHelpers.ResolveUrl(); + if (!await NatsTestHelpers.IsNatsAvailable(natsUrl)) return; + + var stream = $"DEDUP_{Guid.NewGuid():N}"; + var subject = $"dedup.header.{Guid.NewGuid():N}"; + + using var host = await Host.CreateDefaultBuilder() + .ConfigureLogging(l => l.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "DedupHeader"; + opts.UseNats(natsUrl) + .AutoProvision() + .DefineStream(stream, s => s + .WithSubjects(subject) + .WithDeduplicationWindow(5.Minutes())); + + opts.Policies.DisableConventionalLocalRouting(); + opts.PublishMessage().ToNatsSubject(subject).UseJetStream(stream).SendInline(); + }) + .StartAsync(); + + var bus = host.MessageBus(); + var fixedMsgId = Guid.NewGuid().ToString("N"); + + // Two logically-different messages (distinct envelope Ids) but the same explicit Nats-Msg-Id header; + // the explicit header must win over the default and be honored by the server for dedup. + await bus.SendAsync(new OrderPlaced("a"), + new DeliveryOptions { Headers = { [NatsMsgIdHeader] = fixedMsgId } }); + await bus.SendAsync(new OrderPlaced("b"), + new DeliveryOptions { Headers = { [NatsMsgIdHeader] = fixedMsgId } }); + + (await CountStreamMessagesAsync(natsUrl, stream)).ShouldBe(1); + } + + [Fact] + public async Task distinct_messages_are_not_deduplicated_with_default_envelope_id() + { + var natsUrl = NatsTestHelpers.ResolveUrl(); + if (!await NatsTestHelpers.IsNatsAvailable(natsUrl)) return; + + var stream = $"DEDUP_{Guid.NewGuid():N}"; + var subject = $"dedup.default.{Guid.NewGuid():N}"; + + using var host = await Host.CreateDefaultBuilder() + .ConfigureLogging(l => l.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "DedupDefault"; + opts.UseNats(natsUrl) + .AutoProvision() + .DefineStream(stream, s => s + .WithSubjects(subject) + .WithDeduplicationWindow(5.Minutes())); + + opts.Policies.DisableConventionalLocalRouting(); + opts.PublishMessage().ToNatsSubject(subject).UseJetStream(stream).SendInline(); + }) + .StartAsync(); + + var bus = host.MessageBus(); + + // Default Nats-Msg-Id is the (unique) envelope Id, so two distinct sends must both persist — + // guards against over-eager dedup collapsing unrelated messages. + await bus.SendAsync(new OrderPlaced("a")); + await bus.SendAsync(new OrderPlaced("b")); + + (await CountStreamMessagesAsync(natsUrl, stream)).ShouldBe(2); + } + + // Query the stream state over an independent connection (no dependency on Wolverine internals): the + // publish-acks have already been awaited, so the count is authoritative. + private static async Task CountStreamMessagesAsync(string natsUrl, string streamName) + { + await using var connection = new NatsConnection(new NatsOpts { Url = natsUrl }); + await connection.ConnectAsync(); + + var js = connection.CreateJetStreamContext(); + var stream = await js.GetStreamAsync(streamName); + return stream.Info.State.Messages; + } +} diff --git a/src/Transports/NATS/Wolverine.Nats.Tests/NatsPerTenantConnectionTests.cs b/src/Transports/NATS/Wolverine.Nats.Tests/NatsPerTenantConnectionTests.cs new file mode 100644 index 000000000..18a2da0cb --- /dev/null +++ b/src/Transports/NATS/Wolverine.Nats.Tests/NatsPerTenantConnectionTests.cs @@ -0,0 +1,210 @@ +using IntegrationTests; +using JasperFx.Core; +using Microsoft.Extensions.Hosting; +using NATS.Client.Core; +using NATS.Client.JetStream; +using NATS.Net; +using Shouldly; +using Testcontainers.Nats; +using Wolverine.Tracking; +using Wolverine.Transports.Sending; +using Xunit; +using Xunit.Abstractions; + +namespace Wolverine.Nats.Tests; + +/// +/// Integration coverage for per-tenant NATS connections (as opposed to per-tenant subject +/// prefixing on one shared connection, which covers). +/// +/// Answering the practical question "does per-tenant mean a different NATS setup?": yes — to prove a tenant +/// uses its own connection, that connection must point at a genuinely distinguishable server. This +/// test therefore spins up a second NATS broker (server B) via Testcontainers alongside the shared broker +/// (server A) and asserts: +/// +/// a message for the tenant with a dedicated connection is published to server B and not server A, +/// a default (no-tenant) message is published to server A (the shared connection), and +/// a message tagged for the tenant is consumed back over the tenant's own connection (server B), +/// arriving stamped with its tenant id. +/// +/// +/// The publish-side assertions use raw NATS subscribers rather than Wolverine receivers so the exact target +/// server is provable; the inbound test uses a single Wolverine host plus the tracking API. This test needs +/// two brokers, so it is intended as a local sanity check and is not required to run in CI. +/// +[Collection("NATS Integration")] +[Trait("Category", "Integration")] +public class NatsPerTenantConnectionTests : IAsyncLifetime +{ + private readonly ITestOutputHelper _output; + private NatsContainer? _serverB; + private string _serverAUrl = null!; + private string _serverBUrl = null!; + private bool _skip; + + public NatsPerTenantConnectionTests(ITestOutputHelper output) => _output = output; + + public async Task InitializeAsync() + { + _serverAUrl = NatsTestHelpers.ResolveUrl(); + + if (!await NatsTestHelpers.IsNatsAvailable(_serverAUrl)) + { + _skip = true; + return; + } + + // Server B is a second, independent broker so "used the tenant's own connection" is provable: the + // message can only appear on B if the dedicated connection carried it there. + _serverB = new NatsBuilder().WithImage("nats:latest").Build(); + await _serverB.StartAsync(); + _serverBUrl = _serverB.GetConnectionString(); + + _output.WriteLine($"Server A (shared): {_serverAUrl}"); + _output.WriteLine($"Server B (tenant): {_serverBUrl}"); + } + + public async Task DisposeAsync() + { + if (_serverB != null) + { + await _serverB.DisposeAsync(); + } + } + + [Fact] + public async Task tenant_message_is_published_over_the_tenants_own_connection() + { + if (_skip) return; + + var baseSubject = $"pertenant.{Guid.NewGuid():N}"; + var tenantSubject = $"tenantB.{baseSubject}"; // DefaultTenantSubjectMapper prefixes the tenant id + + await using var subOnB = await NatsTestHelpers.SubscribeRawAsync(_serverBUrl, tenantSubject); + await using var subOnA = await NatsTestHelpers.SubscribeRawAsync(_serverAUrl, tenantSubject); + + using var host = await BuildSenderAsync(baseSubject); + + await host.MessageBus().SendAsync(new OrderPlaced("for-tenant-b"), + new DeliveryOptions { TenantId = "tenantB" }); + + // Landed on server B (the tenant's dedicated connection)... + var received = await subOnB.ReadAsync(15.Seconds()); + received.ShouldNotBeNull(); + received!.Value.Subject.ShouldBe(tenantSubject); + + // ...and NOT on the shared server A. + (await subOnA.ReadAsync(2.Seconds())).ShouldBeNull(); + } + + [Fact] + public async Task default_message_uses_the_shared_connection() + { + if (_skip) return; + + var baseSubject = $"pertenant.{Guid.NewGuid():N}"; + + // No tenant prefix for the fallback/default path — it publishes to the base subject on server A. + await using var subOnA = await NatsTestHelpers.SubscribeRawAsync(_serverAUrl, baseSubject); + await using var subOnB = await NatsTestHelpers.SubscribeRawAsync(_serverBUrl, baseSubject); + + using var host = await BuildSenderAsync(baseSubject); + + await host.MessageBus().SendAsync(new OrderPlaced("no-tenant")); + + var received = await subOnA.ReadAsync(15.Seconds()); + received.ShouldNotBeNull(); + received!.Value.Subject.ShouldBe(baseSubject); + + (await subOnB.ReadAsync(2.Seconds())).ShouldBeNull(); + } + + [Fact] + public async Task tenant_message_is_consumed_over_the_tenants_own_connection() + { + if (_skip) return; + + var baseSubject = $"pertenant.{Guid.NewGuid():N}"; + + // A single host both publishes and listens. The tenant listener runs on server B (the dedicated + // connection), so a message tagged for tenantB round-trips back through it, stamped with the tenant id. + using var host = await Host.CreateDefaultBuilder() + .ConfigureLogging(l => l.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "PerTenantInbound"; + opts.UseNats(_serverAUrl) + .ConfigureMultiTenancy(TenantedIdBehavior.FallbackToDefault) + .AddTenant("tenantB", cfg => cfg.ConnectionString = _serverBUrl); + + opts.PublishMessage().ToNatsSubject(baseSubject).SendInline(); + opts.ListenToNatsSubject(baseSubject); + }) + .StartAsync(); + + // Single host both sends and receives via the broker, so explicitly wait for the round-trip receipt + // rather than just the send settling. + var session = await host + .TrackActivity() + .Timeout(30.Seconds()) + .WaitForMessageToBeReceivedAt(host) + .ExecuteAndWaitAsync(c => + c.SendAsync(new OrderPlaced("for-tenant-b"), new DeliveryOptions { TenantId = "tenantB" })); + + var received = session.Received.SingleEnvelope(); + received.TenantId.ShouldBe("tenantB"); + received.Message.ShouldBeOfType().OrderId.ShouldBe("for-tenant-b"); + } + + [Fact] + public async Task streams_are_auto_provisioned_over_a_tenants_own_connection() + { + if (_skip) return; + + var streamName = $"TENANTPROV_{Guid.NewGuid():N}"; + var subject = $"tenantprov.{Guid.NewGuid():N}"; + + // A tenant with its own connection to server B, plus a configured stream + AutoProvision. Stream + // provisioning runs at host start over the tenant's dedicated connection — which is never explicitly + // ConnectAsync'd; the NATS client connects lazily on first use. If that path were broken, StartAsync + // would throw right here. + using var host = await Host.CreateDefaultBuilder() + .ConfigureLogging(l => l.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "PerTenantProvisioning"; + opts.UseNats(_serverAUrl) + .ConfigureMultiTenancy(TenantedIdBehavior.FallbackToDefault) + .AutoProvision() + .DefineStream(streamName, s => s.WithSubjects($"{subject}.>")) + .AddTenant("tenantB", cfg => cfg.ConnectionString = _serverBUrl); + }) + .StartAsync(); + + // Prove the stream was actually created on server B (the tenant's own server), not just server A. + // GetStreamAsync throws if the stream is absent, so a broken provisioning path fails the test. + await using var connToB = new NatsConnection(new NatsOpts { Url = _serverBUrl }); + await connToB.ConnectAsync(); + var streamOnB = await connToB.CreateJetStreamContext().GetStreamAsync(streamName); + streamOnB.ShouldNotBeNull(); + } + + private Task BuildSenderAsync(string baseSubject) + { + return Host.CreateDefaultBuilder() + .ConfigureLogging(l => l.AddXunitLogging(_output)) + .UseWolverine(opts => + { + opts.ServiceName = "PerTenantSender"; + opts.UseNats(_serverAUrl) + .ConfigureMultiTenancy(TenantedIdBehavior.FallbackToDefault) + // The tenant gets its own connection to server B; the action is seeded from the parent + // settings so we only override the URL that differs. + .AddTenant("tenantB", cfg => cfg.ConnectionString = _serverBUrl); + + opts.Policies.DisableConventionalLocalRouting(); + opts.PublishMessage().ToNatsSubject(baseSubject).SendInline(); + }) + .StartAsync(); + } +} diff --git a/src/Transports/NATS/Wolverine.Nats/Configuration/NatsSubscriberConfiguration.cs b/src/Transports/NATS/Wolverine.Nats/Configuration/NatsSubscriberConfiguration.cs index d27268a3e..4e6526664 100644 --- a/src/Transports/NATS/Wolverine.Nats/Configuration/NatsSubscriberConfiguration.cs +++ b/src/Transports/NATS/Wolverine.Nats/Configuration/NatsSubscriberConfiguration.cs @@ -37,4 +37,15 @@ public NatsSubscriberConfiguration UseScheduleSubjectSuffix(string suffix) add(endpoint => endpoint.ScheduleSubjectSuffix = suffix); return this; } + + /// + /// Add a static header written to every message published to this subject. + /// + public NatsSubscriberConfiguration AddOutgoingHeader(string key, string value) + { + ArgumentException.ThrowIfNullOrWhiteSpace(key, nameof(key)); + + add(endpoint => endpoint.CustomHeaders[key] = value); + return this; + } } diff --git a/src/Transports/NATS/Wolverine.Nats/Configuration/NatsTransportConfiguration.cs b/src/Transports/NATS/Wolverine.Nats/Configuration/NatsTransportConfiguration.cs index 0c41d9202..37533e3d5 100644 --- a/src/Transports/NATS/Wolverine.Nats/Configuration/NatsTransportConfiguration.cs +++ b/src/Transports/NATS/Wolverine.Nats/Configuration/NatsTransportConfiguration.cs @@ -52,6 +52,16 @@ public class NatsTransportConfiguration public ITenantIdResolver? TenantIdResolver { get; set; } public ISubjectResolver? SubjectResolver { get; set; } public string? TenantSubjectPrefix { get; set; } + + /// + /// Optional source of the JetStream Nats-Msg-Id used for server-side + /// deduplication (within the stream's duplicate window). Defaults to the Wolverine + /// envelope Id when null; set it to project a domain identity such as + /// {stream}/{version} so non-Wolverine consumers get server-side dedup too. + /// An explicit Nats-Msg-Id header already on the outgoing envelope wins. + /// + [IgnoreDescription] + public Func? MsgIdSource { get; set; } public Dictionary Streams { get; set; } = new(); internal NatsOpts ToNatsOpts() @@ -84,27 +94,32 @@ internal NatsOpts ToNatsOpts() }; } - internal NatsJSOpts? ToJetStreamOpts() - { - if (!EnableJetStream) - { - return null; - } - - return new NatsJSOpts(ToNatsOpts(), JetStreamDomain, JetStreamApiPrefix ?? "$JS.API"); - } } +/// +/// Transport-wide defaults used as the template when Wolverine auto-provisions JetStream streams +/// and consumers. Per-stream overrides these where it sets a value. +/// (AckPolicy is always Explicit for Wolverine consumers.) +/// public class JetStreamDefaults { - public string Retention { get; set; } = "limits"; public TimeSpan? MaxAge { get; set; } = TimeSpan.FromDays(7); public long? MaxMessages { get; set; } = 1_000_000; public long? MaxBytes { get; set; } = 1024 * 1024 * 1024; public int Replicas { get; set; } = 1; - public string AckPolicy { get; set; } = "explicit"; public TimeSpan AckWait { get; set; } = TimeSpan.FromSeconds(30); - public int MaxDeliver { get; set; } = 3; + + /// + /// Default maximum delivery attempts for auto-provisioned JetStream consumers, and the dead-letter + /// threshold. A per-endpoint ConfigureDeadLetterQueue(maxDeliveryAttempts, ...) overrides this. + /// + public int MaxDeliver { get; set; } = 5; + + /// + /// Deduplication window applied to auto-provisioned streams. Within this window JetStream + /// discards messages carrying a duplicate Nats-Msg-Id (see + /// ). + /// public TimeSpan DuplicateWindow { get; set; } = TimeSpan.FromMinutes(2); /// diff --git a/src/Transports/NATS/Wolverine.Nats/Configuration/StreamConfiguration.cs b/src/Transports/NATS/Wolverine.Nats/Configuration/StreamConfiguration.cs index bf56e52f2..1c41903dd 100644 --- a/src/Transports/NATS/Wolverine.Nats/Configuration/StreamConfiguration.cs +++ b/src/Transports/NATS/Wolverine.Nats/Configuration/StreamConfiguration.cs @@ -18,6 +18,13 @@ public class StreamConfiguration public bool AllowDirect { get; set; } public bool DenyDelete { get; set; } public bool DenyPurge { get; set; } + + /// + /// Deduplication window for this stream. Within this window JetStream discards messages carrying a + /// duplicate Nats-Msg-Id. When null the transport-wide + /// is applied. + /// + public TimeSpan? DuplicateWindow { get; set; } /// /// Enable scheduled message delivery (requires NATS Server 2.12+) @@ -95,6 +102,15 @@ public StreamConfiguration WithReplicas(int replicas) return this; } + /// + /// Set the deduplication window for this stream (see ). + /// + public StreamConfiguration WithDeduplicationWindow(TimeSpan window) + { + DuplicateWindow = window; + return this; + } + /// /// Enable scheduled message delivery (requires NATS Server 2.12+). /// Once enabled on a stream, this cannot be disabled. diff --git a/src/Transports/NATS/Wolverine.Nats/Extensions/NatsTransportExtensions.cs b/src/Transports/NATS/Wolverine.Nats/Extensions/NatsTransportExtensions.cs index 72ddb37a6..8cf50e99b 100644 --- a/src/Transports/NATS/Wolverine.Nats/Extensions/NatsTransportExtensions.cs +++ b/src/Transports/NATS/Wolverine.Nats/Extensions/NatsTransportExtensions.cs @@ -4,6 +4,7 @@ using Wolverine.Nats.Configuration; using Wolverine.Nats.Internal; using Wolverine.Runtime.Partitioning; +using Wolverine.Runtime.Routing; namespace Wolverine.Nats; @@ -113,6 +114,27 @@ string subject return new NatsSubscriberConfiguration(endpoint); } + /// + /// Publish messages of type (or castable to it) to a NATS subject + /// computed per message by . Enables per-message dynamic + /// subjects such as orders.events.{id} using Wolverine's generic topic routing + /// ( / Envelope.TopicName), so these messages also + /// participate in . + /// + public static NatsSubscriberConfiguration PublishMessagesToNatsSubject( + this WolverineOptions options, + Func subjectSource + ) + { + var transport = options.NatsTransport(); + var endpoint = transport.NewTopicSender(); + + var routing = new TopicRouting(subjectSource, endpoint); + options.PublishWithMessageRoutingSource(routing); + + return new NatsSubscriberConfiguration(endpoint); + } + /// /// Listen to messages from a NATS subject /// diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/CoreNatsSubscriber.cs b/src/Transports/NATS/Wolverine.Nats/Internal/CoreNatsSubscriber.cs index 49570027a..a21d8c8d6 100644 --- a/src/Transports/NATS/Wolverine.Nats/Internal/CoreNatsSubscriber.cs +++ b/src/Transports/NATS/Wolverine.Nats/Internal/CoreNatsSubscriber.cs @@ -61,7 +61,7 @@ CancellationToken cancellation "Starting Core NATS listener for pattern {Pattern} (base subject: {Subject}) with queue group {QueueGroup}", _subscriptionPattern, _endpoint.Subject, - _endpoint.QueueGroup ?? "(none)" + _endpoint.EffectiveQueueGroup ?? "(none)" ); } @@ -69,11 +69,12 @@ CancellationToken cancellation { IAsyncDisposable subscription; - if (!string.IsNullOrEmpty(_endpoint.QueueGroup)) + var queueGroup = _endpoint.EffectiveQueueGroup; + if (!string.IsNullOrEmpty(queueGroup)) { subscription = await _connection.SubscribeCoreAsync( pattern, - _endpoint.QueueGroup, + queueGroup, cancellationToken: cancellation ); } diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/JetStreamPublisher.cs b/src/Transports/NATS/Wolverine.Nats/Internal/JetStreamPublisher.cs index 45e1fc073..e0d6e1bf8 100644 --- a/src/Transports/NATS/Wolverine.Nats/Internal/JetStreamPublisher.cs +++ b/src/Transports/NATS/Wolverine.Nats/Internal/JetStreamPublisher.cs @@ -10,19 +10,45 @@ namespace Wolverine.Nats.Internal; /// internal class JetStreamPublisher : INatsPublisher { + /// + /// NATS JetStream deduplication header. When present, the server discards duplicate + /// messages carrying the same value within the stream's configured duplicate window. + /// + internal const string NatsMsgIdHeader = "Nats-Msg-Id"; + private readonly NatsConnection _connection; private readonly INatsJSContext _jetStreamContext; private readonly ILogger _logger; private readonly string _scheduleSubjectSuffix; + private readonly Func? _msgIdSource; - public JetStreamPublisher(NatsConnection connection, + public JetStreamPublisher(NatsConnection connection, + INatsJSContext jetStreamContext, ILogger logger, - string scheduleSubjectSuffix = ".scheduled") + string scheduleSubjectSuffix = ".scheduled", + Func? msgIdSource = null) { _connection = connection; + _jetStreamContext = jetStreamContext; _logger = logger; _scheduleSubjectSuffix = scheduleSubjectSuffix; - _jetStreamContext = connection.CreateJetStreamContext(); + _msgIdSource = msgIdSource; + } + + /// + /// Resolve the JetStream Nats-Msg-Id deduplication key for an outgoing message: + /// an explicit Nats-Msg-Id header wins (return null so we don't override it), then + /// the configured MsgIdSource, else the Wolverine envelope Id. + /// + private NatsJSPubOpts? buildDedupOptions(Envelope envelope, NatsHeaders headers) + { + if (headers.ContainsKey(NatsMsgIdHeader)) + { + return null; + } + + var msgId = _msgIdSource?.Invoke(envelope) ?? envelope.Id.ToString(); + return string.IsNullOrEmpty(msgId) ? null : new NatsJSPubOpts { MsgId = msgId }; } public async ValueTask PingAsync(CancellationToken cancellation) @@ -76,6 +102,12 @@ await _connection.PublishAsync( { var publishSubject = subject; + // Server-side dedup only applies to the direct publish path; the native scheduling + // control message is materialized server-side and is not deduplicated by this key. + var pubOpts = envelope.ScheduledTime.HasValue + ? null + : buildDedupOptions(envelope, headers); + if (envelope.ScheduledTime.HasValue) { // NATS rejects a scheduled publish whose subject equals Nats-Schedule-Target ("message @@ -102,6 +134,7 @@ await _connection.PublishAsync( var ack = await _jetStreamContext.PublishAsync( publishSubject, data, + opts: pubOpts, headers: headers, cancellationToken: cancellation ); diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/JetStreamSubscriber.cs b/src/Transports/NATS/Wolverine.Nats/Internal/JetStreamSubscriber.cs index 607921bde..b82c2a775 100644 --- a/src/Transports/NATS/Wolverine.Nats/Internal/JetStreamSubscriber.cs +++ b/src/Transports/NATS/Wolverine.Nats/Internal/JetStreamSubscriber.cs @@ -21,6 +21,7 @@ internal class JetStreamSubscriber : INatsSubscriber public JetStreamSubscriber( NatsEndpoint endpoint, NatsConnection connection, + INatsJSContext jetStreamContext, ILogger logger, JetStreamEnvelopeMapper mapper, string? subscriptionPattern = null @@ -31,7 +32,7 @@ public JetStreamSubscriber( _logger = logger; _mapper = mapper; _subscriptionPattern = subscriptionPattern ?? endpoint.Subject; - _jetStreamContext = connection.CreateJetStreamContext(); + _jetStreamContext = jetStreamContext; } public bool SupportsNativeDeadLetterQueue => _endpoint.DeadLetterQueueEnabled; @@ -55,8 +56,8 @@ CancellationToken cancellation var config = new ConsumerConfig { AckPolicy = ConsumerConfigAckPolicy.Explicit, - MaxDeliver = _endpoint.MaxDeliveryAttempts, - AckWait = TimeSpan.FromSeconds(30) + MaxDeliver = _endpoint.EffectiveMaxDeliveryAttempts, + AckWait = _endpoint.JetStreamDefaults.AckWait }; // Apply the per-endpoint or transport-wide DeliverPolicy override when set. @@ -79,9 +80,9 @@ CancellationToken cancellation config.Name = _endpoint.ConsumerName; config.DurableName = _endpoint.ConsumerName; - if (!string.IsNullOrEmpty(_endpoint.QueueGroup)) + if (!string.IsNullOrEmpty(_endpoint.EffectiveQueueGroup)) { - config.DeliverGroup = _endpoint.QueueGroup; + config.DeliverGroup = _endpoint.EffectiveQueueGroup; } try diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/NatsEndpoint.cs b/src/Transports/NATS/Wolverine.Nats/Internal/NatsEndpoint.cs index 319585045..1c116287f 100644 --- a/src/Transports/NATS/Wolverine.Nats/Internal/NatsEndpoint.cs +++ b/src/Transports/NATS/Wolverine.Nats/Internal/NatsEndpoint.cs @@ -4,6 +4,7 @@ using NATS.Client.JetStream.Models; using NATS.Net; using Wolverine.Configuration; +using Wolverine.Nats.Configuration; using Wolverine.Runtime; using Wolverine.Transports; using Wolverine.Transports.Sending; @@ -39,13 +40,62 @@ public NatsEndpoint(string subject, NatsTransport transport, EndpointRole role) [IgnoreDescription] public object? NatsSerializer { get; set; } public Dictionary CustomHeaders { get; set; } = new(); + + /// + /// Optional transport-wide hook to rewrite the outgoing subject per envelope + /// (headers, tenant, aggregate id) beyond what per-message topic routing can express. + /// Sourced from . + /// + [IgnoreDescription] + internal ISubjectResolver? SubjectResolver => _transport.Configuration.SubjectResolver; + + /// + /// Optional transport-wide source of the JetStream Nats-Msg-Id dedup key. + /// Sourced from . + /// + [IgnoreDescription] + internal Func? MsgIdSource => _transport.Configuration.MsgIdSource; + + /// + /// Transport-wide JetStream stream/consumer template applied when Wolverine auto-provisions. + /// + [IgnoreDescription] + internal JetStreamDefaults JetStreamDefaults => _transport.Configuration.JetStreamDefaults; + + /// + /// Normalize a per-message subject honoring the transport's + /// flag. + /// + internal string NormalizeSubject(string subject) => _transport.NormalizeSubjectIfEnabled(subject); public string? QueueGroup { get; set; } + + /// + /// The queue group actually used for load-balanced delivery: the per-endpoint + /// when set, otherwise the transport-wide + /// . + /// + [IgnoreDescription] + internal string? EffectiveQueueGroup => + string.IsNullOrEmpty(QueueGroup) ? _transport.Configuration.DefaultQueueGroup : QueueGroup; + public string? StreamName { get; set; } public string? ConsumerName { get; set; } public bool UseJetStream { get; set; } public bool DeadLetterQueueEnabled { get; set; } = true; public string? DeadLetterSubject { get; set; } - public int MaxDeliveryAttempts { get; set; } = 5; + + /// + /// Per-endpoint override for the maximum delivery attempts / dead-letter threshold. When null the + /// transport-wide applies (see ). + /// + public int? MaxDeliveryAttempts { get; set; } + + /// + /// Resolved maximum delivery attempts: the per-endpoint when set, + /// otherwise the transport-wide . + /// + [IgnoreDescription] + internal int EffectiveMaxDeliveryAttempts => MaxDeliveryAttempts ?? JetStreamDefaults.MaxDeliver; /// /// Suffix appended to the destination subject to form the NATS JetStream scheduling subject for native @@ -107,9 +157,12 @@ protected override ISender CreateSender(IWolverineRuntime runtime) _transport.Configuration.Streams.TryGetValue(StreamName, out var streamConfig) && streamConfig.AllowMsgSchedules; + var jetStreamContext = useJetStream ? _transport.CreateJetStreamContext() : null; + var baseSender = NatsSender.Create( this, _connection, + jetStreamContext, _logger, _mapper, runtime.Cancellation, @@ -126,6 +179,12 @@ protected override ISender CreateSender(IWolverineRuntime runtime) var subjectMapper = tenant.SubjectMapper ?? _transport.TenantSubjectMapper; var tenantSubject = subjectMapper.MapSubject(Subject, tenant.TenantId); + // Tenants that declare their own connection/credentials publish over their dedicated + // connection (with its own JetStream context); the rest reuse the shared connection. + var tenantConnection = _transport.GetTenantConnection(tenant); + var tenantJetStreamContext = + useJetStream ? _transport.CreateJetStreamContext(tenantConnection) : null; + var tenantEndpoint = new NatsEndpoint(tenantSubject, _transport, Role) { UseJetStream = UseJetStream, @@ -143,12 +202,15 @@ protected override ISender CreateSender(IWolverineRuntime runtime) var tenantSender = NatsSender.Create( tenantEndpoint, - _connection, + tenantConnection, + tenantJetStreamContext, _logger, _mapper, runtime.Cancellation, useJetStream, - supportsScheduledSend + supportsScheduledSend, + subjectMapper, + tenant.TenantId ); tenantedSender.RegisterSender(tenant.TenantId, tenantSender); @@ -175,46 +237,80 @@ IReceiver receiver deadLetterSender = (ISender)runtime.Endpoints.GetOrBuildSendingAgent(dlqEndpoint.Uri); } + var useJetStream = UseJetStream && _transport.Configuration.EnableJetStream; + string subscriptionPattern = Subject; ITenantSubjectMapper? tenantMapper = null; + var tenantAware = _transport.Tenants.Any() && TenancyBehavior == TenancyBehavior.TenantAware; - if (_transport.Tenants.Any() && TenancyBehavior == TenancyBehavior.TenantAware) + if (tenantAware) { tenantMapper = _transport.TenantSubjectMapper; subscriptionPattern = tenantMapper.GetSubscriptionPattern(Subject); } + // The shared listener consumes the default connection plus every subject-prefix tenant, whose messages + // arrive on the shared connection under the wildcard pattern. + var sharedListener = await startListenerAsync( + runtime, receiver, _connection, useJetStream, subscriptionPattern, tenantMapper, deadLetterSender); + + // Tenants with their own connection publish on a separate server/account the shared listener can't see, + // so consume each over its own connection and stamp the tenant id onto inbound envelopes. Mirrors the + // RabbitMQ / Azure Service Bus CompoundListener multi-tenancy pattern; per-envelope completion routes + // back to the right connection via Envelope.Listener. + var dedicatedTenants = tenantAware + ? _transport.Tenants.Where(t => t.HasOwnConnection).ToArray() + : []; + + if (dedicatedTenants.Length == 0) + { + return sharedListener; + } + + var compound = new CompoundListener(Uri); + compound.Inner.Add(sharedListener); + + foreach (var tenant in dedicatedTenants) + { + var tenantReceiver = new ReceiverWithRules(receiver, [new TenantIdRule(tenant.TenantId)]); + var tenantListener = await startListenerAsync( + runtime, tenantReceiver, _transport.GetTenantConnection(tenant), useJetStream, + subscriptionPattern, tenantMapper, deadLetterSender); + compound.Inner.Add(tenantListener); + } + + return compound; + } + + private async ValueTask startListenerAsync( + IWolverineRuntime runtime, + IReceiver receiver, + NatsConnection connection, + bool useJetStream, + string subscriptionPattern, + ITenantSubjectMapper? tenantMapper, + ISender? deadLetterSender) + { + var jetStreamContext = useJetStream ? _transport.CreateJetStreamContext(connection) : null; + var listener = NatsListener.Create( this, - _connection, + connection, + jetStreamContext, runtime, receiver, - _logger, + _logger!, deadLetterSender, runtime.Cancellation, - UseJetStream && _transport.Configuration.EnableJetStream, + useJetStream, subscriptionPattern, tenantMapper ); await listener.StartAsync(); - return listener; } - public NatsHeaders BuildHeaders(Envelope envelope) - { - var headers = new NatsHeaders(); - _mapper?.MapEnvelopeToOutgoing(envelope, headers); - - foreach (var header in CustomHeaders) - { - headers[header.Key] = header.Value; - } - - return headers; - } - public async ValueTask CheckAsync() { _connection ??= _transport.Connection; @@ -231,7 +327,7 @@ public async ValueTask CheckAsync() try { - var js = _connection.CreateJetStreamContext(); + var js = _transport.CreateJetStreamContext(); var stream = await js.GetStreamAsync( StreamName, cancellationToken: CancellationToken.None @@ -262,7 +358,7 @@ public async ValueTask TeardownAsync(ILogger logger) { try { - var js = _connection.CreateJetStreamContext(); + var js = _transport.CreateJetStreamContext(); await js.DeleteConsumerAsync(StreamName, ConsumerName); logger.LogInformation( "Deleted consumer {Consumer} from stream {Stream}", @@ -323,13 +419,16 @@ public async ValueTask SetupAsync(ILogger logger) string.Join(", ", subjects) ); + var defaults = JetStreamDefaults; var config = new StreamConfig(StreamName, subjects) { Retention = StreamConfigRetention.Workqueue, Discard = StreamConfigDiscard.Old, - MaxAge = TimeSpan.FromDays(1), - DuplicateWindow = TimeSpan.FromMinutes(2), - MaxMsgs = 1_000_000 + MaxAge = defaults.MaxAge ?? TimeSpan.Zero, + MaxMsgs = defaults.MaxMessages ?? -1, + MaxBytes = defaults.MaxBytes ?? -1, + NumReplicas = defaults.Replicas, + DuplicateWindow = defaults.DuplicateWindow }; await js.CreateStreamAsync(config); @@ -344,8 +443,8 @@ public async ValueTask SetupAsync(ILogger logger) DurableName = ConsumerName, FilterSubject = Subject, AckPolicy = ConsumerConfigAckPolicy.Explicit, - AckWait = TimeSpan.FromSeconds(30), - MaxDeliver = MaxDeliveryAttempts, + AckWait = JetStreamDefaults.AckWait, + MaxDeliver = EffectiveMaxDeliveryAttempts, ReplayPolicy = ConsumerConfigReplayPolicy.Instant }; diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/NatsListener.cs b/src/Transports/NATS/Wolverine.Nats/Internal/NatsListener.cs index b00ce09fa..edd4fd9b7 100644 --- a/src/Transports/NATS/Wolverine.Nats/Internal/NatsListener.cs +++ b/src/Transports/NATS/Wolverine.Nats/Internal/NatsListener.cs @@ -1,6 +1,7 @@ using JasperFx.Blocks; using Microsoft.Extensions.Logging; using NATS.Client.Core; +using NATS.Client.JetStream; using Wolverine.Runtime; using Wolverine.Transports; using Wolverine.Transports.Sending; @@ -85,38 +86,57 @@ await envelope.JetStreamMsg.NakAsync( public async Task MoveToErrorsAsync(Envelope envelope, Exception exception) { - if (envelope is NatsEnvelope natsEnvelope) + if (envelope is not NatsEnvelope natsEnvelope || !NativeDeadLetterQueueEnabled || + natsEnvelope.JetStreamMsg == null) { - if (NativeDeadLetterQueueEnabled && natsEnvelope.JetStreamMsg != null) - { - var metadata = natsEnvelope.JetStreamMsg.Metadata; - - if (metadata?.NumDelivered >= (ulong)_endpoint.MaxDeliveryAttempts) - { - await natsEnvelope.JetStreamMsg.AckAsync( - cancellationToken: _cancellation.Token - ); + return; + } - if (!string.IsNullOrEmpty(_endpoint.DeadLetterSubject)) - { - envelope.Attempts = (int)(metadata?.NumDelivered ?? 1); + var metadata = natsEnvelope.JetStreamMsg.Metadata; + if (metadata?.NumDelivered < (ulong)_endpoint.EffectiveMaxDeliveryAttempts) + { + return; + } - DeadLetterQueueConstants.StampFailureMetadata(envelope, exception); - envelope.Headers["x-dlq-original-subject"] = _endpoint.Subject; + var attempts = (int)(metadata?.NumDelivered ?? 1); - await _deadLetterSender.SendAsync(envelope); - } + // Retain the poison message by forwarding a copy to the dead-letter subject BEFORE terminating, + // so a terminate failure can't lose it. Terminating without a configured dead-letter subject drops + // the message, so warn loudly in that case. + if (!string.IsNullOrEmpty(_endpoint.DeadLetterSubject)) + { + envelope.Attempts = attempts; + DeadLetterQueueConstants.StampFailureMetadata(envelope, exception); + envelope.Headers["x-dlq-original-subject"] = _endpoint.Subject; - _logger.LogError( - exception, - "Message {MessageId} moved to dead letter queue after {Attempts} attempts. Subject: {Subject}", - envelope.Id, - metadata?.NumDelivered ?? 1, - _endpoint.DeadLetterSubject - ); - } - } + await _deadLetterSender.SendAsync(envelope); + } + else + { + _logger.LogWarning( + exception, + "Message {MessageId} exceeded {Attempts} delivery attempts on subject {Subject} but no dead-letter subject is configured; it will be terminated and dropped. Use DeadLetterTo(...) / ConfigureDeadLetterQueue(...) to retain poison messages.", + envelope.Id, + attempts, + _endpoint.Subject + ); } + + // Terminate delivery on the JetStream consumer with a reason so the server stops redelivering and + // records why the message was dead-lettered. + await natsEnvelope.JetStreamMsg.AckTerminateAsync( + $"wolverine: exceeded {attempts} delivery attempts ({exception.GetType().Name})", + cancellationToken: _cancellation.Token + ); + + _logger.LogError( + exception, + "Message {MessageId} terminated after {Attempts} delivery attempts. Subject: {Subject}, DeadLetter: {DeadLetter}", + envelope.Id, + attempts, + _endpoint.Subject, + _endpoint.DeadLetterSubject ?? "(none)" + ); } public async ValueTask CompleteAsync(Envelope envelope) @@ -168,6 +188,7 @@ public async ValueTask DisposeAsync() internal static NatsListener Create( NatsEndpoint endpoint, NatsConnection connection, + INatsJSContext? jetStreamContext, IWolverineRuntime runtime, IReceiver receiver, ILogger logger, @@ -186,7 +207,7 @@ internal static NatsListener Create( { jsMapper.ReceivesMessage(endpoint.MessageType); } - subscriber = new JetStreamSubscriber(endpoint, connection, logger, jsMapper, subscriptionPattern); + subscriber = new JetStreamSubscriber(endpoint, connection, jetStreamContext!, logger, jsMapper, subscriptionPattern); } else { diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/NatsSender.cs b/src/Transports/NATS/Wolverine.Nats/Internal/NatsSender.cs index 3d6207910..247d0b432 100644 --- a/src/Transports/NATS/Wolverine.Nats/Internal/NatsSender.cs +++ b/src/Transports/NATS/Wolverine.Nats/Internal/NatsSender.cs @@ -1,5 +1,7 @@ +using JasperFx.Core; using Microsoft.Extensions.Logging; using NATS.Client.Core; +using NATS.Client.JetStream; using Wolverine.Transports.Sending; namespace Wolverine.Nats.Internal; @@ -12,6 +14,8 @@ public class NatsSender : ISender private readonly CancellationToken _cancellation; private readonly INatsPublisher _publisher; private readonly bool _supportsNativeScheduledSend; + private readonly ITenantSubjectMapper? _tenantSubjectMapper; + private readonly string? _tenantId; internal NatsSender( NatsEndpoint endpoint, @@ -19,7 +23,9 @@ internal NatsSender( ILogger logger, NatsEnvelopeMapper mapper, CancellationToken cancellation, - bool supportsNativeScheduledSend + bool supportsNativeScheduledSend, + ITenantSubjectMapper? tenantSubjectMapper = null, + string? tenantId = null ) { _endpoint = endpoint; @@ -28,24 +34,30 @@ bool supportsNativeScheduledSend _mapper = mapper; _cancellation = cancellation; _supportsNativeScheduledSend = supportsNativeScheduledSend; + _tenantSubjectMapper = tenantSubjectMapper; + _tenantId = tenantId; Destination = endpoint.Uri; } internal static NatsSender Create( NatsEndpoint endpoint, NatsConnection connection, + INatsJSContext? jetStreamContext, ILogger logger, NatsEnvelopeMapper mapper, CancellationToken cancellation, bool useJetStream, - bool supportsNativeScheduledSend + bool supportsNativeScheduledSend, + ITenantSubjectMapper? tenantSubjectMapper = null, + string? tenantId = null ) { INatsPublisher publisher = useJetStream - ? new JetStreamPublisher(connection, logger, endpoint.ScheduleSubjectSuffix) + ? new JetStreamPublisher(connection, jetStreamContext!, logger, endpoint.ScheduleSubjectSuffix, endpoint.MsgIdSource) : new CoreNatsPublisher(connection, logger); - return new NatsSender(endpoint, publisher, logger, mapper, cancellation, supportsNativeScheduledSend); + return new NatsSender(endpoint, publisher, logger, mapper, cancellation, supportsNativeScheduledSend, + tenantSubjectMapper, tenantId); } public bool SupportsNativeScheduledSend => _supportsNativeScheduledSend; @@ -70,7 +82,7 @@ public async ValueTask SendAsync(Envelope envelope) var data = envelope.Data ?? Array.Empty(); - var targetSubject = _endpoint.Subject; + string targetSubject; string? replyTo = null; if (envelope.IsResponse && envelope.Destination != null) @@ -90,6 +102,38 @@ public async ValueTask SendAsync(Envelope envelope) } else { + // Per-message subject routing: when the endpoint is RoutingMode.ByTopic (see + // PublishMessagesToNatsSubject / IMessageBus.BroadcastToTopicAsync) Wolverine + // stamps the computed subject onto Envelope.TopicName. Static endpoints leave it + // null and fall back to the endpoint's fixed subject. Mirrors RabbitMqSender and + // InlineKafkaSender. + if (envelope.TopicName.IsNotEmpty()) + { + // A per-message subject (RoutingMode.ByTopic). Unlike a static endpoint subject — which is + // tenant-qualified once at construction (see NatsEndpoint.CreateSender) — this computed + // subject arrives un-prefixed, so apply the same tenant mapping here. Otherwise a + // subject-isolation tenant's dynamic-subject sends would publish without the tenant prefix + // and defeat isolation. + targetSubject = _endpoint.NormalizeSubject(envelope.TopicName); + if (_tenantSubjectMapper is not null && !string.IsNullOrEmpty(_tenantId)) + { + targetSubject = _tenantSubjectMapper.MapSubject(targetSubject, _tenantId); + } + } + else + { + targetSubject = _endpoint.Subject; + } + + // Advanced escape hatch: rewrite the subject from envelope-level state (headers, + // tenant, aggregate id) that the strongly-typed subject function can't express. Run the + // result back through NormalizeSubject so a resolver's output honors NormalizeSubjects the + // same way static subjects and TopicName routing do (a no-op beyond trimming when disabled). + if (_endpoint.SubjectResolver is { } resolver) + { + targetSubject = _endpoint.NormalizeSubject(resolver.ResolveSubject(targetSubject, envelope)); + } + if (envelope.ReplyRequested != null && envelope.ReplyUri != null) { replyTo = NatsTransport.ExtractSubjectFromUri(envelope.ReplyUri); diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/NatsTenant.cs b/src/Transports/NATS/Wolverine.Nats/Internal/NatsTenant.cs index 27fba3e2d..ac4cf7720 100644 --- a/src/Transports/NATS/Wolverine.Nats/Internal/NatsTenant.cs +++ b/src/Transports/NATS/Wolverine.Nats/Internal/NatsTenant.cs @@ -1,3 +1,6 @@ +using NATS.Client.Core; +using Wolverine.Nats.Configuration; + namespace Wolverine.Nats.Internal; public class NatsTenant @@ -9,9 +12,24 @@ public NatsTenant(string tenantId) public string TenantId { get; } public ITenantSubjectMapper? SubjectMapper { get; set; } - public string? ConnectionString { get; set; } - public string? CredentialsFile { get; set; } - public string? Username { get; set; } - public string? Password { get; set; } - public string? Token { get; set; } + + /// + /// The tenant's own connection configuration (full auth / TLS surface, via + /// ) when it should use a dedicated NATS connection + /// rather than the shared transport connection. Null means subject-prefix isolation on the shared connection. + /// + public NatsTransportConfiguration? ConnectionConfiguration { get; set; } + + /// + /// The tenant's dedicated NATS connection when is true. Created and owned + /// by the transport during its ConnectAsync, and disposed with the transport. Null when the tenant reuses + /// the shared transport connection. + /// + internal NatsConnection? Connection { get; set; } + + /// + /// True when this tenant declares its own connection configuration, and therefore gets a dedicated NATS + /// connection rather than sharing the transport's connection. + /// + public bool HasOwnConnection => ConnectionConfiguration != null; } diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/NatsTransport.cs b/src/Transports/NATS/Wolverine.Nats/Internal/NatsTransport.cs index 8b1339294..0e1db2683 100644 --- a/src/Transports/NATS/Wolverine.Nats/Internal/NatsTransport.cs +++ b/src/Transports/NATS/Wolverine.Nats/Internal/NatsTransport.cs @@ -38,7 +38,7 @@ public NatsTransport() { _endpoints.OnMissing = subject => { - var normalized = NormalizeSubject(subject); + var normalized = NormalizeSubjectIfEnabled(subject); return new NatsEndpoint(normalized, this, EndpointRole.Application); }; } @@ -144,14 +144,33 @@ public override async ValueTask ConnectAsync(IWolverineRuntime runtime) } } + var autoProvisionStreams = Configuration.AutoProvision && Configuration.Streams.Any(); + if (Configuration.EnableJetStream) { - _jetStreamContext = _connection.CreateJetStreamContext(); + _jetStreamContext = CreateJetStreamContext(); _logger.LogInformation("JetStream context initialized"); - if (Configuration.AutoProvision && Configuration.Streams.Any()) + if (autoProvisionStreams) + { + await ProvisionStreamsAsync(_jetStreamContext); + } + } + + // Tenants that declare their own connection string / credentials get a dedicated connection they + // own for the lifetime of the transport; the NATS client connects lazily on first use. Tenants + // without their own connection reuse the shared connection above (subject-prefix isolation only). + foreach (var tenant in Tenants.Where(x => x.HasOwnConnection)) + { + var tenantConnection = new NatsConnection(buildTenantNatsOpts(tenant)); + tenant.Connection = tenantConnection; + _logger.LogInformation("Created dedicated NATS connection for tenant {TenantId}", tenant.TenantId); + + // Each tenant server is its own JetStream instance, so mirror the configured streams onto it + // (the streams the shared connection just provisioned don't exist on the tenant's server). + if (Configuration.EnableJetStream && autoProvisionStreams) { - await ProvisionStreamsAsync(); + await ProvisionStreamsAsync(CreateJetStreamContext(tenantConnection)); } } } @@ -174,6 +193,65 @@ public static string NormalizeSubject(string subject) return subject.Trim().Replace('/', '.'); } + /// + /// Normalize a subject honoring : when the flag + /// is enabled (the default) '/' separators are converted to NATS '.' tokens; when disabled the subject is + /// only trimmed, so callers can use literal subjects containing '/'. + /// + internal string NormalizeSubjectIfEnabled(string subject) + { + return Configuration.NormalizeSubjects ? NormalizeSubject(subject) : subject.Trim(); + } + + /// + /// Create a JetStream context on the shared connection honoring the configured + /// / . + /// + internal INatsJSContext CreateJetStreamContext() => CreateJetStreamContext(Connection); + + /// + /// Create a JetStream context on the given connection honoring the configured JetStream domain / API prefix. + /// All JetStream context creation flows through this factory so domain / leaf-node setups work uniformly + /// (including per-tenant connections). When neither is configured the result is identical to the client + /// default (connection.CreateJetStreamContext()). + /// + internal INatsJSContext CreateJetStreamContext(NatsConnection connection) + { + var domain = Configuration.JetStreamDomain; + var apiPrefix = Configuration.JetStreamApiPrefix; + + if (string.IsNullOrWhiteSpace(domain) && string.IsNullOrWhiteSpace(apiPrefix)) + { + return connection.CreateJetStreamContext(); + } + + // NatsJSOpts forbids setting both ApiPrefix and Domain; when both are supplied domain wins. + var jsOpts = string.IsNullOrWhiteSpace(domain) + ? new NatsJSOpts(connection.Opts, apiPrefix: apiPrefix) + : new NatsJSOpts(connection.Opts, domain: domain); + + return connection.CreateJetStreamContext(jsOpts); + } + + /// + /// Resolve the NATS connection for a tenant: the tenant's own dedicated connection (created during + /// ) when it declares its own connection string / credentials, otherwise the + /// shared transport connection. + /// + internal NatsConnection GetTenantConnection(NatsTenant tenant) + { + return tenant.HasOwnConnection ? tenant.Connection ?? Connection : Connection; + } + + private static NatsOpts buildTenantNatsOpts(NatsTenant tenant) + { + // The tenant carries its own full connection configuration (URL + any of the NATS auth mechanisms + + // TLS), so we reuse the same ToNatsOpts() the shared connection uses rather than privileging one + // credential kind. Only the client name is decorated so tenant connections are distinguishable. + var opts = tenant.ConnectionConfiguration!.ToNatsOpts(); + return opts with { Name = $"{opts.Name}-tenant-{tenant.TenantId}" }; + } + public static string ExtractSubjectFromUri(Uri uri) { if (uri.Scheme != "nats") @@ -187,12 +265,48 @@ public static string ExtractSubjectFromUri(Uri uri) public NatsEndpoint EndpointForSubject(string subject) { - var normalized = NormalizeSubject(subject); + var normalized = NormalizeSubjectIfEnabled(subject); return _endpoints[normalized]; } + /// + /// Base subject for on-demand topic-routed sending endpoints created by + /// PublishMessagesToNatsSubject<T>. The real destination is the per-message + /// subject stamped onto ; this is only a base/fallback. + /// + internal const string TopicSenderSubject = "wolverine.topics"; + + private int _topicSenderIndex; + + /// + /// Create a new topic-routed () sending endpoint so + /// messages can be published to a per-message subject computed at send time. Mirrors the + /// MQTT transport's NewTopicSender; each call returns a distinct endpoint so multiple + /// subject-source functions can coexist. Being ByTopic also enrolls the endpoint in + /// . + /// + internal NatsEndpoint NewTopicSender() + { + var subject = $"{TopicSenderSubject}.{++_topicSenderIndex}"; + var endpoint = _endpoints[subject]; + endpoint.RoutingType = RoutingMode.ByTopic; + return endpoint; + } + public async ValueTask DisposeAsync() { + foreach (var tenant in Tenants.Where(x => x.Connection != null)) + { + try + { + await tenant.Connection!.DisposeAsync(); + } + catch (Exception ex) + { + _logger?.LogError(ex, "Error disposing NATS connection for tenant {TenantId}", tenant.TenantId); + } + } + try { if (_connection != null) @@ -206,7 +320,7 @@ public async ValueTask DisposeAsync() } } - private async Task ProvisionStreamsAsync() + private async Task ProvisionStreamsAsync(INatsJSContext js) { _logger?.LogInformation( "Provisioning {Count} configured streams", @@ -220,7 +334,7 @@ private async Task ProvisionStreamsAsync() var exists = false; try { - await JetStreamContext.GetStreamAsync(name); + await js.GetStreamAsync(name); exists = true; _logger?.LogDebug("Stream {StreamName} already exists", name); } @@ -240,6 +354,7 @@ private async Task ProvisionStreamsAsync() MaxMsgsPerSubject = config.MaxMessagesPerSubject ?? 0, Discard = config.DiscardPolicy, NumReplicas = config.Replicas, + DuplicateWindow = config.DuplicateWindow ?? Configuration.JetStreamDefaults.DuplicateWindow, AllowRollupHdrs = config.AllowRollup, AllowDirect = config.AllowDirect, DenyDelete = config.DenyDelete, @@ -247,7 +362,7 @@ private async Task ProvisionStreamsAsync() AllowMsgSchedules = config.AllowMsgSchedules }; - await JetStreamContext.CreateStreamAsync(streamConfig); + await js.CreateStreamAsync(streamConfig); _logger?.LogInformation( "Created stream {StreamName} with subjects: {Subjects}", name, diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/NatsTransportExpression.cs b/src/Transports/NATS/Wolverine.Nats/Internal/NatsTransportExpression.cs index 028b6645d..bfad7c874 100644 --- a/src/Transports/NATS/Wolverine.Nats/Internal/NatsTransportExpression.cs +++ b/src/Transports/NATS/Wolverine.Nats/Internal/NatsTransportExpression.cs @@ -27,6 +27,18 @@ public NatsTransportExpression UseJetStream(Action configure) return this; } + /// + /// Project a domain identity into the JetStream deduplication key (Nats-Msg-Id) instead + /// of the default Wolverine envelope Id, so non-Wolverine consumers get server-side dedup within + /// the stream's duplicate window. Example: e => $"{stream}/{version}". An explicit + /// Nats-Msg-Id header already on the outgoing envelope still wins. + /// + public NatsTransportExpression DeduplicateUsing(Func msgIdSource) + { + Transport.Configuration.MsgIdSource = msgIdSource; + return this; + } + /// /// Set the identifier prefix for all NATS subjects /// @@ -221,39 +233,104 @@ public NatsTransportExpression AddTenant(string tenantId, ITenantSubjectMapper m } /// - /// Add a tenant with custom credentials + /// Add a tenant that uses its own dedicated NATS connection, configured with the full connection / auth / + /// TLS surface. The action receives a configuration seeded from the transport's own settings, so you only + /// override what differs for this tenant — e.g. a different server or account, a token, JWT / NKey creds, + /// a credentials file, or a client certificate. /// - public NatsTransportExpression AddTenantWithCredentials( + public NatsTransportExpression AddTenant( string tenantId, - string username, - string password + Action configureConnection ) { - var tenant = new NatsTenant(tenantId) - { - Username = username, - Password = password - }; - Transport.Tenants[tenantId] = tenant; - return this; + return AddTenant(tenantId, null, configureConnection); } /// - /// Add a tenant with JWT credentials file + /// Add a tenant with a dedicated connection (see + /// ) and a custom subject mapper. /// - public NatsTransportExpression AddTenantWithCredentialsFile( + public NatsTransportExpression AddTenant( string tenantId, - string credentialsFile + ITenantSubjectMapper? mapper, + Action configureConnection ) { - var tenant = new NatsTenant(tenantId) + ArgumentNullException.ThrowIfNull(configureConnection); + + var configuration = cloneConnectionConfiguration(); + configureConnection(configuration); + + Transport.Tenants[tenantId] = new NatsTenant(tenantId) { - CredentialsFile = credentialsFile + SubjectMapper = mapper, + ConnectionConfiguration = configuration }; - Transport.Tenants[tenantId] = tenant; return this; } + // Seed a tenant connection configuration from the transport's own settings so a tenant only needs to + // override what differs (mirrors RabbitMqTenant.Compile copying the parent's connection settings). + private NatsTransportConfiguration cloneConnectionConfiguration() + { + var source = Transport.Configuration; + + var clone = new NatsTransportConfiguration(); + foreach (var property in typeof(NatsTransportConfiguration).GetProperties()) + { + if (property is { CanRead: true, CanWrite: true }) + { + property.SetValue(clone, property.GetValue(source)); + } + } + + // The reflective copy above aliases the two mutable reference members by reference, so a tenant action + // that mutates them in place (e.g. cfg.JetStreamDefaults.DuplicateWindow = ... or cfg.Streams[...] = ...) + // would leak into the shared transport config and every other tenant. Give each tenant its own copy. + clone.JetStreamDefaults = cloneJetStreamDefaults(source.JetStreamDefaults); + clone.Streams = source.Streams.ToDictionary(pair => pair.Key, pair => cloneStream(pair.Value)); + + return clone; + } + + private static JetStreamDefaults cloneJetStreamDefaults(JetStreamDefaults source) + { + return new JetStreamDefaults + { + MaxAge = source.MaxAge, + MaxMessages = source.MaxMessages, + MaxBytes = source.MaxBytes, + Replicas = source.Replicas, + AckWait = source.AckWait, + MaxDeliver = source.MaxDeliver, + DuplicateWindow = source.DuplicateWindow, + DeliverPolicy = source.DeliverPolicy + }; + } + + private static StreamConfiguration cloneStream(StreamConfiguration source) + { + return new StreamConfiguration + { + Name = source.Name, + Subjects = new List(source.Subjects), + Retention = source.Retention, + Storage = source.Storage, + MaxMessages = source.MaxMessages, + MaxBytes = source.MaxBytes, + MaxAge = source.MaxAge, + MaxMessagesPerSubject = source.MaxMessagesPerSubject, + DiscardPolicy = source.DiscardPolicy, + Replicas = source.Replicas, + AllowRollup = source.AllowRollup, + AllowDirect = source.AllowDirect, + DenyDelete = source.DenyDelete, + DenyPurge = source.DenyPurge, + DuplicateWindow = source.DuplicateWindow, + AllowMsgSchedules = source.AllowMsgSchedules + }; + } + protected override NatsListenerConfiguration createListenerExpression( NatsEndpoint listenerEndpoint ) diff --git a/src/Transports/NATS/Wolverine.Nats/Internal/TenantAwareNatsSender.cs b/src/Transports/NATS/Wolverine.Nats/Internal/TenantAwareNatsSender.cs deleted file mode 100644 index 73b647820..000000000 --- a/src/Transports/NATS/Wolverine.Nats/Internal/TenantAwareNatsSender.cs +++ /dev/null @@ -1,43 +0,0 @@ -using Wolverine.Transports.Sending; - -namespace Wolverine.Nats.Internal; - -/// -/// A wrapper around a base NATS sender that applies tenant-specific subject mapping -/// -internal class TenantAwareNatsSender : ISender -{ - private readonly ISender _innerSender; - private readonly string _tenantId; - private readonly ITenantSubjectMapper _subjectMapper; - - public TenantAwareNatsSender(ISender innerSender, string tenantId, ITenantSubjectMapper subjectMapper) - { - _innerSender = innerSender ?? throw new ArgumentNullException(nameof(innerSender)); - _tenantId = tenantId ?? throw new ArgumentNullException(nameof(tenantId)); - _subjectMapper = subjectMapper ?? throw new ArgumentNullException(nameof(subjectMapper)); - } - - public Uri Destination => _innerSender.Destination; - - public bool SupportsNativeScheduledSend => _innerSender.SupportsNativeScheduledSend; - - public Task PingAsync() => _innerSender.PingAsync(); - - public async ValueTask SendAsync(Envelope envelope) - { - var tenantEnvelope = new Envelope(envelope) - { - TenantId = _tenantId - }; - - if (tenantEnvelope.Destination != null) - { - var originalSubject = NatsTransport.ExtractSubjectFromUri(tenantEnvelope.Destination); - var tenantSubject = _subjectMapper.MapSubject(originalSubject, _tenantId); - tenantEnvelope.Destination = new Uri($"nats://subject/{tenantSubject}"); - } - - await _innerSender.SendAsync(tenantEnvelope); - } -} \ No newline at end of file