From 173b45948810bfa8f26d99f8f8014532719dd4a1 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Tue, 15 Jul 2025 13:44:46 +0100 Subject: [PATCH 1/3] Add JSON serialization with options --- ...atsJsonContextOptionsSerializerRegistry.cs | 69 +++++++++++++++++++ .../NATS.Client.Core2.Tests/SerializerTest.cs | 49 +++++++++++++ 2 files changed, 118 insertions(+) create mode 100644 src/NATS.Client.Serializers.Json/NatsJsonContextOptionsSerializerRegistry.cs diff --git a/src/NATS.Client.Serializers.Json/NatsJsonContextOptionsSerializerRegistry.cs b/src/NATS.Client.Serializers.Json/NatsJsonContextOptionsSerializerRegistry.cs new file mode 100644 index 000000000..1e730cdb7 --- /dev/null +++ b/src/NATS.Client.Serializers.Json/NatsJsonContextOptionsSerializerRegistry.cs @@ -0,0 +1,69 @@ +using System.Buffers; +using System.Text.Json; +using NATS.Client.Core; + +namespace NATS.Client.Serializers.Json; + +public class NatsJsonContextOptionsSerializerRegistry(JsonSerializerOptions options) : INatsSerializerRegistry +{ + public INatsSerialize GetSerializer() => new NatsJsonOptionsSerializer(options); + + public INatsDeserialize GetDeserializer() => new NatsJsonOptionsSerializer(options); +} + +public class NatsJsonOptionsSerializer(JsonSerializerOptions options) : INatsSerializer +{ + // ReSharper disable once StaticMemberInGenericType + private static readonly JsonWriterOptions JsonWriterOpts = new() { Indented = false, SkipValidation = true }; + + // ReSharper disable once StaticMemberInGenericType + [ThreadStatic] + private static Utf8JsonWriter? jsonWriter; + + public void Serialize(IBufferWriter bufferWriter, T value) + { + Utf8JsonWriter writer; + if (jsonWriter == null) + { + writer = jsonWriter = new Utf8JsonWriter(bufferWriter, JsonWriterOpts); + } + else + { + writer = jsonWriter; + writer.Reset(bufferWriter); + } + + JsonSerializer.Serialize(writer, value, options); + + writer.Reset(NullBufferWriter.Instance); + } + + public T? Deserialize(in ReadOnlySequence buffer) + { + if (buffer.Length == 0) + { + return default; + } + + var reader = new Utf8JsonReader(buffer); // Utf8JsonReader is ref struct, no allocate. + return JsonSerializer.Deserialize(ref reader, options); + } + + public INatsSerializer CombineWith(INatsSerializer next) + { + throw new NotSupportedException(); + } + + private sealed class NullBufferWriter : IBufferWriter + { + public static readonly IBufferWriter Instance = new NullBufferWriter(); + + public void Advance(int count) + { + } + + public Memory GetMemory(int sizeHint = 0) => Array.Empty(); + + public Span GetSpan(int sizeHint = 0) => []; + } +} diff --git a/tests/NATS.Client.Core2.Tests/SerializerTest.cs b/tests/NATS.Client.Core2.Tests/SerializerTest.cs index f63c7a746..1c976563c 100644 --- a/tests/NATS.Client.Core2.Tests/SerializerTest.cs +++ b/tests/NATS.Client.Core2.Tests/SerializerTest.cs @@ -1,6 +1,9 @@ using System.Buffers; using System.Text; +using System.Text.Json; +using System.Text.Json.Serialization; using NATS.Client.Core2.Tests; +using NATS.Client.Serializers.Json; // ReSharper disable RedundantTypeArgumentsOfMethod // ReSharper disable ReturnTypeCanBeNotNullable @@ -283,6 +286,40 @@ public async Task Deserialize_chained_with_empty() Assert.Equal("something", result2.Data.Name); } + [Fact] + public async Task Deserialize_using_json_stream_serializer_registry() + { + var jsonSerializerOptions = new JsonSerializerOptions(); + jsonSerializerOptions.TypeInfoResolverChain.Add(TestSerializerContext1.Default); + jsonSerializerOptions.TypeInfoResolverChain.Add(TestSerializerContext2.Default); + + await using var nats = new NatsConnection(new NatsOpts + { + Url = _server.Url, + SerializerRegistry = new NatsJsonContextOptionsSerializerRegistry(jsonSerializerOptions), + }); + var prefix = _server.GetNextId(); + + using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); + var cancellationToken = cts.Token; + + await nats.ConnectAsync(); + + var sub1 = await nats.SubscribeCoreAsync($"{prefix}.1", cancellationToken: cancellationToken); + var sub2 = await nats.SubscribeCoreAsync($"{prefix}.2", cancellationToken: cancellationToken); + + await nats.PublishAsync($"{prefix}.1", new TestMessage1("one"), cancellationToken: cancellationToken); + await nats.PublishAsync($"{prefix}.2", new TestMessage2("two"), cancellationToken: cancellationToken); + + var result1 = await sub1.Msgs.ReadAsync(cancellationToken); + Assert.NotNull(result1.Data); + Assert.Equal("one", result1.Data.Name); + + var result2 = await sub2.Msgs.ReadAsync(cancellationToken); + Assert.NotNull(result2.Data); + Assert.Equal("two", result2.Data.Name); + } + private static void AssertByteArray(byte[] expected, byte[] actual) { Assert.Equal(expected.Length, actual.Length); @@ -323,3 +360,15 @@ public class TestSerializerWithEmpty : INatsSerializer } public record TestData(string Name); + +public record TestMessage1(string Name); + +public record TestMessage2(string Name); + +[JsonSerializable(typeof(TestMessage1))] +[JsonSourceGenerationOptions(DefaultIgnoreCondition = JsonIgnoreCondition.WhenWritingDefault, WriteIndented = false)] +internal partial class TestSerializerContext1 : JsonSerializerContext; + +[JsonSerializable(typeof(TestMessage2))] +[JsonSourceGenerationOptions(DefaultIgnoreCondition = JsonIgnoreCondition.WhenWritingDefault, WriteIndented = false)] +internal partial class TestSerializerContext2 : JsonSerializerContext; From 9f9281f7068e9a4f4b0ffd1ddd6f42e18c679910 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Tue, 15 Jul 2025 14:24:44 +0100 Subject: [PATCH 2/3] wip --- tests/NATS.Client.Core2.Tests/SerializerTest.cs | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/tests/NATS.Client.Core2.Tests/SerializerTest.cs b/tests/NATS.Client.Core2.Tests/SerializerTest.cs index 1c976563c..fbd987e79 100644 --- a/tests/NATS.Client.Core2.Tests/SerializerTest.cs +++ b/tests/NATS.Client.Core2.Tests/SerializerTest.cs @@ -2,6 +2,7 @@ using System.Text; using System.Text.Json; using System.Text.Json.Serialization; +using System.Text.Json.Serialization.Metadata; using NATS.Client.Core2.Tests; using NATS.Client.Serializers.Json; @@ -289,9 +290,12 @@ public async Task Deserialize_chained_with_empty() [Fact] public async Task Deserialize_using_json_stream_serializer_registry() { - var jsonSerializerOptions = new JsonSerializerOptions(); - jsonSerializerOptions.TypeInfoResolverChain.Add(TestSerializerContext1.Default); - jsonSerializerOptions.TypeInfoResolverChain.Add(TestSerializerContext2.Default); + var jsonSerializerOptions = new JsonSerializerOptions + { + TypeInfoResolver = JsonTypeInfoResolver.Combine( + TestSerializerContext1.Default, + TestSerializerContext2.Default), + }; await using var nats = new NatsConnection(new NatsOpts { From a0c7592e5a7e68406f22a06ee4c8d200cbc9fdc7 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Tue, 15 Jul 2025 14:26:50 +0100 Subject: [PATCH 3/3] Fix format --- .../NatsJsonContextOptionsSerializerRegistry.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/NATS.Client.Serializers.Json/NatsJsonContextOptionsSerializerRegistry.cs b/src/NATS.Client.Serializers.Json/NatsJsonContextOptionsSerializerRegistry.cs index 1e730cdb7..c205a9213 100644 --- a/src/NATS.Client.Serializers.Json/NatsJsonContextOptionsSerializerRegistry.cs +++ b/src/NATS.Client.Serializers.Json/NatsJsonContextOptionsSerializerRegistry.cs @@ -1,4 +1,4 @@ -using System.Buffers; +using System.Buffers; using System.Text.Json; using NATS.Client.Core;