From a2c95172d14caf6bc10617f38b6efb94ba7444ce Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Fri, 21 Nov 2025 13:13:23 +0000 Subject: [PATCH] Update INatsJsConsumer to return INatsJsMsg --- src/NATS.Client.JetStream/INatsJSConsumer.cs | 8 ++++---- src/NATS.Client.JetStream/NatsJSConsumer.cs | 8 ++++---- src/NATS.Client.JetStream/NatsJSOrderedConsumer.cs | 8 ++++---- tests/NATS.Client.CheckNativeAot/Program.cs | 2 +- tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs | 4 ++-- tests/NATS.Net.DocsExamples/JetStream/ConsumePage.cs | 2 +- 6 files changed, 16 insertions(+), 16 deletions(-) diff --git a/src/NATS.Client.JetStream/INatsJSConsumer.cs b/src/NATS.Client.JetStream/INatsJSConsumer.cs index 15f2ce1d0..ac69af217 100644 --- a/src/NATS.Client.JetStream/INatsJSConsumer.cs +++ b/src/NATS.Client.JetStream/INatsJSConsumer.cs @@ -20,7 +20,7 @@ public interface INatsJSConsumer /// Message type to deserialize. /// Async enumerable of messages which can be used in a await foreach loop. /// Consumer is deleted, it's push based or request sent to server is invalid. - IAsyncEnumerable> ConsumeAsync( + IAsyncEnumerable> ConsumeAsync( INatsDeserialize? serializer = default, NatsJSConsumeOpts? opts = default, CancellationToken cancellationToken = default); @@ -56,7 +56,7 @@ IAsyncEnumerable> ConsumeAsync( /// } /// /// - ValueTask?> NextAsync(INatsDeserialize? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default); + ValueTask?> NextAsync(INatsDeserialize? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default); /// /// Consume a set number of messages from the stream using this consumer. @@ -68,7 +68,7 @@ IAsyncEnumerable> ConsumeAsync( /// Async enumerable of messages which can be used in a await foreach loop. /// Consumer is deleted, it's push based or request sent to server is invalid. /// There is an error sending the message or this consumer object isn't valid anymore because it was deleted earlier. - IAsyncEnumerable> FetchAsync( + IAsyncEnumerable> FetchAsync( NatsJSFetchOpts opts, INatsDeserialize? serializer = default, CancellationToken cancellationToken = default); @@ -126,7 +126,7 @@ IAsyncEnumerable> FetchAsync( /// } /// /// - IAsyncEnumerable> FetchNoWaitAsync( + IAsyncEnumerable> FetchNoWaitAsync( NatsJSFetchOpts opts, INatsDeserialize? serializer = default, CancellationToken cancellationToken = default); diff --git a/src/NATS.Client.JetStream/NatsJSConsumer.cs b/src/NATS.Client.JetStream/NatsJSConsumer.cs index 90fe4e24d..16b2b1c5d 100644 --- a/src/NATS.Client.JetStream/NatsJSConsumer.cs +++ b/src/NATS.Client.JetStream/NatsJSConsumer.cs @@ -56,7 +56,7 @@ public async ValueTask DeleteAsync(CancellationToken cancellationToken = d /// Consumer is deleted, it's push based or request sent to server is invalid. /// Connection has permanently failed and cannot be recovered. /// Consumer-related errors, such as the consumer being deleted after too many consecutive 503 errors. - public async IAsyncEnumerable> ConsumeAsync( + public async IAsyncEnumerable> ConsumeAsync( INatsDeserialize? serializer = default, NatsJSConsumeOpts? opts = default, [EnumeratorCancellation] CancellationToken cancellationToken = default) @@ -158,7 +158,7 @@ public async IAsyncEnumerable> ConsumeAsync( /// } /// /// - public async ValueTask?> NextAsync(INatsDeserialize? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default) + public async ValueTask?> NextAsync(INatsDeserialize? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default) { ThrowIfDeleted(); opts ??= _context.Opts.DefaultNextOpts; @@ -185,7 +185,7 @@ public async IAsyncEnumerable> ConsumeAsync( } /// - public async IAsyncEnumerable> FetchAsync( + public async IAsyncEnumerable> FetchAsync( NatsJSFetchOpts opts, INatsDeserialize? serializer = default, [EnumeratorCancellation] CancellationToken cancellationToken = default) @@ -282,7 +282,7 @@ public async IAsyncEnumerable> FetchAsync( /// } /// /// - public async IAsyncEnumerable> FetchNoWaitAsync( + public async IAsyncEnumerable> FetchNoWaitAsync( NatsJSFetchOpts opts, INatsDeserialize? serializer = default, [EnumeratorCancellation] CancellationToken cancellationToken = default) diff --git a/src/NATS.Client.JetStream/NatsJSOrderedConsumer.cs b/src/NATS.Client.JetStream/NatsJSOrderedConsumer.cs index b60f73fea..b7d3afef3 100644 --- a/src/NATS.Client.JetStream/NatsJSOrderedConsumer.cs +++ b/src/NATS.Client.JetStream/NatsJSOrderedConsumer.cs @@ -59,7 +59,7 @@ public NatsJSOrderedConsumer(string stream, NatsJSContext context, NatsJSOrdered /// Serialized message data type. /// Asynchronous enumeration which can be used in a await foreach loop. /// There was a JetStream server error. - public async IAsyncEnumerable> ConsumeAsync( + public async IAsyncEnumerable> ConsumeAsync( INatsDeserialize? serializer = default, NatsJSConsumeOpts? opts = default, [EnumeratorCancellation] CancellationToken cancellationToken = default) @@ -193,7 +193,7 @@ public async IAsyncEnumerable> ConsumeAsync( /// A used to cancel fetch operation. /// Serialized message data type. /// Asynchronous enumeration which can be used in a await foreach loop. - public async IAsyncEnumerable> FetchAsync( + public async IAsyncEnumerable> FetchAsync( NatsJSFetchOpts opts, INatsDeserialize? serializer = default, [EnumeratorCancellation] CancellationToken cancellationToken = default) @@ -261,7 +261,7 @@ public async IAsyncEnumerable> FetchAsync( } /// - public async IAsyncEnumerable> FetchNoWaitAsync( + public async IAsyncEnumerable> FetchNoWaitAsync( NatsJSFetchOpts opts, INatsDeserialize? serializer = default, [EnumeratorCancellation] CancellationToken cancellationToken = default) @@ -335,7 +335,7 @@ public async IAsyncEnumerable> FetchNoWaitAsync( /// A used to cancel the underlying fetch operation. /// Serialized message data type. /// The next NATS JetStream message in order. - public async ValueTask?> NextAsync(INatsDeserialize? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default) + public async ValueTask?> NextAsync(INatsDeserialize? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default) { opts ??= _context.Opts.DefaultNextOpts; diff --git a/tests/NATS.Client.CheckNativeAot/Program.cs b/tests/NATS.Client.CheckNativeAot/Program.cs index 26e48b337..bd65f2f2e 100644 --- a/tests/NATS.Client.CheckNativeAot/Program.cs +++ b/tests/NATS.Client.CheckNativeAot/Program.cs @@ -126,7 +126,7 @@ async Task JetStreamTests() // Consume var cts2 = new CancellationTokenSource(TimeSpan.FromSeconds(10)); - var messages = new List>(); + var messages = new List>(); await foreach (var msg in consumer.ConsumeAsync(serializer: TestDataJsonSerializer.Default, new NatsJSConsumeOpts { MaxMsgs = 100 }, cancellationToken: cts2.Token)) { messages.Add(msg); diff --git a/tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs b/tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs index 6fe6634b7..02e653720 100644 --- a/tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs +++ b/tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs @@ -572,10 +572,10 @@ public async Task Rename_object_should_perge_old_named_meta() var s1 = await js.GetStreamAsync("OBJ_b1", cancellationToken: cancellationToken); - async Task>> GetAllMsgs() + async Task>> GetAllMsgs() { var c = await s1.CreateOrderedConsumerAsync(cancellationToken: cancellationToken); - List> msgs = new(); + List> msgs = new(); await foreach (var msg in c.ConsumeAsync(cancellationToken: cancellationToken)) { msgs.Add(msg); diff --git a/tests/NATS.Net.DocsExamples/JetStream/ConsumePage.cs b/tests/NATS.Net.DocsExamples/JetStream/ConsumePage.cs index f97c336b5..4cba940c5 100644 --- a/tests/NATS.Net.DocsExamples/JetStream/ConsumePage.cs +++ b/tests/NATS.Net.DocsExamples/JetStream/ConsumePage.cs @@ -62,7 +62,7 @@ public async Task Run() { #region consumer-next - NatsJSMsg? next = await consumer.NextAsync(); + INatsJSMsg? next = await consumer.NextAsync(); if (next is { } msg) {