Rebus.Nats
0.0.4
dotnet add package Rebus.Nats --version 0.0.4
NuGet\Install-Package Rebus.Nats -Version 0.0.4
<PackageReference Include="Rebus.Nats" Version="0.0.4" />
<PackageVersion Include="Rebus.Nats" Version="0.0.4" />
<PackageReference Include="Rebus.Nats" />
paket add Rebus.Nats --version 0.0.4
#r "nuget: Rebus.Nats, 0.0.4"
#:package Rebus.Nats@0.0.4
#addin nuget:?package=Rebus.Nats&version=0.0.4
#tool nuget:?package=Rebus.Nats&version=0.0.4
Rebus NATS
This is a port of the Rebus.Redis library to use NATS, supporting saga persistence, outbox, and async messaging, as well as a transport implementation which was not present in the original Rebus.Redis implementation.
Async Support
There is an existing async library for Rebus which uses the normal Rebus transport to send a reply, however it is marked experimental and the reasoning given for this is quite a reasonable one that durable messages are not suitable for the ephemeral state of an async request. A pending async await is by nature ephemeral, using a persistent queue to send a reply is undesirable. In place of using the normal Rebus transport, this library uses NATS publish/subscribe to send the reply, only currently subscribed listeners will receive the reply making it well suited for this use case.
Using async you can make a call like this on the client:
var response = await bus.SendRequest<ReplyMessage>(request);
This will send the message to the server as normal, as well as add a task to the NATS subscription. On the server you can reply to the pending task like this:
await bus.ReplyAsync(replyMessage);
There are some additional methods to allow flexibility in cases like sagas where the handler is not ready to reply until some further action is taken. Calling GetReplyContext in the context of a NATS async request will return a context to allow you to send a reply at some later date. This context is just an identifier, it is safe to store and can be added to a saga state, allowing a future message to reply to the original request.
var replyContext = messageContext.GetReplyContext();
await replyContext.ReplyAsync(replyMessage);
In addition, the timeout for the caller is sent along with the request so that the recipient of a message can determine how long the caller will be waiting for a response, which may be useful for cancelling a task or determining whether to send a response to the caller.
// in the client (the default timeout is 15 seconds if not specified)
var response = await bus.SendRequest<ReplyMessage>(request, timeout: TimeSpan.FromSeconds(30));
// on the handler
var timeout = messageContext.GetReplyTimeout();
Async Configuration
To configure async messaging, you need to enable NATS and configure the async messaging. This can be done as follows:
Configure.With(activationHandler)
.Options(o =>
{
o.SetBusName("main");
o.EnableNats("nats://localhost:4222", r => r.EnableAsync());
})
// ...
By default, both client and server mode will be active. This means that a listener will be started to listen for replies from dispatched requests and that a step handler will be registered to redirect replies sent from a NATS request to the NATS publish channel. If you only want to use async messaging in one direction, you can disable the other mode as follows:
Configure.With(activationHandler)
.Options(o =>
{
o.SetBusName("main");
o.EnableNats("nats://localhost:4222", r => r.EnableAsync(AsyncMode.Client)); // or AsyncMode.Host
})
// ...
Typically only one service would use NATS async messaging, e.g. a client facing service. If however you need to send replies via NATS from one service to another and if each one has its own NATS server, you can configure the replies to be routed based on the sender address. This can be done using the RouteRepliesTo method on the NATS configuration:
Configure.With(activationHandler)
.Options(o =>
{
o.SetBusName("main");
o.EnableNats("nats://main-nats:4222", r => r.EnableAsync()
.RouteRepliesTo("other-service", "nats://other-nats:4222"));
})
// ...
Note that this impacts only the reply routing, all other NATS components will use the main NATS connection configured when calling EnableNats.
Saga Storage and Outbox
This library also provides a NATS implementation of the saga storage and an outbox implementation modeled after the Postgres implementation in Rebus. The outbox is implemented using NATS JetStream, and saga data is stored using NATS key-value stores.
Subscriptions (pub/sub) are handled natively by the NATS transport using JetStream subjects and consumers, so no separate subscription storage configuration is needed.
Basic configuration for the saga storage and outbox is as follows:
using var activationHandler = new BuiltinHandlerActivator();
Configure.With(activationHandler)
.Options(o =>
{
o.SetBusName("main");
o.EnableNats("nats://localhost:4222", r => r.EnableAsync());
})
.Transport(t => t.UseNatsJetStream("my-queue")) // Subscriptions handled automatically
.Outbox(o => o.StoreInNats())
.Sagas(s => s.StoreInNats());
This outbox only makes sense to use when the activity being performed is also stored in NATS, e.g. for sagas that use NATS storage.
Internally the outbox is implemented using NATS JetStream using the NATS .NET client. Messages are consumed using push or pull consumers depending on the configuration, with automatic acknowledgment handling to ensure reliable delivery.
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net5.0 was computed. net5.0-windows was computed. net6.0 was computed. net6.0-android was computed. net6.0-ios was computed. net6.0-maccatalyst was computed. net6.0-macos was computed. net6.0-tvos was computed. net6.0-windows was computed. net7.0 was computed. net7.0-android was computed. net7.0-ios was computed. net7.0-maccatalyst was computed. net7.0-macos was computed. net7.0-tvos was computed. net7.0-windows was computed. net8.0 is compatible. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 was computed. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 was computed. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
| .NET Core | netcoreapp2.0 was computed. netcoreapp2.1 was computed. netcoreapp2.2 was computed. netcoreapp3.0 was computed. netcoreapp3.1 was computed. |
| .NET Standard | netstandard2.0 is compatible. netstandard2.1 was computed. |
| .NET Framework | net461 was computed. net462 was computed. net463 was computed. net47 was computed. net471 was computed. net472 was computed. net48 was computed. net481 was computed. |
| MonoAndroid | monoandroid was computed. |
| MonoMac | monomac was computed. |
| MonoTouch | monotouch was computed. |
| Tizen | tizen40 was computed. tizen60 was computed. |
| Xamarin.iOS | xamarinios was computed. |
| Xamarin.Mac | xamarinmac was computed. |
| Xamarin.TVOS | xamarintvos was computed. |
| Xamarin.WatchOS | xamarinwatchos was computed. |
-
.NETStandard 2.0
- NATS.Client.Core (>= 2.6.11)
- NATS.Client.JetStream (>= 2.6.11)
- NATS.Client.KeyValueStore (>= 2.6.11)
- Rebus (>= 8.9.0)
-
net8.0
- NATS.Client.Core (>= 2.6.11)
- NATS.Client.JetStream (>= 2.6.11)
- NATS.Client.KeyValueStore (>= 2.6.11)
- Rebus (>= 8.9.0)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.