Skip to content
Merged
2 changes: 1 addition & 1 deletion src/NATS.Client.Core/Internal/HeaderWriter.cs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ internal class HeaderWriter

private static ReadOnlySpan<byte> ColonSpace => new[] { ByteColon, ByteSpace };

internal static long GetBytesLength(NatsHeaders headers, Encoding encoding)
internal static int GetBytesLength(NatsHeaders headers, Encoding encoding)
Comment thread
mtmk marked this conversation as resolved.
{
var len = CommandConstants.NatsHeaders10NewLine.Length;
foreach (var kv in headers)
Expand Down
2 changes: 1 addition & 1 deletion src/NATS.Client.Core/NatsHeaders.cs
Original file line number Diff line number Diff line change
Expand Up @@ -301,7 +301,7 @@ public void CopyTo(KeyValuePair<string, StringValues>[] array, int arrayIndex)
/// </summary>
/// <param name="encoding">Encoding used. Default to utf8 if not provided</param>
/// <returns>Bytes length of the header</returns>
public long GetBytesLength(Encoding? encoding = null)
public int GetBytesLength(Encoding? encoding = null)
{
// if null set to utf-8
encoding = encoding ?? Encoding.UTF8;
Expand Down
108 changes: 107 additions & 1 deletion src/NATS.Client.Core/NatsMsg.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
using System.Buffers;
using System.Diagnostics.CodeAnalysis;
using System.Xml.Linq;
using System.Text;
using NATS.Client.Core.Commands;
using NATS.Client.Core.Internal;

namespace NATS.Client.Core;
Expand Down Expand Up @@ -436,6 +437,111 @@ private void CheckReplyPreconditions()
}
}

/// <summary>
/// Builder class for creating <see cref="NatsMsg{T}" />
/// </summary>
public class NatsMsgBuilder<T>
{
/// <summary>
/// The destination subject to publish to.
/// </summary>
public string Subject { get; set; } = null!;

/// <summary>
/// The reply subject that subscribers can use to send a response back to the publisher/requester.
/// </summary>
public string? ReplyTo { get; set; }

/// <summary>
/// Pass additional information using name-value pairs.
/// </summary>
public NatsHeaders? Headers { get; set; }

/// <summary>
/// Gets or sets the message payload.
/// </summary>
public T? Data { get; set; }

/// <summary>
/// NATS connection this message is associated to.
/// </summary>
public INatsConnection? Connection { get; set; }

/// <summary>
/// Message flags to indicate no responders and empty payloads.
/// </summary>
public NatsMsgFlags Flags { get; set; } = default;

/// <summary>
/// Serializer to use for the calculate data size.
/// </summary>
/// <remarks>This results in double serialization: once for size calculation, once when publishing.</remarks>
public INatsSerializer<T>? Serializer { get; set; }

/// <summary>
/// The serializer buffer size
/// </summary>
/// <remarks>This results in double serialization: once for size calculation, once when publishing.</remarks>
public int SerializationBufferSize { get; set; } = 256;

/// <summary>
/// Encoding used. Default to utf8 if not provided
/// </summary>
public Encoding Encoding { get; set; } = Encoding.UTF8;

/// <summary>
/// Builds and returns the <see cref="NatsMsg{T}" /> instance.
/// </summary>
/// <exception cref="InvalidOperationException">Thrown if <see cref="Subject"/> is null or Empty</exception>
public NatsMsg<T> Msg
{
get
{
if (string.IsNullOrWhiteSpace(Subject))
throw new InvalidOperationException("Subject must be specified and non-empty.");
var size = 0; // Default size is 0 if not calculated

switch (Data)
{
case ReadOnlyMemory<byte> memoryData:
size = memoryData.Length;
break;
case byte[] byteArray:
size = byteArray.Length; // Directly use the length of byte[]
break;
case string str:
size = Encoding.GetBytes(str).Length;
break;
default:
if (Serializer != null && Data != null)
{
var bufferWriter = new NatsPooledBufferWriter<byte>(SerializationBufferSize);
Serializer.Serialize(bufferWriter, Data);
size = bufferWriter.WrittenMemory.Length;
}

break;
}

if (size > 0)
{
size += Subject.Length
+ (ReplyTo?.Length ?? 0)
+ (Headers?.GetBytesLength() ?? 0);
}

return new NatsMsg<T>(
Subject,
ReplyTo,
size,
Headers,
Data,
Connection,
Flags);
}
}
}

