diff --git a/src/NATS.Client.JetStream/Models/StreamConfig.cs b/src/NATS.Client.JetStream/Models/StreamConfig.cs index a5bdb2197..6c6e88c4a 100644 --- a/src/NATS.Client.JetStream/Models/StreamConfig.cs +++ b/src/NATS.Client.JetStream/Models/StreamConfig.cs @@ -301,4 +301,12 @@ public StreamConfig() [System.Text.Json.Serialization.JsonConverter(typeof(NatsJSJsonStringEnumConverter))] #endif public StreamConfigPersistMode? PersistMode { get; set; } + + /// + /// AllowBatchPublish allows fast batch publishing into the stream. + /// + /// Requires nats-server v2.14.0 or later. + [System.Text.Json.Serialization.JsonPropertyName("allow_batched")] + [System.Text.Json.Serialization.JsonIgnore(Condition = System.Text.Json.Serialization.JsonIgnoreCondition.WhenWritingDefault)] + public bool AllowBatchPublish { get; set; } } diff --git a/tests/NATS.Client.JetStream.Tests/ManageStreamTest.cs b/tests/NATS.Client.JetStream.Tests/ManageStreamTest.cs index 9a3961a59..e1e8ffaa5 100644 --- a/tests/NATS.Client.JetStream.Tests/ManageStreamTest.cs +++ b/tests/NATS.Client.JetStream.Tests/ManageStreamTest.cs @@ -297,6 +297,36 @@ public async Task AllowAtomicPublish_property_should_be_set_on_stream() Assert.False(reRetrievedStream.Info.Config.AllowAtomicPublish); } + [SkipIfNatsServer(versionEarlierThan: "2.14")] + public async Task AllowBatchPublish_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)); + + var streamConfig = new StreamConfig($"{prefix}batched", [$"{prefix}batched.*"]) + { + AllowBatchPublish = true, + }; + + var stream = await js.CreateStreamAsync(streamConfig, cts.Token); + Assert.True(stream.Info.Config.AllowBatchPublish); + + var retrievedStream = await js.GetStreamAsync($"{prefix}batched", cancellationToken: cts.Token); + Assert.True(retrievedStream.Info.Config.AllowBatchPublish); + + var updatedConfig = streamConfig with { AllowBatchPublish = false }; + var updatedStream = await js.UpdateStreamAsync(updatedConfig, cts.Token); + Assert.False(updatedStream.Info.Config.AllowBatchPublish); + + var reRetrievedStream = await js.GetStreamAsync($"{prefix}batched", cancellationToken: cts.Token); + Assert.False(reRetrievedStream.Info.Config.AllowBatchPublish); + } + [SkipIfNatsServer(versionEarlierThan: "2.12")] public async Task PersistMode_property_should_be_set_on_stream() {