From 5e1eac5b3ee2cb6bc568847d29fcee934e714029 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=90=BD=E7=AC=94wys?= <46271592+Skymly@users.noreply.github.com> Date: Sun, 30 Aug 2026 17:25:24 +0800 Subject: [PATCH 1/3] Concentrate NATS publish/subscribe/request behind thin adapters. R3 and Reactive cloned the NATS pump. Keep R3 FromPublish(ReadOnlyMemory) as the existing public-surface divergence. Closes #254 --- .../SystemReactiveNatsAdapter.cs | 28 ++---- .../Observables.Nats/NatsObservable.cs | 33 ++----- .../Observables.Nats/NatsProtocol.cs | 91 +++++++++++++++++++ .../Observables.Nats/Observables.Nats.csproj | 4 + 4 files changed, 110 insertions(+), 46 deletions(-) create mode 100644 Observables.Nats/Observables.Nats/NatsProtocol.cs 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..60ad8fc --- /dev/null +++ b/Observables.Nats/Observables.Nats/NatsProtocol.cs @@ -0,0 +1,91 @@ +using NATS.Client.Core; +#if NET8_0_OR_GREATER +using System.Diagnostics.CodeAnalysis; +#endif + +namespace Observables.Nats; + +/// Feature protocol bridge for NATS (publish / subscribe / request). +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 @@ + + + + From 1100ed2b7b69af39c97a6a65af68aadb2376f8c3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=90=BD=E7=AC=94wys?= <46271592+Skymly@users.noreply.github.com> Date: Sun, 30 Aug 2026 17:34:55 +0800 Subject: [PATCH 2/3] Omit XML docs on internal protocol types so loc parity ignores them. --- Observables.Nats/Observables.Nats/NatsProtocol.cs | 2 -- 1 file changed, 2 deletions(-) diff --git a/Observables.Nats/Observables.Nats/NatsProtocol.cs b/Observables.Nats/Observables.Nats/NatsProtocol.cs index 60ad8fc..0229a08 100644 --- a/Observables.Nats/Observables.Nats/NatsProtocol.cs +++ b/Observables.Nats/Observables.Nats/NatsProtocol.cs @@ -4,8 +4,6 @@ #endif namespace Observables.Nats; - -/// Feature protocol bridge for NATS (publish / subscribe / request). internal static class NatsProtocol { internal static async Task PublishEmptyAsync( From 6b69f05ec470e61a6559188261e4faa14f8f3a66 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=90=BD=E7=AC=94wys?= <46271592+Skymly@users.noreply.github.com> Date: Sun, 30 Aug 2026 17:41:48 +0800 Subject: [PATCH 3/3] Fix whitespace to satisfy CI format verify. --- Observables.Nats/Observables.Nats/NatsProtocol.cs | 1 + 1 file changed, 1 insertion(+) diff --git a/Observables.Nats/Observables.Nats/NatsProtocol.cs b/Observables.Nats/Observables.Nats/NatsProtocol.cs index 0229a08..a671000 100644 --- a/Observables.Nats/Observables.Nats/NatsProtocol.cs +++ b/Observables.Nats/Observables.Nats/NatsProtocol.cs @@ -4,6 +4,7 @@ #endif namespace Observables.Nats; + internal static class NatsProtocol { internal static async Task PublishEmptyAsync(