-
Notifications
You must be signed in to change notification settings - Fork 52
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
NatsSub simplification and JetStream support (#108)
* NatsSubX simplification and JetStream support * Remove INatsSubBuilder interfaces: This is to reduce unnecessary complexity since NatsSubBase can be used as the contract in SubscriptionManager * Create a separation of public interfaces and internal handling for subscriptions by using INatsSub* interfaces for public methods and using NatsSubBase as the internal handling interface with internals exposed such as Pending messages. * Push responsibility of reconnections to SubscriptionManager so that not only we can renew subscriptions but in the case of JetStream consumers recovery, we can issue pull requests for example. * Fixed req-reply * Tidy-up * Tidy-up * Test flapping * Fixed warnings
- Loading branch information
Showing
23 changed files
with
231 additions
and
283 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,57 @@ | ||
using System.Threading.Channels; | ||
|
||
namespace NATS.Client.Core; | ||
|
||
public interface INatsSub : IAsyncDisposable | ||
{ | ||
/// <summary> | ||
/// Access incoming messages for your subscription. | ||
/// </summary> | ||
ChannelReader<NatsMsg> Msgs { get; } | ||
|
||
/// <summary> | ||
/// The subject name to subscribe to. | ||
/// </summary> | ||
string Subject { get; } | ||
|
||
/// <summary> | ||
/// If specified, the subscriber will join this queue group. Subscribers with the same queue group name, | ||
/// become a queue group, and only one randomly chosen subscriber of the queue group will | ||
/// consume a message each time a message is received by the queue group. | ||
/// </summary> | ||
string? QueueGroup { get; } | ||
|
||
/// <summary> | ||
/// Complete the message channel, stop timers if they were used and send an unsubscribe | ||
/// message to the server. | ||
/// </summary> | ||
/// <returns>A <see cref="ValueTask"/> that represents the asynchronous server UNSUB operation.</returns> | ||
public ValueTask UnsubscribeAsync(); | ||
} | ||
|
||
public interface INatsSub<T> : IAsyncDisposable | ||
{ | ||
/// <summary> | ||
/// Access incoming messages for your subscription. | ||
/// </summary> | ||
ChannelReader<NatsMsg<T?>> Msgs { get; } | ||
|
||
/// <summary> | ||
/// The subject name to subscribe to. | ||
/// </summary> | ||
string Subject { get; } | ||
|
||
/// <summary> | ||
/// If specified, the subscriber will join this queue group. Subscribers with the same queue group name, | ||
/// become a queue group, and only one randomly chosen subscriber of the queue group will | ||
/// consume a message each time a message is received by the queue group. | ||
/// </summary> | ||
string? QueueGroup { get; } | ||
|
||
/// <summary> | ||
/// Complete the message channel, stop timers if they were used and send an unsubscribe | ||
/// message to the server. | ||
/// </summary> | ||
/// <returns>A <see cref="ValueTask"/> that represents the asynchronous server UNSUB operation.</returns> | ||
public ValueTask UnsubscribeAsync(); | ||
} |
This file was deleted.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.