public class NatsDeserializeException : NatsException
{
public NatsDeserializeException(byte[] rawData, Exception inner)
Expand Down
113 changes: 113 additions & 0 deletions tests/NATS.Client.Core.Tests/NatsMsgTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
using System.Text;
using NATS.Client.Serializers.Json;

namespace NATS.Client.Core.Tests;

public class NatsMsgTests
{
[Fact]
public void Create_nats_msg_by_string()
{
// Arrange
const string subject = "test";
const string data = "12345";
const string replyTo = "reply";
var headers = new NatsHeaders();

// Act
var builder = new NatsMsgBuilder<string>
{
Subject = subject,
Data = data,
ReplyTo = replyTo,
Headers = headers,
};
var msg = builder.Msg;

// Assert
var expectedSize = subject.Length + replyTo.Length + headers.GetBytesLength() + Encoding.UTF8.GetByteCount(data);
msg.Size.Should().Be(expectedSize);
}

[Fact]
public void Create_nats_msg_by_byte_array()
{
// Arrange
const string subject = "test";
const string replyTo = "reply";
var data = new byte[] { 1, 2, 3, 4, 5 };
var headers = new NatsHeaders();

// Act
var builder = new NatsMsgBuilder<byte[]>
{
Subject = subject,
Data = data,
ReplyTo = replyTo,
Headers = headers,
};
var msg = builder.Msg;

// Assert
var expectedSize = subject.Length + (replyTo?.Length ?? 0) + headers.GetBytesLength() + data.Length;
msg.Size.Should().Be(expectedSize);
}

[Fact]
public void Create_nats_msg_by_object()
{
// Arrange
const string subject = "test";
const string replyTo = "reply";
var data = new TestData { Id = 1, Name = "example" };
var headers = new NatsHeaders();

var serializer = new NatsJsonSerializer<TestData>();

// Act
var builder = new NatsMsgBuilder<TestData>
{
Subject = subject,
Data = data,
Serializer = serializer,
Headers = headers,
ReplyTo = replyTo,
};
var msg = builder.Msg;

var bufferWriter = new NatsPooledBufferWriter<byte>(256);
serializer.Serialize(bufferWriter, data);
var serializedSize = bufferWriter.WrittenCount;

var expectedSize = subject.Length + (replyTo?.Length ?? 0) + headers.GetBytesLength() + serializedSize;

// Assert
msg.Size.Should().Be(expectedSize);
}

[Fact]
public void Create_nats_msg_by_object_without_serializer()
{
// Arrange
const string subject = "test";
var data = new TestData { Id = 1, Name = "example" };

// Act
var builder = new NatsMsgBuilder<TestData>
{
Subject = subject,
Data = data,
};
var msg = builder.Msg;

// Assert
msg.Size.Should().Be(0);
}

private class TestData
{
public int Id { get; set; }

public string Name { get; set; } = null!;
}
}
2 changes: 1 addition & 1 deletion tests/NATS.Client.CoreUnit.Tests/NatsHeaderTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -210,7 +210,7 @@ public void GetBytesLengthTest()
var bytesLength = headers.GetBytesLength(encoding);

var text = "NATS/1.0\r\nk1: v1\r\nk2: v2-0\r\nk2: v2-1\r\na-long-header-key: value\r\nkey: a-long-header-value\r\n\r\n";
var expected = new ReadOnlySequence<byte>(Encoding.UTF8.GetBytes(text));
var expected = new ReadOnlySpan<byte>(Encoding.UTF8.GetBytes(text));

Assert.Equal(expected.Length, bytesLength);
}
Expand Down