diff --git a/src/NATS.Client.JetStream/Models/StreamConfig.cs b/src/NATS.Client.JetStream/Models/StreamConfig.cs index e1b2a16c8..66a06a752 100644 --- a/src/NATS.Client.JetStream/Models/StreamConfig.cs +++ b/src/NATS.Client.JetStream/Models/StreamConfig.cs @@ -270,4 +270,11 @@ public StreamConfig() [System.Text.Json.Serialization.JsonPropertyName("allow_msg_counter")] [System.Text.Json.Serialization.JsonIgnore(Condition = System.Text.Json.Serialization.JsonIgnoreCondition.WhenWritingDefault)] public bool AllowMsgCounter { get; set; } + + /// + /// AllowAtomicPublish allows atomic batch publishing into the stream. + /// + [System.Text.Json.Serialization.JsonPropertyName("allow_atomic")] + [System.Text.Json.Serialization.JsonIgnore(Condition = System.Text.Json.Serialization.JsonIgnoreCondition.WhenWritingDefault)] + public bool AllowAtomicPublish { get; set; } } diff --git a/tests/NATS.Client.JetStream.Tests/ManageStreamTest.cs b/tests/NATS.Client.JetStream.Tests/ManageStreamTest.cs index 2f938df3a..56e287def 100644 --- a/tests/NATS.Client.JetStream.Tests/ManageStreamTest.cs +++ b/tests/NATS.Client.JetStream.Tests/ManageStreamTest.cs @@ -2,6 +2,7 @@ using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Models; using NATS.Client.Platform.Windows.Tests; +using NATS.Client.TestUtilities; using NATS.Client.TestUtilities2; namespace NATS.Client.JetStream.Tests; @@ -219,4 +220,42 @@ public async Task Create_or_update_stream_should_be_throwing_update_operation_er await js.CreateOrUpdateStreamAsync(streamConfig, cts.Token); await Assert.ThrowsAsync(async () => await js.CreateOrUpdateStreamAsync(streamConfigForUpdated, cts.Token)); } + + [SkipIfNatsServer(versionEarlierThan: "2.12")] + public async Task AllowAtomicPublish_property_should_be_set_on_stream() + { + await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + var prefix = _server.GetNextId(); + await nats.ConnectRetryAsync(); + + var js = new NatsJSContext(nats); + + var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); + + // Create a stream with AllowAtomicPublish enabled + var streamConfig = new StreamConfig($"{prefix}atomic", [$"{prefix}atomic.*"]) + { + AllowAtomicPublish = true, + }; + + var stream = await js.CreateStreamAsync(streamConfig, cts.Token); + + // Verify the property is set on the created stream + Assert.True(stream.Info.Config.AllowAtomicPublish); + + // Get the stream and verify the property is persisted + var retrievedStream = await js.GetStreamAsync($"{prefix}atomic", cancellationToken: cts.Token); + Assert.True(retrievedStream.Info.Config.AllowAtomicPublish); + + // Update stream with AllowAtomicPublish disabled + var updatedConfig = streamConfig with { AllowAtomicPublish = false }; + var updatedStream = await js.UpdateStreamAsync(updatedConfig, cts.Token); + + // Verify the property is updated + Assert.False(updatedStream.Info.Config.AllowAtomicPublish); + + // Get the stream and verify the update is persisted + var reRetrievedStream = await js.GetStreamAsync($"{prefix}atomic", cancellationToken: cts.Token); + Assert.False(reRetrievedStream.Info.Config.AllowAtomicPublish); + } }