diff --git a/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs b/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs index 815e2dba6..29d696538 100644 --- a/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs +++ b/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs @@ -165,6 +165,8 @@ public ValueTask CallMsgNextAsync(string origin, ConsumerGetnextRequest request, cancellationToken: cancellationToken); } + public void StopHeartbeatTimer() => _timer.Change(Timeout.Infinite, Timeout.Infinite); + public void ResetHeartbeatTimer() => _timer.Change(_hbTimeout, _hbTimeout); public void Delivered(int msgSize) diff --git a/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs b/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs index c881329db..98e218360 100644 --- a/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs +++ b/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs @@ -137,6 +137,16 @@ public ValueTask CallMsgNextAsync(ConsumerGetnextRequest request, CancellationTo serializer: NatsJSJsonSerializer.Default, cancellationToken: cancellationToken); + public void StopHeartbeatTimer() + { + // if we don't have an idle timeout, we don't need to reset the timer + // because we don't expect any heartbeats. + if (_idle == TimeSpan.Zero) + return; + + _hbTimer.Change(Timeout.Infinite, Timeout.Infinite); + } + public void ResetHeartbeatTimer() { // if we don't have an idle timeout, we don't need to reset the timer diff --git a/src/NATS.Client.JetStream/NatsJSConsumer.cs b/src/NATS.Client.JetStream/NatsJSConsumer.cs index ca122936c..c36c15b4b 100644 --- a/src/NATS.Client.JetStream/NatsJSConsumer.cs +++ b/src/NATS.Client.JetStream/NatsJSConsumer.cs @@ -95,7 +95,11 @@ public async IAsyncEnumerable> ConsumeAsync( if (!read) break; + // if yield is blocked by the application code, we don't want + // heartbeat timer kicking in and issuing unnecessary pull requests. + cc.StopHeartbeatTimer(); yield return jsMsg; + cc.ResetHeartbeatTimer(); cc.Delivered(jsMsg.Size); } } @@ -202,7 +206,11 @@ public async IAsyncEnumerable> FetchAsync( if (!read) break; + // if yield is blocked by the application code, we don't want + // heartbeat timer kicking in and issuing unnecessary pull requests. + fc.StopHeartbeatTimer(); yield return jsMsg; + fc.ResetHeartbeatTimer(); } } } @@ -263,7 +271,11 @@ public async IAsyncEnumerable> FetchNoWaitAsync( await using var fc = await FetchInternalAsync(opts with { NoWait = true }, serializer, cancellationToken).ConfigureAwait(false); await foreach (var jsMsg in fc.Msgs.ReadAllAsync(cancellationToken).ConfigureAwait(false)) { + // if yield is blocked by the application code, we don't want + // heartbeat timer kicking in and issuing unnecessary pull requests. + fc.StopHeartbeatTimer(); yield return jsMsg; + fc.ResetHeartbeatTimer(); } }