diff --git a/src/NATS.Client.Core/NatsHeaderParser.cs b/src/NATS.Client.Core/NatsHeaderParser.cs index e26077aa5..bfbe1f5df 100644 --- a/src/NATS.Client.Core/NatsHeaderParser.cs +++ b/src/NATS.Client.Core/NatsHeaderParser.cs @@ -119,6 +119,11 @@ private bool TryParseHeaderLine(ReadOnlySpan headerLine, NatsHeaders heade headers.Message = NatsHeaders.Messages.MessageSizeExceedsMaxBytes; headers.MessageText = NatsHeaders.MessageMessageSizeExceedsMaxBytesStr; } + else if (headerLine.SequenceEqual(NatsHeaders.MessageRequestsPendingBytes)) + { + headers.Message = NatsHeaders.Messages.RequestsPending; + headers.MessageText = NatsHeaders.MessageRequestsPendingStr; + } else { headers.Message = NatsHeaders.Messages.Text; diff --git a/src/NATS.Client.Core/NatsHeaders.cs b/src/NATS.Client.Core/NatsHeaders.cs index 1c5833c8e..15d09be45 100644 --- a/src/NATS.Client.Core/NatsHeaders.cs +++ b/src/NATS.Client.Core/NatsHeaders.cs @@ -27,6 +27,7 @@ public enum Messages NoMessages, RequestTimeout, MessageSizeExceedsMaxBytes, + RequestsPending, } // Uses C# compiler's optimization for static byte[] data @@ -58,6 +59,10 @@ public enum Messages internal static ReadOnlySpan MessageMessageSizeExceedsMaxBytes => new byte[] { 77, 101, 115, 115, 97, 103, 101, 32, 83, 105, 122, 101, 32, 69, 120, 99, 101, 101, 100, 115, 32, 77, 97, 120, 66, 121, 116, 101, 115 }; internal static readonly string MessageMessageSizeExceedsMaxBytesStr = "Message Size Exceeds MaxBytes"; + // Requests Pending + internal static ReadOnlySpan MessageRequestsPendingBytes => new byte[] { 82, 101, 113, 117, 101, 115, 116, 115, 32, 80, 101, 110, 100, 105, 110, 103 }; + internal static readonly string MessageRequestsPendingStr = "Requests Pending"; + private static readonly string[] EmptyKeys = Array.Empty(); private static readonly StringValues[] EmptyValues = Array.Empty(); diff --git a/src/NATS.Client.Core/NatsSubBase.cs b/src/NATS.Client.Core/NatsSubBase.cs index 2119dc5d7..cbee5290d 100644 --- a/src/NATS.Client.Core/NatsSubBase.cs +++ b/src/NATS.Client.Core/NatsSubBase.cs @@ -22,6 +22,7 @@ public enum NatsSubEndReason Cancelled, Exception, JetStreamError, + RequestsPending, } /// diff --git a/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs b/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs index e79781c1b..239c2c503 100644 --- a/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs +++ b/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs @@ -220,6 +220,10 @@ protected override async ValueTask ReceiveInternalAsync( { EndSubscription(NatsSubEndReason.NoMsgs); } + else if (headers is { Code: 408, Message: NatsHeaders.Messages.RequestsPending }) + { + EndSubscription(NatsSubEndReason.RequestsPending); + } else if (headers is { Code: 408, Message: NatsHeaders.Messages.RequestTimeout }) { EndSubscription(NatsSubEndReason.Timeout);