Skip to content
Closed
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
7 changes: 6 additions & 1 deletion src/NATS.Client.Core/NatsConnection.RequestReply.cs
Original file line number Diff line number Diff line change
Expand Up @@ -62,8 +62,13 @@ public async IAsyncEnumerable<NatsMsg<TReply>> RequestManyAsync<TRequest, TReply
{
while (sub.Msgs.TryRead(out var msg))
{
if (msg.IsNoRespondersError)
{
throw new NatsNoRespondersException();
}

// Received end of stream sentinel
if (msg.Data is null)
if (msg.HasNoPayload)
{
yield break;
}
Expand Down
33 changes: 27 additions & 6 deletions src/NATS.Client.Core/NatsMsg.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,14 @@

namespace NATS.Client.Core;

[Flags]
public enum NatsMsgFlags : byte
{
None = 0,
NoPayload = 1,
NoResponders = 2,
}

/// <summary>
/// NATS message structure as defined by the protocol.
/// </summary>
Expand All @@ -12,6 +20,7 @@ namespace NATS.Client.Core;
/// <param name="Headers">Pass additional information using name-value pairs.</param>
/// <param name="Data">Serializable data object.</param>
/// <param name="Connection">NATS connection this message is associated to.</param>
/// <param name="Flags">Message flags.</param>
/// <typeparam name="T">Specifies the type of data that may be sent to the NATS Server.</typeparam>
/// <remarks>
/// <para>Connection property is used to provide reply functionality.</para>
Expand All @@ -28,9 +37,12 @@ public readonly record struct NatsMsg<T>(
int Size,
NatsHeaders? Headers,
T? Data,
INatsConnection? Connection)
INatsConnection? Connection,
NatsMsgFlags Flags = NatsMsgFlags.None)
{
public bool IsNoRespondersError => Headers?.Code == 503;
public bool IsNoRespondersError => Flags.HasFlag(NatsMsgFlags.NoResponders);

public bool HasNoPayload => Flags.HasFlag(NatsMsgFlags.NoPayload);

internal static NatsMsg<T> Build(
string subject,
Expand All @@ -41,9 +53,13 @@ internal static NatsMsg<T> Build(
NatsHeaderParser headerParser,
INatsDeserialize<T> serializer)
{
// Consider an empty payload as null or default value for value types. This way we are able to
// receive sentinels as nulls or default values. This might cause an issue with where we are not
// able to differentiate between an empty sentinel and actual default value of a struct e.g. 0 (zero).
var flags = NatsMsgFlags.None;

if (payloadBuffer.Length == 0)
{
flags |= NatsMsgFlags.NoPayload;
}

var data = payloadBuffer.Length > 0
? serializer.Deserialize(payloadBuffer)
: default;
Expand All @@ -58,6 +74,11 @@ internal static NatsMsg<T> Build(
throw new NatsException("Error parsing headers");
}

if (payloadBuffer.Length == 0 && headers is { Count: 0, Code: 503 })
{
flags |= NatsMsgFlags.NoResponders;
}

headers.SetReadOnly();
}

Expand All @@ -66,7 +87,7 @@ internal static NatsMsg<T> Build(
+ (headersBuffer?.Length ?? 0)
+ payloadBuffer.Length;

return new NatsMsg<T>(subject, replyTo, (int)size, headers, data, connection);
return new NatsMsg<T>(subject, replyTo, (int)size, headers, data, connection, flags);
}

/// <summary>
Expand Down
65 changes: 65 additions & 0 deletions tests/NATS.Client.Core.Tests/NatsMsgTest.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
using System.Buffers;
using System.Text;

namespace NATS.Client.Core.Tests;

public class NatsMsgTest
{
[Fact]
public void Empty_payload()
{
var msg1 = NatsMsg<int>.Build(
subject: "foo",
replyTo: "bar",
headersBuffer: null,
payloadBuffer: new ReadOnlySequence<byte>(new byte[] { }),
connection: null,
headerParser: new NatsHeaderParser(Encoding.UTF8),
NatsDefaultSerializer<int>.Default);
Assert.True(msg1.HasNoPayload);

var msg2 = NatsMsg<int>.Build(
subject: "foo",
replyTo: "bar",
headersBuffer: null,
payloadBuffer: new ReadOnlySequence<byte>(new[] { (byte)'0' }),
connection: null,
headerParser: new NatsHeaderParser(Encoding.UTF8),
NatsDefaultSerializer<int>.Default);
Assert.False(msg2.HasNoPayload);
}

[Fact]
public void No_responders()
{
var msg1 = NatsMsg<int>.Build(
subject: "foo",
replyTo: "bar",
headersBuffer: new ReadOnlySequence<byte>(Encoding.UTF8.GetBytes("NATS/1.0 503\r\n\r\n")),
payloadBuffer: new ReadOnlySequence<byte>(new byte[] { }),
connection: null,
headerParser: new NatsHeaderParser(Encoding.UTF8),
NatsDefaultSerializer<int>.Default);
Assert.True(msg1.IsNoRespondersError);

var msg2 = NatsMsg<int>.Build(
subject: "foo",
replyTo: "bar",
headersBuffer: new ReadOnlySequence<byte>(Encoding.UTF8.GetBytes("NATS/1.0 503\r\n\r\n")),
payloadBuffer: new ReadOnlySequence<byte>(new[] { (byte)'0' }),
connection: null,
headerParser: new NatsHeaderParser(Encoding.UTF8),
NatsDefaultSerializer<int>.Default);
Assert.False(msg2.IsNoRespondersError);

var msg3 = NatsMsg<int>.Build(
subject: "foo",
replyTo: "bar",
headersBuffer: new ReadOnlySequence<byte>(Encoding.UTF8.GetBytes("NATS/1.0 503\r\nk: v\r\n\r\n")),
payloadBuffer: new ReadOnlySequence<byte>(new byte[] { }),
connection: null,
headerParser: new NatsHeaderParser(Encoding.UTF8),
NatsDefaultSerializer<int>.Default);
Assert.False(msg3.IsNoRespondersError);
}
}