21#include <boost/enable_shared_from_this.hpp>
22#include <boost/weak_ptr.hpp>
37using namespace boost::posix_time;
39namespace ph = std::placeholders;
44typedef std::function<void(boost::system::error_code ec,
size_t length)>
60 SocketCallback(SocketCallbackFunction socket_callback)
61 : callback_(socket_callback) {
70 void operator()(boost::system::error_code ec,
size_t length = 0) {
71 if (ec.value() == boost::asio::error::operation_aborted) {
74 callback_(ec, length);
87typedef boost::shared_ptr<ConnectionPool> ConnectionPoolPtr;
104class Connection :
public boost::enable_shared_from_this<Connection> {
117 const ConnectionPoolPtr& conn_pool,
118 const IOAddress& address,
119 const uint16_t port);
143 const bool persistent,
144 const long request_timeout,
157 bool isTransactionOngoing()
const {
164 bool isClosed()
const {
172 void isClosedByPeer();
179 bool isMySocket(
int socket_fd)
const;
196 bool checkPrematureTimeout(
const uint64_t transid);
219 void doTransactionInternal(
const WireDataPtr& request,
221 const bool persistent,
222 const long request_timeout,
232 void closeInternal();
240 void isClosedByPeerInternal();
259 bool checkPrematureTimeoutInternal(
const uint64_t transid);
278 void terminate(
const boost::system::error_code& ec,
279 const std::string& error_msg =
"");
292 void terminateInternal(
const boost::system::error_code& ec,
293 std::string error_msg =
"");
301 bool runCompleteCheck(
const boost::system::error_code& ec,
size_t length);
311 bool runCompleteCheckInternal(
const boost::system::error_code& ec,
size_t length);
316 void scheduleTimer(
const long request_timeout);
323 void doHandshake(
const uint64_t transid);
330 void doSend(
const uint64_t transid);
337 void doReceive(
const uint64_t transid);
350 const uint64_t transid,
351 const boost::system::error_code& ec);
363 const uint64_t transid,
364 const boost::system::error_code& ec);
376 void sendCallback(
const uint64_t transid,
const boost::system::error_code& ec,
385 void receiveCallback(
const uint64_t transid,
386 const boost::system::error_code& ec,
390 void timerCallback();
401 void closeCallback(
const bool clear =
false);
410 boost::weak_ptr<ConnectionPool> conn_pool_;
422 std::shared_ptr<TCPSocket<SocketCallback>> tcp_socket_;
425 std::shared_ptr<TLSSocket<SocketCallback>> tls_socket_;
437 bool current_persistent_;
440 bool current_response_complete_;
449 std::vector<uint8_t> buf_;
455 std::array<uint8_t, 32768> input_buf_;
458 uint64_t current_transid_;
467 std::atomic<bool> started_;
470 std::atomic<bool> need_handshake_;
473 std::atomic<bool> closed_;
480typedef boost::shared_ptr<Connection> ConnectionPtr;
489class ConnectionPool :
public boost::enable_shared_from_this<ConnectionPool> {
498 explicit ConnectionPool(
const IOServicePtr& io_service,
size_t max_addr_connections)
499 : io_service_(io_service), destinations_(), pool_mutex_(),
500 max_addr_connections_(max_addr_connections) {
516 void processNextRequest(
const IOAddress& address,
520 std::lock_guard<std::mutex> lk(pool_mutex_);
521 return (processNextRequestInternal(address, port, tls_context));
523 return (processNextRequestInternal(address, port, tls_context));
533 void postProcessNextRequest(
const IOAddress& address,
536 io_service_->post(std::bind(&ConnectionPool::processNextRequest,
565 void queueRequest(
const IOAddress& address,
570 const bool persistent,
571 const long request_timeout,
578 std::lock_guard<std::mutex> lk(pool_mutex_);
579 return (queueRequestInternal(address, port, tls_context,
580 request, response, persistent,
581 request_timeout, complete_check,
582 request_callback, connect_callback,
583 handshake_callback, close_callback));
585 return (queueRequestInternal(address, port, tls_context,
586 request, response, persistent,
587 request_timeout, complete_check,
588 request_callback, connect_callback,
589 handshake_callback, close_callback));
597 std::lock_guard<std::mutex> lk(pool_mutex_);
616 void closeIfOutOfBand(
int socket_fd) {
618 std::lock_guard<std::mutex> lk(pool_mutex_);
619 closeIfOutOfBandInternal(socket_fd);
621 closeIfOutOfBandInternal(socket_fd);
635 void processNextRequestInternal(
const IOAddress& address,
640 DestinationPtr destination = findDestination(address, port, tls_context);
643 destination->garbageCollectConnections();
644 if (!destination->queueEmpty()) {
647 ConnectionPtr connection = destination->getIdleConnection();
650 if (destination->connectionsFull()) {
655 connection.reset(
new Connection(io_service_,
659 destination->addConnection(connection);
664 RequestDescriptor desc = destination->popNextRequest();
665 connection->doTransaction(desc.request_,
668 desc.request_timeout_,
669 desc.complete_check_,
671 desc.connect_callback_,
672 desc.handshake_callback_,
673 desc.close_callback_);
703 void queueRequestInternal(
const IOAddress& address,
708 const bool persistent,
709 const long request_timeout,
715 ConnectionPtr connection;
717 DestinationPtr destination = findDestination(address, port, tls_context);
720 destination->garbageCollectConnections();
722 connection = destination->getIdleConnection();
725 destination = addDestination(address, port, tls_context);
729 if (destination->connectionsFull()) {
731 destination->pushRequest(RequestDescriptor(request,
744 connection.reset(
new Connection(io_service_, tls_context,
747 destination->addConnection(connection);
751 connection->doTransaction(request, response, persistent,
752 request_timeout, complete_check,
753 request_callback, connect_callback,
754 handshake_callback, close_callback);
761 void closeAllInternal() {
762 for (
auto const& destination : destinations_) {
763 destination.second->closeAllConnections();
766 destinations_.clear();
783 void closeIfOutOfBandInternal(
int socket_fd) {
784 for (
auto const& destination : destinations_) {
786 ConnectionPtr connection = destination.second->findBySocketFd(socket_fd);
788 if (!connection->isTransactionOngoing()) {
794 destination.second->closeConnection(connection);
804 struct RequestDescriptor {
822 const bool persistent,
823 const long& request_timeout,
831 persistent_(persistent),
832 request_timeout_(request_timeout),
833 complete_check_(complete_check),
835 connect_callback_(connect_callback),
836 handshake_callback_(handshake_callback),
837 close_callback_(close_callback) {
850 long request_timeout_;
869 struct DestinationDescriptor {
871 DestinationDescriptor(
const IOAddress& address,
874 : address_(address), port_(port), tls_context_(tls_context) {
883 bool operator<(
const DestinationDescriptor& other)
const {
884 return ((address_ < other.address_) ||
885 ((address_ == other.address_) && (port_ < other.port_)) ||
886 ((address_ == other.address_) && (port_ == other.port_) &&
887 (tls_context_ < other.tls_context_)));
895 const size_t QUEUE_SIZE_THRESHOLD = 2048;
897 const int QUEUE_WARN_SECS = 5;
906 Destination(IOAddress
const& address,
909 size_t max_connections)
910 : address_(address), port_(port), tls_context_(tls_context),
911 max_connections_(max_connections), connections_(), queue_(),
912 last_queue_warn_time_(min_date_time), last_queue_size_(0) {
917 closeAllConnections();
927 void addConnection(ConnectionPtr connection) {
928 if (connectionsFull()) {
930 <<
", already at maximum connections: "
931 << max_connections_);
934 connections_.push_back(connection);
941 void closeConnection(ConnectionPtr connection) {
942 for (
auto it = connections_.begin(); it != connections_.end(); ++it) {
943 if (*it == connection) {
945 connections_.erase(it);
953 void closeAllConnections() {
955 while (!queue_.empty()) {
959 for (
auto const& connection : connections_) {
963 connections_.clear();
989 void garbageCollectConnections() {
990 for (
auto it = connections_.begin(); it != connections_.end();) {
991 (*it)->isClosedByPeer();
992 if (!(*it)->isClosed()) {
995 it = connections_.erase(it);
1011 ConnectionPtr getIdleConnection() {
1012 for (
auto const& connection : connections_) {
1013 if (!connection->isTransactionOngoing() &&
1014 !connection->isClosed()) {
1015 return (connection);
1019 return (ConnectionPtr());
1028 ConnectionPtr findBySocketFd(
int socket_fd) {
1029 for (
auto const& connection : connections_) {
1030 if (connection->isMySocket(socket_fd)) {
1031 return (connection);
1035 return (ConnectionPtr());
1041 bool connectionsEmpty() {
1042 return (connections_.empty());
1048 bool connectionsFull() {
1049 return (connections_.size() >= max_connections_);
1055 size_t connectionCount() {
1056 return (connections_.size());
1062 size_t getMaxConnections()
const {
1063 return (max_connections_);
1069 bool queueEmpty()
const {
1070 return (queue_.empty());
1079 void pushRequest(RequestDescriptor
const& desc) {
1081 size_t size = queue_.size();
1084 if ((size > QUEUE_SIZE_THRESHOLD) && (size > last_queue_size_)) {
1085 ptime now = microsec_clock::universal_time();
1086 if ((now - last_queue_warn_time_) > seconds(QUEUE_WARN_SECS)) {
1092 last_queue_warn_time_ = now;
1097 last_queue_size_ = size;
1103 RequestDescriptor popNextRequest() {
1104 if (queue_.empty()) {
1105 isc_throw(InvalidOperation,
"cannot pop, queue is empty");
1108 RequestDescriptor desc = queue_.front();
1124 size_t max_connections_;
1127 std::list<ConnectionPtr> connections_;
1130 std::queue<RequestDescriptor> queue_;
1133 ptime last_queue_warn_time_;
1136 size_t last_queue_size_;
1140 typedef boost::shared_ptr<Destination> DestinationPtr;
1150 DestinationPtr addDestination(
const IOAddress& address,
1151 const uint16_t port,
1153 DestinationDescriptor desc(address, port, tls_context);
1154 DestinationPtr destination(
new Destination(address, port, tls_context,
1155 max_addr_connections_));
1156 destinations_[desc] = destination;
1157 return (destination);
1169 DestinationPtr findDestination(
const IOAddress& address,
1170 const uint16_t port,
1172 DestinationDescriptor desc(address, port, tls_context);
1173 auto it = destinations_.find(desc);
1174 if (it != destinations_.end()) {
1175 return (it->second);
1178 return (DestinationPtr());
1193 void removeDestination(
const IOAddress& address,
1194 const uint16_t port,
1196 DestinationDescriptor desc(address, port, tls_context);
1197 auto it = destinations_.find(desc);
1198 if (it != destinations_.end()) {
1199 it->second->closeAllConnections();
1200 destinations_.erase(it);
1208 std::map<DestinationDescriptor, DestinationPtr> destinations_;
1211 std::mutex pool_mutex_;
1214 size_t max_addr_connections_;
1219 const ConnectionPoolPtr& conn_pool,
1220 const IOAddress& address,
1221 const uint16_t port)
1222 : io_service_(io_service), conn_pool_(conn_pool), address_(address),
1223 port_(port), tls_context_(tls_context), tcp_socket_(), tls_socket_(),
1224 timer_(new IntervalTimer(io_service)), current_request_(),
1225 current_response_(), current_persistent_(false),
1226 current_response_complete_(false), current_complete_check_(),
1227 current_callback_(), buf_(), position_(0), input_buf_(),
1228 current_transid_(0), close_callback_(), started_(false),
1229 need_handshake_(false), closed_(false) {
1235 need_handshake_ =
true;
1239Connection::~Connection() {
1244Connection::resetState() {
1246 current_request_.reset();
1247 current_response_.reset();
1248 current_persistent_ =
false;
1249 current_response_complete_ =
false;
1254Connection::closeCallback(
const bool clear) {
1255 if (close_callback_) {
1258 close_callback_(tcp_socket_->getNative());
1259 }
else if (tls_socket_) {
1260 close_callback_(tls_socket_->getNative());
1263 "internal error: can't find a socket to close");
1276Connection::isClosedByPeer() {
1278 if (started_ || closed_) {
1283 std::lock_guard<std::mutex> lk(mutex_);
1284 isClosedByPeerInternal();
1286 isClosedByPeerInternal();
1291Connection::isClosedByPeerInternal() {
1300 if (tcp_socket_->getASIOSocket().is_open() &&
1301 !tcp_socket_->isUsable()) {
1304 tcp_socket_->close();
1306 }
else if (tls_socket_) {
1307 if (tls_socket_->getASIOSocket().is_open() &&
1308 !tls_socket_->isUsable()) {
1311 tls_socket_->close();
1314 isc_throw(Unexpected,
"internal error: can't find the sending socket");
1319Connection::doTransaction(
const WireDataPtr& request,
1321 const bool persistent,
1322 const long request_timeout,
1329 std::lock_guard<std::mutex> lk(mutex_);
1330 doTransactionInternal(request, response, persistent, request_timeout,
1331 complete_check, callback, connect_callback,
1332 handshake_callback, close_callback);
1334 doTransactionInternal(request, response, persistent, request_timeout,
1335 complete_check, callback, connect_callback,
1336 handshake_callback, close_callback);
1341Connection::doTransactionInternal(
const WireDataPtr& request,
1343 const bool persistent,
1344 const long request_timeout,
1352 current_request_ = request;
1353 current_response_ = response;
1354 current_persistent_ = persistent;
1355 current_complete_check_ = complete_check;
1356 current_callback_ = callback;
1357 handshake_callback_ = handshake_callback;
1358 close_callback_ = close_callback;
1369 size_t to_dump = request->size();
1370 bool truncated =
false;
1371 if (to_dump > 100) {
1378 (truncated ?
"..." :
""))
1383 scheduleTimer(request_timeout);
1388 TCPEndpoint endpoint(address_, port_);
1389 SocketCallback socket_cb(std::bind(&Connection::connectCallback,
1397 tcp_socket_->open(&endpoint, socket_cb);
1401 tls_socket_->open(&endpoint, socket_cb);
1406 isc_throw(Unexpected,
"internal error: can't find a socket to open");
1408 }
catch (
const std::exception& ex) {
1415Connection::close() {
1417 std::lock_guard<std::mutex> lk(mutex_);
1418 return (closeInternal());
1420 return (closeInternal());
1425Connection::closeInternal() {
1427 closeCallback(
true);
1432 tcp_socket_->close();
1435 tls_socket_->close();
1442Connection::isMySocket(
int socket_fd)
const {
1444 return (tcp_socket_->getNative() == socket_fd);
1445 }
else if (tls_socket_) {
1446 return (tls_socket_->getNative() == socket_fd);
1449 std::cerr <<
"internal error: can't find my socket\n";
1454Connection::checkPrematureTimeout(
const uint64_t transid) {
1456 std::lock_guard<std::mutex> lk(mutex_);
1457 return (checkPrematureTimeoutInternal(transid));
1459 return (checkPrematureTimeoutInternal(transid));
1464Connection::checkPrematureTimeoutInternal(
const uint64_t transid) {
1470 if (!isTransactionOngoing() || (transid != current_transid_)) {
1472 .arg(isTransactionOngoing())
1474 .arg(current_transid_);
1482Connection::terminate(
const boost::system::error_code& ec,
1483 const std::string& error_msg) {
1485 std::lock_guard<std::mutex> lk(mutex_);
1486 terminateInternal(ec, error_msg);
1488 terminateInternal(ec, error_msg);
1493Connection::terminateInternal(
const boost::system::error_code& ec,
1494 std::string error_msg) {
1496 if (isTransactionOngoing()) {
1500 tcp_socket_->cancel();
1503 tls_socket_->cancel();
1506 if (!ec && current_response_complete_) {
1507 response = current_response_;
1514 if (error_msg.empty()) {
1515 error_msg = ec.message();
1524 if (!current_response_->empty()) {
1525 size_t to_dump = current_response_->size();
1526 bool truncated =
false;
1527 if (to_dump > 100) {
1536 (truncated ?
"..." :
""));
1544 UnlockGuard<std::mutex> lock(mutex_);
1545 current_callback_(ec, response, error_msg);
1547 current_callback_(ec, response, error_msg);
1555 (!current_persistent_ || (ec == boost::asio::error::timed_out))) {
1564 ConnectionPoolPtr conn_pool = conn_pool_.lock();
1566 conn_pool->postProcessNextRequest(address_, port_, tls_context_);
1571Connection::scheduleTimer(
const long request_timeout) {
1572 if (request_timeout > 0) {
1573 timer_->setup(std::bind(&Connection::timerCallback,
this), request_timeout,
1579Connection::doHandshake(
const uint64_t transid) {
1581 if (!need_handshake_) {
1586 SocketCallback socket_cb(std::bind(&Connection::handshakeCallback,
1588 handshake_callback_,
1592 tls_socket_->handshake(socket_cb);
1595 terminate(boost::asio::error::not_connected);
1600Connection::doSend(
const uint64_t transid) {
1601 SocketCallback socket_cb(std::bind(&Connection::sendCallback,
1610 size_t remaining = buf_.size() - position_;
1612 tcp_socket_->asyncSend(&buf_[position_], remaining, socket_cb);
1617 tls_socket_->asyncSend(&buf_[position_], remaining, socket_cb);
1622 std::cerr <<
"internal error: can't find a socket to send to\n";
1624 "internal error: can't find a socket to send to");
1626 terminate(boost::asio::error::not_connected);
1631Connection::doReceive(
const uint64_t transid) {
1632 TCPEndpoint endpoint;
1633 SocketCallback socket_cb(std::bind(&Connection::receiveCallback,
1640 tcp_socket_->asyncReceive(
static_cast<void*
>(input_buf_.data()),
1641 input_buf_.size(), 0,
1642 &endpoint, socket_cb);
1646 tls_socket_->asyncReceive(
static_cast<void*
>(input_buf_.data()),
1647 input_buf_.size(), 0,
1648 &endpoint, socket_cb);
1652 std::cerr <<
"internal error: can't find a socket to receive from\n";
1654 "internal error: can't find a socket to receive from");
1657 terminate(boost::asio::error::not_connected);
1663 const uint64_t transid,
1664 const boost::system::error_code& ec) {
1665 if (checkPrematureTimeout(transid)) {
1670 if (connect_callback) {
1674 if (!connect_callback(ec, tcp_socket_->getNative())) {
1677 }
else if (tls_socket_) {
1678 if (!connect_callback(ec, tls_socket_->getNative())) {
1683 std::cerr <<
"internal error: can't find a socket to connect\n";
1687 if (ec && (ec.value() == boost::asio::error::operation_aborted)) {
1695 (ec.value() != boost::asio::error::in_progress) &&
1696 (ec.value() != boost::asio::error::already_connected)) {
1701 doHandshake(transid);
1707 const uint64_t transid,
1708 const boost::system::error_code& ec) {
1709 need_handshake_ =
false;
1710 if (checkPrematureTimeout(transid)) {
1715 if (handshake_callback) {
1719 if (!handshake_callback(ec, tls_socket_->getNative())) {
1724 std::cerr <<
"internal error: can't find TLS socket\n";
1728 if (ec && (ec.value() == boost::asio::error::operation_aborted)) {
1740Connection::sendCallback(
const uint64_t transid,
1741 const boost::system::error_code& ec,
1743 if (checkPrematureTimeout(transid)) {
1748 if (ec.value() == boost::asio::error::operation_aborted) {
1753 }
else if ((ec.value() == boost::asio::error::would_block) ||
1754 (ec.value() == boost::asio::error::try_again)) {
1765 scheduleTimer(timer_->getInterval());
1769 if (length >= buf_.size() - position_) {
1770 position_ = buf_.size();
1772 position_ += length;
1777 if (position_ == buf_.size()) {
1786Connection::receiveCallback(
const uint64_t transid,
1787 const boost::system::error_code& ec,
1789 if (checkPrematureTimeout(transid)) {
1794 if (ec.value() == boost::asio::error::operation_aborted) {
1800 if ((ec.value() != boost::asio::error::try_again) &&
1801 (ec.value() != boost::asio::error::would_block)) {
1813 scheduleTimer(timer_->getInterval());
1815 if (runCompleteCheck(ec, length)) {
1821Connection::runCompleteCheck(
const boost::system::error_code& ec,
size_t length) {
1823 std::lock_guard<std::mutex> lk(mutex_);
1824 return (runCompleteCheckInternal(ec, length));
1826 return (runCompleteCheckInternal(ec, length));
1831Connection::runCompleteCheckInternal(
const boost::system::error_code& ec,
1835 current_response_->insert(current_response_->end(),
1837 input_buf_.begin() + length);
1842 std::string err =
"";
1843 if (current_complete_check_) {
1844 status = current_complete_check_(current_response_, err);
1846 err =
"Internal error: no completion checker?";
1850 }
else if (status > 0) {
1852 current_response_complete_ =
true;
1853 terminateInternal(ec);
1856 terminateInternal(ec, err);
1863Connection::timerCallback() {
1865 terminate(boost::asio::error::timed_out);
1900 bool defer_thread_start =
false)
1901 : thread_pool_size_(thread_pool_size), thread_pool_() {
1902 if (thread_pool_size_ > 0) {
1904 thread_io_service_.reset(
new IOService());
1908 conn_pool_.reset(
new ConnectionPool(thread_io_service_, thread_pool_size_));
1912 defer_thread_start));
1916 .arg(thread_pool_size_);
1920 conn_pool_.reset(
new ConnectionPool(io_service, 1));
1938 thread_pool_->checkPausePermissions();
1945 thread_pool_->run();
1957 thread_pool_->stop();
1960 if (thread_io_service_) {
1961 thread_io_service_->stopAndPoll();
1962 thread_io_service_->stop();
1971 if (!thread_pool_) {
1976 thread_pool_->pause();
1984 if (!thread_pool_) {
1989 thread_pool_->run();
1998 return (thread_pool_->isRunning());
2010 return (thread_pool_->isStopped());
2022 return (thread_pool_->isPaused());
2033 return (thread_io_service_);
2040 return (thread_pool_size_);
2047 if (!thread_pool_) {
2050 return (thread_pool_->getThreadCount());
2059 size_t thread_pool_size_;
2070 size_t thread_pool_size,
bool defer_thread_start) {
2071 if (!multi_threading_enabled && thread_pool_size) {
2073 "TcpClient thread_pool_size must be zero "
2074 "when Kea core multi-threading is disabled");
2078 defer_thread_start));
2087 const uint16_t port,
2091 const bool persistent,
2102 if (request->empty()) {
2110 if (!complete_check) {
2114 if (!request_callback) {
2118 impl_->conn_pool_->queueRequest(address, port, tls_context,
2119 request, response, persistent,
2122 request_callback, connect_callback,
2123 handshake_callback, close_callback);
2128 return (impl_->conn_pool_->closeIfOutOfBand(socket_fd));
2138 impl_->checkPermissions();
2158 return (impl_->getThreadIOService());
2163 return (impl_->getThreadPoolSize());
2168 return (impl_->getThreadCount());
2173 return (impl_->isRunning());
2178 return (impl_->isStopped());
2183 return (impl_->isPaused());
A generic exception that is thrown if a function is called in a prohibited way.
The IOAddress class represents an IP addresses (version agnostic).
std::string toText() const
Convert the address to a string.
The IOService class is a wrapper for the ASIO io_context class.
Implements a pausable pool of IOService driven threads.
The TCPSocket class is a concrete derived class of IOAsioSocket that represents a TCP socket.
The TLSSocket class is a concrete derived class of IOAsioSocket that represents a TLS socket.
A generic error raised by the TcpClient class.
TcpClient implementation.
void resume()
Resumes running the client's thread pool.
uint16_t getThreadPoolSize()
Fetches the maximum size of the thread pool.
void start()
Starts running the client's thread pool, if multi-threaded.
ConnectionPoolPtr conn_pool_
Holds a pointer to the connection pool.
bool isPaused()
Indicates if the thread pool is paused.
void checkPermissions()
Check if the current thread can perform thread pool state transition.
uint16_t getThreadCount()
Fetches the number of threads in the pool.
bool isRunning()
Indicates if the thread pool is running.
void pause()
Pauses the client's thread pool.
TcpClientImpl(const IOServicePtr &io_service, size_t thread_pool_size=0, bool defer_thread_start=false)
Constructor.
asiolink::IOServicePtr getThreadIOService()
Fetches the internal IOService used in multi-threaded mode.
void stop()
Close all connections, and if multi-threaded, stops the client's thread pool.
bool isStopped()
Indicates if the thread pool is stopped.
~TcpClientImpl()
Destructor.
std::function< int(const WireDataPtr &, std::string &)> CompleteCheck
Completion check type.
void closeIfOutOfBand(int socket_fd)
Closes a connection if it has an out-of-band socket event.
void start()
Starts running the client's thread pool, if multi-threaded.
uint16_t getThreadCount() const
Fetches the number of threads in the pool.
std::function< void(const boost::system::error_code &, const WireDataPtr &, const std::string &)> RequestHandler
Callback type used in call to TcpClient::asyncSendRequest.
void stop()
Halts client-side IO activity.
bool isPaused()
Indicates if the thread pool is paused.
uint16_t getThreadPoolSize() const
Fetches the maximum size of the thread pool.
void checkPermissions()
Check if the current thread can perform thread pool state transition.
void asyncSendRequest(const asiolink::IOAddress &address, const uint16_t port, const asiolink::TlsContextPtr &tls_context, const WireDataPtr &request, const WireDataPtr &response, const bool persistent, const CompleteCheck &complete_check, const RequestHandler &request_callback, const RequestTimeout &request_timeout=RequestTimeout(10000), const ConnectHandler &connect_callback=ConnectHandler(), const HandshakeHandler &handshake_callback=HandshakeHandler(), const CloseHandler &close_callback=CloseHandler())
Queues new asynchronous TCP request for a given address.
bool isRunning()
Indicates if the thread pool is running.
std::function< bool(const boost::system::error_code &, const int)> ConnectHandler
Optional handler invoked when client connects to the server.
void pause()
Pauses the client's thread pool.
bool isStopped()
Indicates if the thread pool is stopped.
TcpClient(const asiolink::IOServicePtr &io_service, bool multi_threading_enabled, size_t thread_pool_size=0, bool defer_thread_start=false)
Constructor.
std::function< bool(const boost::system::error_code &, const int)> HandshakeHandler
Optional handler invoked when client performs the TLS handshake with the server.
void resume()
Resumes running the client's thread pool.
const asiolink::IOServicePtr getThreadIOService() const
Fetches a pointer to the internal IOService used to drive the thread-pool in multi-threaded mode.
std::function< void(const int)> CloseHandler
Optional handler invoked when client closes the connection to the server.
static MultiThreadingMgr & instance()
Returns a single instance of Multi Threading Manager.
#define isc_throw(type, stream)
A shortcut macro to insert known values into exception arguments.
#define LOG_ERROR(LOGGER, MESSAGE)
Macro to conveniently test error output and log it.
#define LOG_WARN(LOGGER, MESSAGE)
Macro to conveniently test warn output and log it.
#define LOG_DEBUG(LOGGER, LEVEL, MESSAGE)
Macro to conveniently test debug output and log it.
boost::shared_ptr< TlsContext > TlsContextPtr
The type of shared pointers to TlsContext objects.
boost::shared_ptr< IoServiceThreadPool > IoServiceThreadPoolPtr
Defines a pointer to a thread pool.
boost::shared_ptr< isc::asiolink::IntervalTimer > IntervalTimerPtr
boost::shared_ptr< IOService > IOServicePtr
Defines a smart pointer to an IOService instance.
bool operator<(Element const &a, Element const &b)
Test less than.
const int DBGLVL_TRACE_BASIC
Trace basic operations.
const int DBGLVL_TRACE_BASIC_DATA
Trace data associated with the basic operations.
const int DBGLVL_TRACE_DETAIL
Trace detailed operations.
std::function< void(boost::system::error_code ec, size_t length)> SocketCallbackFunction
Type of the function implementing a callback invoked by the SocketCallback functor.
const isc::log::MessageID TCP_CLIENT_PREMATURE_CONNECTION_TIMEOUT_OCCURRED
const isc::log::MessageID TCP_CLIENT_BAD_SERVER_RESPONSE_RECEIVED
const isc::log::MessageID TCP_CLIENT_BAD_SERVER_RESPONSE_RECEIVED_DETAILS
const isc::log::MessageID TCP_CLIENT_SERVER_RESPONSE_RECEIVED
isc::log::Logger tcp_logger("tcp")
Defines the logger used within libkea-tcp library.
boost::shared_ptr< WireData > WireDataPtr
const isc::log::MessageID TCP_CLIENT_QUEUE_SIZE_GROWING
const isc::log::MessageID TCP_CLIENT_CONNECTION_CLOSE_CALLBACK_FAILED
const isc::log::MessageID TCP_CLIENT_REQUEST_SEND
const isc::log::MessageID TCP_CLIENT_MT_STARTED
string dumpAsHex(const uint8_t *data, size_t length)
Dumps a buffer of bytes as a string of hexadecimal digits.
Defines the logger used by the top-level component of kea-lfc.
TCP request/response timeout value.
long value_
Timeout value specified.