From 2529d31008b53accfc53d90ff971da61cf21d171 Mon Sep 17 00:00:00 2001 From: Masahiro Sakamoto Date: Fri, 27 Sep 2019 17:06:59 +0900 Subject: [PATCH] Fix memory leak caused by not being executed ClientConnection destructor --- pulsar-client-cpp/lib/ClientConnection.cc | 13 ++++++++++--- pulsar-client-cpp/lib/ClientImpl.cc | 1 + pulsar-client-cpp/lib/ConnectionPool.cc | 2 +- pulsar-client-cpp/lib/ConnectionPool.h | 2 +- 4 files changed, 13 insertions(+), 5 deletions(-) diff --git a/pulsar-client-cpp/lib/ClientConnection.cc b/pulsar-client-cpp/lib/ClientConnection.cc index 077f36fdf583d..bacbe5583715f 100644 --- a/pulsar-client-cpp/lib/ClientConnection.cc +++ b/pulsar-client-cpp/lib/ClientConnection.cc @@ -136,8 +136,8 @@ ClientConnection::ClientConnection(const std::string& logicalAddress, const std: serverProtocolVersion_(ProtocolVersion_MIN), maxMessageSize_(Commands::DefaultMaxMessageSize), executor_(executor), - resolver_(executor->createTcpResolver()), - socket_(executor->createSocket()), + resolver_(executor_->createTcpResolver()), + socket_(executor_->createSocket()), #if BOOST_VERSION >= 107000 strand_(boost::asio::make_strand(executor_->io_service_.get_executor())), #elif BOOST_VERSION >= 106600 @@ -222,7 +222,7 @@ ClientConnection::ClientConnection(const std::string& logicalAddress, const std: } } - tlsSocket_ = executor->createTlsSocket(socket_, ctx); + tlsSocket_ = executor_->createTlsSocket(socket_, ctx); } } @@ -1325,6 +1325,9 @@ void ClientConnection::handleConsumerStatsTimeout(const boost::system::error_cod void ClientConnection::close() { Lock lock(mutex_); + if (isClosed()) { + return; + } state_ = Disconnected; boost::system::error_code err; socket_->close(err); @@ -1385,6 +1388,10 @@ void ClientConnection::close() { if (tlsSocket_) { tlsSocket_->lowest_layer().close(); } + + if (executor_) { + executor_.reset(); + } } bool ClientConnection::isClosed() const { return state_ == Disconnected; } diff --git a/pulsar-client-cpp/lib/ClientImpl.cc b/pulsar-client-cpp/lib/ClientImpl.cc index 0eb9034902777..f77aec17b3dda 100644 --- a/pulsar-client-cpp/lib/ClientImpl.cc +++ b/pulsar-client-cpp/lib/ClientImpl.cc @@ -561,6 +561,7 @@ void ClientImpl::shutdown() { } } + pool_.close(); ioExecutorProvider_->close(); listenerExecutorProvider_->close(); partitionListenerExecutorProvider_->close(); diff --git a/pulsar-client-cpp/lib/ConnectionPool.cc b/pulsar-client-cpp/lib/ConnectionPool.cc index 124f81ff855ca..586acec4dd235 100644 --- a/pulsar-client-cpp/lib/ConnectionPool.cc +++ b/pulsar-client-cpp/lib/ConnectionPool.cc @@ -42,7 +42,7 @@ ConnectionPool::ConnectionPool(const ClientConfiguration& conf, ExecutorServiceP poolConnections_(poolConnections), mutex_() {} -ConnectionPool::~ConnectionPool() { +void ConnectionPool::close() { std::unique_lock lock(mutex_); if (poolConnections_) { for (auto cnxIt = pool_.begin(); cnxIt != pool_.end(); cnxIt++) { diff --git a/pulsar-client-cpp/lib/ConnectionPool.h b/pulsar-client-cpp/lib/ConnectionPool.h index 21011c5bdbf85..8b3504447b64a 100644 --- a/pulsar-client-cpp/lib/ConnectionPool.h +++ b/pulsar-client-cpp/lib/ConnectionPool.h @@ -36,7 +36,7 @@ class PULSAR_PUBLIC ConnectionPool { ConnectionPool(const ClientConfiguration& conf, ExecutorServiceProviderPtr executorProvider, const AuthenticationPtr& authentication, bool poolConnections = true); - ~ConnectionPool(); + void close(); /** * Get a connection from the pool.