diff --git a/src/NATS.Client.Core/NatsConnection.RequestReply.cs b/src/NATS.Client.Core/NatsConnection.RequestReply.cs index a3576cbf5..e82fd5589 100644 --- a/src/NATS.Client.Core/NatsConnection.RequestReply.cs +++ b/src/NATS.Client.Core/NatsConnection.RequestReply.cs @@ -62,8 +62,13 @@ public async IAsyncEnumerable> RequestManyAsync /// NATS message structure as defined by the protocol. /// @@ -12,6 +20,7 @@ namespace NATS.Client.Core; /// Pass additional information using name-value pairs. /// Serializable data object. /// NATS connection this message is associated to. +/// Message flags. /// Specifies the type of data that may be sent to the NATS Server. /// /// Connection property is used to provide reply functionality. @@ -28,9 +37,12 @@ public readonly record struct NatsMsg( 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 Build( string subject, @@ -41,9 +53,13 @@ internal static NatsMsg Build( NatsHeaderParser headerParser, INatsDeserialize 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; @@ -58,6 +74,11 @@ internal static NatsMsg Build( throw new NatsException("Error parsing headers"); } + if (payloadBuffer.Length == 0 && headers is { Count: 0, Code: 503 }) + { + flags |= NatsMsgFlags.NoResponders; + } + headers.SetReadOnly(); } @@ -66,7 +87,7 @@ internal static NatsMsg Build( + (headersBuffer?.Length ?? 0) + payloadBuffer.Length; - return new NatsMsg(subject, replyTo, (int)size, headers, data, connection); + return new NatsMsg(subject, replyTo, (int)size, headers, data, connection, flags); } /// diff --git a/tests/NATS.Client.Core.Tests/NatsMsgTest.cs b/tests/NATS.Client.Core.Tests/NatsMsgTest.cs new file mode 100644 index 000000000..7cc2394bf --- /dev/null +++ b/tests/NATS.Client.Core.Tests/NatsMsgTest.cs @@ -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.Build( + subject: "foo", + replyTo: "bar", + headersBuffer: null, + payloadBuffer: new ReadOnlySequence(new byte[] { }), + connection: null, + headerParser: new NatsHeaderParser(Encoding.UTF8), + NatsDefaultSerializer.Default); + Assert.True(msg1.HasNoPayload); + + var msg2 = NatsMsg.Build( + subject: "foo", + replyTo: "bar", + headersBuffer: null, + payloadBuffer: new ReadOnlySequence(new[] { (byte)'0' }), + connection: null, + headerParser: new NatsHeaderParser(Encoding.UTF8), + NatsDefaultSerializer.Default); + Assert.False(msg2.HasNoPayload); + } + + [Fact] + public void No_responders() + { + var msg1 = NatsMsg.Build( + subject: "foo", + replyTo: "bar", + headersBuffer: new ReadOnlySequence(Encoding.UTF8.GetBytes("NATS/1.0 503\r\n\r\n")), + payloadBuffer: new ReadOnlySequence(new byte[] { }), + connection: null, + headerParser: new NatsHeaderParser(Encoding.UTF8), + NatsDefaultSerializer.Default); + Assert.True(msg1.IsNoRespondersError); + + var msg2 = NatsMsg.Build( + subject: "foo", + replyTo: "bar", + headersBuffer: new ReadOnlySequence(Encoding.UTF8.GetBytes("NATS/1.0 503\r\n\r\n")), + payloadBuffer: new ReadOnlySequence(new[] { (byte)'0' }), + connection: null, + headerParser: new NatsHeaderParser(Encoding.UTF8), + NatsDefaultSerializer.Default); + Assert.False(msg2.IsNoRespondersError); + + var msg3 = NatsMsg.Build( + subject: "foo", + replyTo: "bar", + headersBuffer: new ReadOnlySequence(Encoding.UTF8.GetBytes("NATS/1.0 503\r\nk: v\r\n\r\n")), + payloadBuffer: new ReadOnlySequence(new byte[] { }), + connection: null, + headerParser: new NatsHeaderParser(Encoding.UTF8), + NatsDefaultSerializer.Default); + Assert.False(msg3.IsNoRespondersError); + } +}