Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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;
});

Expand All @@ -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;
});

Expand All @@ -47,13 +44,9 @@ public static IObservable<T> FromSubscribe<T>(INatsConnection connection, string
{
try
{
await foreach (var msg in connection.SubscribeAsync<T>(subject, cancellationToken: ct)
.ConfigureAwait(false))
{
observer.OnNext(msg.Data!);
}

observer.OnCompleted();
await NatsProtocol
.SubscribeAsync<T>(connection, subject, observer.OnNext, observer.OnCompleted, ct)
.ConfigureAwait(false);
}
catch (Exception ex) when (ex is not OperationCanceledException)
{
Expand All @@ -71,11 +64,6 @@ public static IObservable<TResponse> FromRequest<TRequest, TResponse>(
TRequest request,
CancellationToken cancellationToken = default) =>
Observable.FromAsync(async ct =>
{
using var linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, ct);
var reply = await connection
.RequestAsync<TRequest, TResponse>(subject, request, cancellationToken: linked.Token)
.ConfigureAwait(false);
return reply.Data ?? throw new InvalidOperationException("NATS request returned null payload.");
});
await NatsProtocol.RequestAsync<TRequest, TResponse>(connection, subject, request, cancellationToken, ct)
.ConfigureAwait(false));
}
33 changes: 7 additions & 26 deletions Observables.Nats/Observables.Nats/NatsObservable.cs
Original file line number Diff line number Diff line change
Expand Up @@ -22,17 +22,7 @@ public static Observable<Unit> 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;
});

Expand All @@ -47,8 +37,7 @@ public static Observable<Unit> FromPublish<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 Unit.Default;
});

Expand All @@ -61,12 +50,9 @@ public static Observable<T> FromSubscribe<T>(INatsConnection connection, string
{
try
{
await foreach (var msg in connection.SubscribeAsync<T>(subject, cancellationToken: ct).ConfigureAwait(false))
{
observer.OnNext(msg.Data!);
}

observer.OnCompleted();
await NatsProtocol
.SubscribeAsync<T>(connection, subject, observer.OnNext, observer.OnCompleted, ct)
.ConfigureAwait(false);
}
catch (Exception ex) when (ex is not OperationCanceledException)
{
Expand All @@ -84,11 +70,6 @@ public static Observable<TResponse> FromRequest<TRequest, TResponse>(
TRequest request,
CancellationToken cancellationToken = default) =>
Observable.FromAsync(async ct =>
{
using var linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, ct);
var reply = await connection
.RequestAsync<TRequest, TResponse>(subject, request, cancellationToken: linked.Token)
.ConfigureAwait(false);
return reply.Data ?? throw new InvalidOperationException("NATS request returned null payload.");
});
await NatsProtocol.RequestAsync<TRequest, TResponse>(connection, subject, request, cancellationToken, ct)
.ConfigureAwait(false));
}
90 changes: 90 additions & 0 deletions Observables.Nats/Observables.Nats/NatsProtocol.cs
Original file line number Diff line number Diff line change
@@ -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<byte> 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<T>(
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<T>(
INatsConnection connection,
string subject,
Action<T> onNext,
Action onCompleted,
CancellationToken cancellationToken)
{
await foreach (var msg in connection.SubscribeAsync<T>(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<TResponse> RequestAsync<TRequest, TResponse>(
INatsConnection connection,
string subject,
TRequest request,
CancellationToken userToken,
CancellationToken pumpToken)
{
using var linked = CancellationTokenSource.CreateLinkedTokenSource(userToken, pumpToken);
var reply = await connection
.RequestAsync<TRequest, TResponse>(subject, request, cancellationToken: linked.Token)
.ConfigureAwait(false);
return reply.Data ?? throw new InvalidOperationException("NATS request returned null payload.");
}
}
4 changes: 4 additions & 0 deletions Observables.Nats/Observables.Nats/Observables.Nats.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,10 @@

<Import Project="..\..\eng\NetStandard20.props" Condition="Exists('..\..\eng\NetStandard20.props')" />

<ItemGroup>
<InternalsVisibleTo Include="Observables.Nats.Reactive" />
</ItemGroup>

<ItemGroup>
<PackageReference Include="NATS.Client.Core" />
<PackageReference Include="NATS.Client.Serializers.Json" />
Expand Down
Loading