slachiewicz commented on code in PR #3977:
URL: https://github.com/apache/thrift/pull/3977#discussion_r4152236983
##########
lib/cpp/test/TServerSocketTest.cpp:
##########
@@ -20,9 +20,18 @@
#include <boost/test/unit_test.hpp>
#include <thrift/transport/TSocket.h>
#include <thrift/transport/TServerSocket.h>
+#include <chrono>
#include <memory>
+#include <thread>
+#include <vector>
#include "TTransportCheckThrow.h"
#include <iostream>
+#ifndef _WIN32
+#include <cerrno>
+#include <fcntl.h>
+#include <sys/resource.h>
Review Comment:
Guarded with `HAVE_SYS_RESOURCE_H` in 5534aca43.
##########
lib/cpp/src/thrift/transport/TServerSocket.cpp:
##########
@@ -94,6 +98,64 @@ void destroyer_of_fine_sockets(THRIFT_SOCKET* ssock) {
delete ssock;
}
+namespace {
+
+// accept() failed because of the one connection it tried to take: the peer
went away before it
+// was accepted or, on Linux, a firewall rule refused it or a network error
arrived that accept(2)
+// says to treat like EAGAIN. The listening socket is fine, so the next
connection can be accepted.
+bool isConnectionError(int err) {
+ switch (err) {
+ case THRIFT_EINTR:
+ case THRIFT_EAGAIN:
+#if THRIFT_EWOULDBLOCK != THRIFT_EAGAIN
Review Comment:
5534aca43 moves the errno classification into `TServerSocketErrors.h` and
tests each case. The call site in `acceptImpl()` is not tested, which would
need a hook around `accept()`.
##########
lib/cpp/src/thrift/transport/TServerSocket.cpp:
##########
@@ -650,60 +712,94 @@ shared_ptr<TTransport> TServerSocket::acceptImpl() {
int maxEintrs = 5;
int numEintrs = 0;
- while (true) {
- std::memset(fds, 0, sizeof(fds));
- fds[0].fd = serverSocket_;
- fds[0].events = THRIFT_POLLIN;
- if (interruptSockReader_ != THRIFT_INVALID_SOCKET) {
- fds[1].fd = interruptSockReader_;
- fds[1].events = THRIFT_POLLIN;
- }
- /*
- TODO: if THRIFT_EINTR is received, we'll restart the timeout.
- To be accurate, we need to fix this in the future.
- */
- int ret = THRIFT_POLL(fds, 2, accTimeout_);
-
- if (ret < 0) {
- // error cases
- if (THRIFT_GET_SOCKET_ERROR == THRIFT_EINTR && (numEintrs++ <
maxEintrs)) {
- // THRIFT_EINTR needs to be handled manually and we can tolerate
- // a certain number
- continue;
+ struct sockaddr_storage clientAddress;
+ int size = 0;
+ THRIFT_SOCKET clientSocket = THRIFT_INVALID_SOCKET;
+ int backoffMs = 0;
+
+ while (clientSocket == THRIFT_INVALID_SOCKET) {
+ while (true) {
+ std::memset(fds, 0, sizeof(fds));
+ fds[0].fd = serverSocket_;
+ fds[0].events = THRIFT_POLLIN;
+ if (interruptSockReader_ != THRIFT_INVALID_SOCKET) {
+ fds[1].fd = interruptSockReader_;
+ fds[1].events = THRIFT_POLLIN;
}
- int errno_copy = THRIFT_GET_SOCKET_ERROR;
- TOutput::instance().perror("TServerSocket::acceptImpl() THRIFT_POLL() ",
errno_copy);
- throw TTransportException(TTransportException::UNKNOWN, "Unknown",
errno_copy);
- } else if (ret > 0) {
- // Check for an interrupt signal
- if (interruptSockReader_ != THRIFT_INVALID_SOCKET && (fds[1].revents &
THRIFT_POLLIN)) {
- int8_t buf;
- if (-1 == recv(interruptSockReader_, cast_sockopt(&buf),
sizeof(int8_t), 0)) {
- TOutput::instance().perror("TServerSocket::acceptImpl() recv()
interrupt ",
- THRIFT_GET_SOCKET_ERROR);
+ /*
+ TODO: if THRIFT_EINTR is received, we'll restart the timeout.
+ To be accurate, we need to fix this in the future.
+ */
+ int ret = THRIFT_POLL(fds, 2, accTimeout_);
+
+ if (ret < 0) {
+ // error cases
+ if (THRIFT_GET_SOCKET_ERROR == THRIFT_EINTR && (numEintrs++ <
maxEintrs)) {
+ // THRIFT_EINTR needs to be handled manually and we can tolerate
+ // a certain number
+ continue;
+ }
+ int errno_copy = THRIFT_GET_SOCKET_ERROR;
+ TOutput::instance().perror("TServerSocket::acceptImpl() THRIFT_POLL()
", errno_copy);
+ throw TTransportException(TTransportException::UNKNOWN, "Unknown",
errno_copy);
+ } else if (ret > 0) {
+ // Check for an interrupt signal
+ if (interruptSockReader_ != THRIFT_INVALID_SOCKET && (fds[1].revents &
THRIFT_POLLIN)) {
+ int8_t buf;
+ if (-1 == recv(interruptSockReader_, cast_sockopt(&buf),
sizeof(int8_t), 0)) {
+ TOutput::instance().perror("TServerSocket::acceptImpl() recv()
interrupt ",
+ THRIFT_GET_SOCKET_ERROR);
+ }
+ throw TTransportException(TTransportException::INTERRUPTED);
}
- throw TTransportException(TTransportException::INTERRUPTED);
- }
- // Check for the actual server socket being ready
- if (fds[0].revents & THRIFT_POLLIN) {
- break;
+ // Check for the actual server socket being ready
+ if (fds[0].revents & THRIFT_POLLIN) {
+ break;
+ }
+ } else {
+ TOutput::instance()("TServerSocket::acceptImpl() THRIFT_POLL 0");
+ throw TTransportException(TTransportException::UNKNOWN);
}
- } else {
- TOutput::instance()("TServerSocket::acceptImpl() THRIFT_POLL 0");
- throw TTransportException(TTransportException::UNKNOWN);
}
- }
- struct sockaddr_storage clientAddress;
- int size = sizeof(clientAddress);
- THRIFT_SOCKET clientSocket
- = ::accept(serverSocket_, (struct sockaddr*)&clientAddress,
(socklen_t*)&size);
+ size = sizeof(clientAddress);
+ clientSocket = ::accept(serverSocket_, (struct sockaddr*)&clientAddress,
(socklen_t*)&size);
+ if (clientSocket != THRIFT_INVALID_SOCKET) {
+ break;
+ }
- if (clientSocket == THRIFT_INVALID_SOCKET) {
int errno_copy = THRIFT_GET_SOCKET_ERROR;
- TOutput::instance().perror("TServerSocket::acceptImpl() ::accept() ",
errno_copy);
- throw TTransportException(TTransportException::UNKNOWN, "accept()",
errno_copy);
+ if (isConnectionError(errno_copy)) {
+ // Only that connection is lost; wait for the next one. Not logged,
because a peer can
+ // cause this as often as it likes.
+ continue;
+ }
+ if (!isResourceExhaustion(errno_copy)) {
+ TOutput::instance().perror("TServerSocket::acceptImpl() ::accept() ",
errno_copy);
+ throw TTransportException(TTransportException::UNKNOWN, "accept()",
errno_copy);
+ }
+
+ // Out of descriptors or memory: back off before trying again, but stay
interruptible. An
+ // interrupt that arrives meanwhile is left for the poll above to report.
Every retry polls
+ // again with the full accept timeout, so the whole call can take longer
than that timeout.
+ // Logged at most once a minute: connections that come and go keep
bringing it back.
+ std::chrono::steady_clock::time_point now =
std::chrono::steady_clock::now();
+ if (now >= nextExhaustionLog_) {
+ nextExhaustionLog_ = now + std::chrono::minutes(1);
+ TOutput::instance().perror("TServerSocket::acceptImpl() ::accept()
waiting to retry ",
+ errno_copy);
+ }
+ backoffMs = (backoffMs == 0) ? 5 : (std::min)(backoffMs * 2, 1000);
+ if (interruptSockReader_ != THRIFT_INVALID_SOCKET) {
+ struct THRIFT_POLLFD interruptFd;
+ std::memset(&interruptFd, 0, sizeof(interruptFd));
+ interruptFd.fd = interruptSockReader_;
+ interruptFd.events = THRIFT_POLLIN;
+ THRIFT_POLL(&interruptFd, 1, backoffMs);
Review Comment:
5534aca43 moves the backoff step into `nextBackoffMs()` and tests its growth
and the 1 s cap, without timing.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]