diff --git a/src/NATS.Client.Core/Internal/NatsReadProtocolProcessor.cs b/src/NATS.Client.Core/Internal/NatsReadProtocolProcessor.cs index c0d426d15..b82498adf 100644 --- a/src/NATS.Client.Core/Internal/NatsReadProtocolProcessor.cs +++ b/src/NATS.Client.Core/Internal/NatsReadProtocolProcessor.cs @@ -297,7 +297,9 @@ await _connection.PublishToClientHandlersAsync(subject, replyTo, sid, headerSlic } catch (SocketClosedException e) { + _logger.LogDebug(NatsLogEvents.Protocol, e, "Socket closed during read loop"); _waitForInfoSignal.TrySetException(e); + _waitForPongOrErrorSignal.TrySetException(e); return; } catch (Exception ex) diff --git a/src/NATS.Client.Core/NatsConnection.cs b/src/NATS.Client.Core/NatsConnection.cs index 267dc7238..dfbf7b36f 100644 --- a/src/NATS.Client.Core/NatsConnection.cs +++ b/src/NATS.Client.Core/NatsConnection.cs @@ -496,7 +496,17 @@ private async ValueTask SetupReaderWriterAsync(bool reconnect) } // receive COMMAND response (PONG or ERROR) - await waitForPongOrErrorSignal.Task.ConfigureAwait(false); + try + { + await waitForPongOrErrorSignal.Task + .WaitAsync(Opts.ConnectTimeout) + .ConfigureAwait(false); + } + catch (TimeoutException) + { + _logger.LogDebug(NatsLogEvents.Connection, "Timeout waiting for initial pong"); + throw; + } if (reconnectTask != null) {