Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions src/NATS.Client.JetStream/INatsJSConsumer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ public interface INatsJSConsumer
/// <typeparam name="T">Message type to deserialize.</typeparam>
/// <returns>Async enumerable of messages which can be used in a <c>await foreach</c> loop.</returns>
/// <exception cref="NatsJSProtocolException">Consumer is deleted, it's push based or request sent to server is invalid.</exception>
IAsyncEnumerable<NatsJSMsg<T>> ConsumeAsync<T>(
IAsyncEnumerable<INatsJSMsg<T>> ConsumeAsync<T>(
INatsDeserialize<T>? serializer = default,
NatsJSConsumeOpts? opts = default,
CancellationToken cancellationToken = default);
Expand Down Expand Up @@ -56,7 +56,7 @@ IAsyncEnumerable<NatsJSMsg<T>> ConsumeAsync<T>(
/// }
/// </code>
/// </example>
ValueTask<NatsJSMsg<T>?> NextAsync<T>(INatsDeserialize<T>? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default);
ValueTask<INatsJSMsg<T>?> NextAsync<T>(INatsDeserialize<T>? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default);

/// <summary>
/// Consume a set number of messages from the stream using this consumer.
Expand All @@ -68,7 +68,7 @@ IAsyncEnumerable<NatsJSMsg<T>> ConsumeAsync<T>(
/// <returns>Async enumerable of messages which can be used in a <c>await foreach</c> loop.</returns>
/// <exception cref="NatsJSProtocolException">Consumer is deleted, it's push based or request sent to server is invalid.</exception>
/// <exception cref="NatsJSException">There is an error sending the message or this consumer object isn't valid anymore because it was deleted earlier.</exception>
IAsyncEnumerable<NatsJSMsg<T>> FetchAsync<T>(
IAsyncEnumerable<INatsJSMsg<T>> FetchAsync<T>(
NatsJSFetchOpts opts,
INatsDeserialize<T>? serializer = default,
CancellationToken cancellationToken = default);
Expand Down Expand Up @@ -126,7 +126,7 @@ IAsyncEnumerable<NatsJSMsg<T>> FetchAsync<T>(
/// }
/// </code>
/// </example>
IAsyncEnumerable<NatsJSMsg<T>> FetchNoWaitAsync<T>(
IAsyncEnumerable<INatsJSMsg<T>> FetchNoWaitAsync<T>(
NatsJSFetchOpts opts,
INatsDeserialize<T>? serializer = default,
CancellationToken cancellationToken = default);
Expand Down
8 changes: 4 additions & 4 deletions src/NATS.Client.JetStream/NatsJSConsumer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ public async ValueTask<bool> DeleteAsync(CancellationToken cancellationToken = d
/// <exception cref="NatsJSProtocolException">Consumer is deleted, it's push based or request sent to server is invalid.</exception>
/// <exception cref="NatsConnectionFailedException">Connection has permanently failed and cannot be recovered.</exception>
/// <exception cref="NatsJSException">Consumer-related errors, such as the consumer being deleted after too many consecutive 503 errors.</exception>
public async IAsyncEnumerable<NatsJSMsg<T>> ConsumeAsync<T>(
public async IAsyncEnumerable<INatsJSMsg<T>> ConsumeAsync<T>(
INatsDeserialize<T>? serializer = default,
NatsJSConsumeOpts? opts = default,
[EnumeratorCancellation] CancellationToken cancellationToken = default)
Expand Down Expand Up @@ -158,7 +158,7 @@ public async IAsyncEnumerable<NatsJSMsg<T>> ConsumeAsync<T>(
/// }
/// </code>
/// </example>
public async ValueTask<NatsJSMsg<T>?> NextAsync<T>(INatsDeserialize<T>? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default)
public async ValueTask<INatsJSMsg<T>?> NextAsync<T>(INatsDeserialize<T>? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default)
{
ThrowIfDeleted();
opts ??= _context.Opts.DefaultNextOpts;
Expand All @@ -185,7 +185,7 @@ public async IAsyncEnumerable<NatsJSMsg<T>> ConsumeAsync<T>(
}

/// <inheritdoc />
public async IAsyncEnumerable<NatsJSMsg<T>> FetchAsync<T>(
public async IAsyncEnumerable<INatsJSMsg<T>> FetchAsync<T>(
NatsJSFetchOpts opts,
INatsDeserialize<T>? serializer = default,
[EnumeratorCancellation] CancellationToken cancellationToken = default)
Expand Down Expand Up @@ -282,7 +282,7 @@ public async IAsyncEnumerable<NatsJSMsg<T>> FetchAsync<T>(
/// }
/// </code>
/// </example>
public async IAsyncEnumerable<NatsJSMsg<T>> FetchNoWaitAsync<T>(
public async IAsyncEnumerable<INatsJSMsg<T>> FetchNoWaitAsync<T>(
NatsJSFetchOpts opts,
INatsDeserialize<T>? serializer = default,
[EnumeratorCancellation] CancellationToken cancellationToken = default)
Expand Down
8 changes: 4 additions & 4 deletions src/NATS.Client.JetStream/NatsJSOrderedConsumer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ public NatsJSOrderedConsumer(string stream, NatsJSContext context, NatsJSOrdered
/// <typeparam name="T">Serialized message data type.</typeparam>
/// <returns>Asynchronous enumeration which can be used in a <c>await foreach</c> loop.</returns>
/// <exception cref="NatsJSProtocolException">There was a JetStream server error.</exception>
public async IAsyncEnumerable<NatsJSMsg<T>> ConsumeAsync<T>(
public async IAsyncEnumerable<INatsJSMsg<T>> ConsumeAsync<T>(
INatsDeserialize<T>? serializer = default,
NatsJSConsumeOpts? opts = default,
[EnumeratorCancellation] CancellationToken cancellationToken = default)
Expand Down Expand Up @@ -193,7 +193,7 @@ public async IAsyncEnumerable<NatsJSMsg<T>> ConsumeAsync<T>(
/// <param name="cancellationToken">A <see cref="CancellationToken"/> used to cancel fetch operation.</param>
/// <typeparam name="T">Serialized message data type.</typeparam>
/// <returns>Asynchronous enumeration which can be used in a <c>await foreach</c> loop.</returns>
public async IAsyncEnumerable<NatsJSMsg<T>> FetchAsync<T>(
public async IAsyncEnumerable<INatsJSMsg<T>> FetchAsync<T>(
NatsJSFetchOpts opts,
INatsDeserialize<T>? serializer = default,
[EnumeratorCancellation] CancellationToken cancellationToken = default)
Expand Down Expand Up @@ -261,7 +261,7 @@ public async IAsyncEnumerable<NatsJSMsg<T>> FetchAsync<T>(
}

/// <inheritdoc />
public async IAsyncEnumerable<NatsJSMsg<T>> FetchNoWaitAsync<T>(
public async IAsyncEnumerable<INatsJSMsg<T>> FetchNoWaitAsync<T>(
NatsJSFetchOpts opts,
INatsDeserialize<T>? serializer = default,
[EnumeratorCancellation] CancellationToken cancellationToken = default)
Expand Down Expand Up @@ -335,7 +335,7 @@ public async IAsyncEnumerable<NatsJSMsg<T>> FetchNoWaitAsync<T>(
/// <param name="cancellationToken">A <see cref="CancellationToken"/> used to cancel the underlying fetch operation.</param>
/// <typeparam name="T">Serialized message data type.</typeparam>
/// <returns>The next NATS JetStream message in order.</returns>
public async ValueTask<NatsJSMsg<T>?> NextAsync<T>(INatsDeserialize<T>? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default)
public async ValueTask<INatsJSMsg<T>?> NextAsync<T>(INatsDeserialize<T>? serializer = default, NatsJSNextOpts? opts = default, CancellationToken cancellationToken = default)
{
opts ??= _context.Opts.DefaultNextOpts;

Expand Down
2 changes: 1 addition & 1 deletion tests/NATS.Client.CheckNativeAot/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,7 @@ async Task JetStreamTests()

// Consume
var cts2 = new CancellationTokenSource(TimeSpan.FromSeconds(10));
var messages = new List<NatsJSMsg<TestData>>();
var messages = new List<INatsJSMsg<TestData>>();
await foreach (var msg in consumer.ConsumeAsync(serializer: TestDataJsonSerializer<TestData>.Default, new NatsJSConsumeOpts { MaxMsgs = 100 }, cancellationToken: cts2.Token))
{
messages.Add(msg);
Expand Down
4 changes: 2 additions & 2 deletions tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<List<NatsJSMsg<byte[]>>> GetAllMsgs()
async Task<List<INatsJSMsg<byte[]>>> GetAllMsgs()
{
var c = await s1.CreateOrderedConsumerAsync(cancellationToken: cancellationToken);
List<NatsJSMsg<byte[]>> msgs = new();
List<INatsJSMsg<byte[]>> msgs = new();
await foreach (var msg in c.ConsumeAsync<byte[]>(cancellationToken: cancellationToken))
{
msgs.Add(msg);
Expand Down
2 changes: 1 addition & 1 deletion tests/NATS.Net.DocsExamples/JetStream/ConsumePage.cs
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ public async Task Run()

{
#region consumer-next
NatsJSMsg<Order>? next = await consumer.NextAsync<Order>();
INatsJSMsg<Order>? next = await consumer.NextAsync<Order>();

if (next is { } msg)
{
Expand Down
Loading