diff --git a/Observables.Nats/Observables.Nats.Reactive/SystemReactiveNatsAdapter.cs b/Observables.Nats/Observables.Nats.Reactive/SystemReactiveNatsAdapter.cs index c53b55d..f487970 100644 --- a/Observables.Nats/Observables.Nats.Reactive/SystemReactiveNatsAdapter.cs +++ b/Observables.Nats/Observables.Nats.Reactive/SystemReactiveNatsAdapter.cs @@ -1,6 +1,6 @@ using System.Reactive.Linq; -using System.Threading; using NATS.Client.Core; +using Observables.Nats; #if NET8_0_OR_GREATER using System.Diagnostics.CodeAnalysis; #endif @@ -16,9 +16,7 @@ public static class SystemReactiveNatsAdapter CancellationToken cancellationToken = default) => Observable.FromAsync(async ct => { - using var linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, ct); - await connection.PublishAsync(subject, string.Empty, cancellationToken: linked.Token) - .ConfigureAwait(false); + await NatsProtocol.PublishEmptyAsync(connection, subject, cancellationToken, ct).ConfigureAwait(false); return System.Reactive.Unit.Default; }); @@ -33,8 +31,7 @@ await connection.PublishAsync(subject, string.Empty, cancellationToken: linked.T CancellationToken cancellationToken = default) => Observable.FromAsync(async ct => { - using var linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, ct); - await connection.PublishAsync(subject, payload, cancellationToken: linked.Token).ConfigureAwait(false); + await NatsProtocol.PublishAsync(connection, subject, payload, cancellationToken, ct).ConfigureAwait(false); return System.Reactive.Unit.Default; }); @@ -47,13 +44,9 @@ public static IObservable FromSubscribe(INatsConnection connection, string { try { - await foreach (var msg in connection.SubscribeAsync(subject, cancellationToken: ct) - .ConfigureAwait(false)) - { - observer.OnNext(msg.Data!); - } - - observer.OnCompleted(); + await NatsProtocol + .SubscribeAsync(connection, subject, observer.OnNext, observer.OnCompleted, ct) + .ConfigureAwait(false); } catch (Exception ex) when (ex is not OperationCanceledException) { @@ -71,11 +64,6 @@ public static IObservable FromRequest( TRequest request, CancellationToken cancellationToken = default) => Observable.FromAsync(async ct => - { - using var linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, ct); - var reply = await connection - .RequestAsync(subject, request, cancellationToken: linked.Token) - .ConfigureAwait(false); - return reply.Data ?? throw new InvalidOperationException("NATS request returned null payload."); - }); + await NatsProtocol.RequestAsync(connection, subject, request, cancellationToken, ct) + .ConfigureAwait(false)); } diff --git a/Observables.Nats/Observables.Nats/NatsObservable.cs b/Observables.Nats/Observables.Nats/NatsObservable.cs index ba3ef09..bb82c3b 100644 --- a/Observables.Nats/Observables.Nats/NatsObservable.cs +++ b/Observables.Nats/Observables.Nats/NatsObservable.cs @@ -22,17 +22,7 @@ public static Observable FromPublish( CancellationToken cancellationToken = default) => Observable.FromAsync(async ct => { - using var linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, ct); - if (payload.IsEmpty) - { - await connection.PublishAsync(subject, string.Empty, cancellationToken: linked.Token) - .ConfigureAwait(false); - } - else - { - await connection.PublishAsync(subject, payload, cancellationToken: linked.Token).ConfigureAwait(false); - } - + await NatsProtocol.PublishBytesAsync(connection, subject, payload, cancellationToken, ct).ConfigureAwait(false); return Unit.Default; }); @@ -47,8 +37,7 @@ public static Observable FromPublish( CancellationToken cancellationToken = default) => Observable.FromAsync(async ct => { - using var linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, ct); - await connection.PublishAsync(subject, payload, cancellationToken: linked.Token).ConfigureAwait(false); + await NatsProtocol.PublishAsync(connection, subject, payload, cancellationToken, ct).ConfigureAwait(false); return Unit.Default; }); @@ -61,12 +50,9 @@ public static Observable FromSubscribe(INatsConnection connection, string { try { - await foreach (var msg in connection.SubscribeAsync(subject, cancellationToken: ct).ConfigureAwait(false)) - { - observer.OnNext(msg.Data!); - } - - observer.OnCompleted(); + await NatsProtocol + .SubscribeAsync(connection, subject, observer.OnNext, observer.OnCompleted, ct) + .ConfigureAwait(false); } catch (Exception ex) when (ex is not OperationCanceledException) { @@ -84,11 +70,6 @@ public static Observable FromRequest( TRequest request, CancellationToken cancellationToken = default) => Observable.FromAsync(async ct => - { - using var linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, ct); - var reply = await connection - .RequestAsync(subject, request, cancellationToken: linked.Token) - .ConfigureAwait(false); - return reply.Data ?? throw new InvalidOperationException("NATS request returned null payload."); - }); + await NatsProtocol.RequestAsync(connection, subject, request, cancellationToken, ct) + .ConfigureAwait(false)); } diff --git a/Observables.Nats/Observables.Nats/NatsProtocol.cs b/Observables.Nats/Observables.Nats/NatsProtocol.cs new file mode 100644 index 0000000..a671000 --- /dev/null +++ b/Observables.Nats/Observables.Nats/NatsProtocol.cs @@ -0,0 +1,90 @@ +using NATS.Client.Core; +#if NET8_0_OR_GREATER +using System.Diagnostics.CodeAnalysis; +#endif + +namespace Observables.Nats; + +internal static class NatsProtocol +{ + internal static async Task PublishEmptyAsync( + INatsConnection connection, + string subject, + CancellationToken userToken, + CancellationToken pumpToken) + { + using var linked = CancellationTokenSource.CreateLinkedTokenSource(userToken, pumpToken); + await connection.PublishAsync(subject, string.Empty, cancellationToken: linked.Token).ConfigureAwait(false); + } + + internal static async Task PublishBytesAsync( + INatsConnection connection, + string subject, + ReadOnlyMemory payload, + CancellationToken userToken, + CancellationToken pumpToken) + { + using var linked = CancellationTokenSource.CreateLinkedTokenSource(userToken, pumpToken); + if (payload.IsEmpty) + { + await connection.PublishAsync(subject, string.Empty, cancellationToken: linked.Token).ConfigureAwait(false); + } + else + { + await connection.PublishAsync(subject, payload, cancellationToken: linked.Token).ConfigureAwait(false); + } + } + +#if NET8_0_OR_GREATER + [RequiresUnreferencedCode("NATS payload serialization may use reflection. Preserve payload type members when trimming.")] + [RequiresDynamicCode("NATS payload serialization may use reflection.")] +#endif + internal static async Task PublishAsync( + INatsConnection connection, + string subject, + T payload, + CancellationToken userToken, + CancellationToken pumpToken) + { + using var linked = CancellationTokenSource.CreateLinkedTokenSource(userToken, pumpToken); + await connection.PublishAsync(subject, payload, cancellationToken: linked.Token).ConfigureAwait(false); + } + +#if NET8_0_OR_GREATER + [RequiresUnreferencedCode("NATS payload serialization may use reflection. Preserve payload type members when trimming.")] + [RequiresDynamicCode("NATS payload serialization may use reflection.")] +#endif + internal static async Task SubscribeAsync( + INatsConnection connection, + string subject, + Action onNext, + Action onCompleted, + CancellationToken cancellationToken) + { + await foreach (var msg in connection.SubscribeAsync(subject, cancellationToken: cancellationToken) + .ConfigureAwait(false)) + { + onNext(msg.Data!); + } + + onCompleted(); + } + +#if NET8_0_OR_GREATER + [RequiresUnreferencedCode("NATS payload serialization may use reflection. Preserve payload type members when trimming.")] + [RequiresDynamicCode("NATS payload serialization may use reflection.")] +#endif + internal static async Task RequestAsync( + INatsConnection connection, + string subject, + TRequest request, + CancellationToken userToken, + CancellationToken pumpToken) + { + using var linked = CancellationTokenSource.CreateLinkedTokenSource(userToken, pumpToken); + var reply = await connection + .RequestAsync(subject, request, cancellationToken: linked.Token) + .ConfigureAwait(false); + return reply.Data ?? throw new InvalidOperationException("NATS request returned null payload."); + } +} diff --git a/Observables.Nats/Observables.Nats/Observables.Nats.csproj b/Observables.Nats/Observables.Nats/Observables.Nats.csproj index 0fe8d53..725e948 100644 --- a/Observables.Nats/Observables.Nats/Observables.Nats.csproj +++ b/Observables.Nats/Observables.Nats/Observables.Nats.csproj @@ -8,6 +8,10 @@ + + + +