Merge branch 'poll' of github.com:eclipse/paho.mqtt.c into poll

This commit is contained in:
Ian Craggs 2022-02-11 17:40:52 +00:00
commit edf99c19b2
6 changed files with 300 additions and 14 deletions

View File

@ -36,10 +36,6 @@ INCLUDE(GNUInstallDirs)
STRING(TIMESTAMP BUILD_TIMESTAMP UTC)
MESSAGE(STATUS "Timestamp is ${BUILD_TIMESTAMP}")
IF (PAHO_HIGH_PERFORMANCE)
ADD_DEFINITIONS(-DHIGH_PERFORMANCE=1)
ENDIF()
IF(WIN32)
ADD_DEFINITIONS(-D_CRT_SECURE_NO_DEPRECATE -DWIN32_LEAN_AND_MEAN -MD)
ELSEIF(${CMAKE_SYSTEM_NAME} STREQUAL "Darwin")
@ -55,6 +51,16 @@ SET(PAHO_BUILD_SAMPLES FALSE CACHE BOOL "Build sample programs")
SET(PAHO_BUILD_DEB_PACKAGE FALSE CACHE BOOL "Build debian package")
SET(PAHO_ENABLE_TESTING TRUE CACHE BOOL "Build tests and run")
SET(PAHO_ENABLE_CPACK TRUE CACHE BOOL "Enable CPack")
SET(PAHO_HIGH_PERFORMANCE FALSE CACHE BOOL "Disable tracing and heap tracking")
SET(PAHO_USE_SELECT FALSE CACHE BOOL "Revert to select system call instead of poll")
IF (PAHO_HIGH_PERFORMANCE)
ADD_DEFINITIONS(-DHIGH_PERFORMANCE=1)
ENDIF()
IF (PAHO_USE_SELECT)
ADD_DEFINITIONS(-DUSE_SELECT=1)
ENDIF()
IF (NOT PAHO_BUILD_SHARED AND NOT PAHO_BUILD_STATIC)
MESSAGE(FATAL_ERROR "You must set either PAHO_BUILD_SHARED, PAHO_BUILD_STATIC, or both")

View File

@ -1,5 +1,5 @@
/*******************************************************************************
* Copyright (c) 2009, 2013 IBM Corp.
* Copyright (c) 2009, 2022 IBM Corp.
*
* All rights reserved. This program and the accompanying materials
* are made available under the terms of the Eclipse Public License v2.0

View File

@ -1014,7 +1014,9 @@ int SSLSocket_putdatas(SSL* ssl, SOCKET socket, char* buf0, size_t buf0len, Pack
SocketBuffer_pendingWrite(socket, ssl, 1, &iovec, &free, iovec.iov_len, 0);
*sockmem = socket;
ListAppend(mod_s.write_pending, sockmem, sizeof(int));
//FD_SET(socket, &(mod_s.pending_wset));
#if defined(USE_SELECT)
FD_SET(socket, &(mod_s.pending_wset));
#endif
rc = TCPSOCKET_INTERRUPTED;
}
else

View File

@ -44,14 +44,19 @@
#include "Heap.h"
#if defined(USE_SELECT)
int isReady(int socket, fd_set* read_set, fd_set* write_set);
int Socket_continueWrites(fd_set* pwset, int* socket, mutex_type mutex);
#else
int isReady(int index);
int Socket_continueWrites(SOCKET* socket, mutex_type mutex);
#endif
int Socket_setnonblocking(SOCKET sock);
int Socket_error(char* aString, SOCKET sock);
int Socket_addSocket(SOCKET newSd);
int isReady(int index);
int Socket_writev(SOCKET socket, iobuf* iovecs, int count, unsigned long* bytes);
int Socket_close_only(SOCKET socket);
int Socket_continueWrite(SOCKET socket);
int Socket_continueWrites(SOCKET* socket, mutex_type mutex);
char* Socket_getaddrname(struct sockaddr* sa, SOCKET sock);
int Socket_abortWrite(SOCKET socket);
@ -65,6 +70,9 @@ int Socket_abortWrite(SOCKET socket);
* Structure to hold all socket data for this module
*/
Sockets mod_s;
#if defined(USE_SELECT)
static fd_set wset;
#endif
/**
* Set a socket non-blocking, OS independently
@ -135,13 +143,22 @@ void Socket_outInitialize(void)
SocketBuffer_initialize();
mod_s.connect_pending = ListInitialize();
mod_s.write_pending = ListInitialize();
#if defined(USE_SELECT)
mod_s.clientsds = ListInitialize();
mod_s.cur_clientsds = NULL;
FD_ZERO(&(mod_s.rset)); /* Initialize the descriptor set */
FD_ZERO(&(mod_s.pending_wset));
mod_s.maxfdp1 = 0;
memcpy((void*)&(mod_s.rset_saved), (void*)&(mod_s.rset), sizeof(mod_s.rset_saved));
#else
mod_s.nfds = 0;
mod_s.fds = NULL;
mod_s.saved.cur_fd = -1;
mod_s.saved.fds = NULL;
mod_s.saved.nfds = 0;
#endif
FUNC_EXIT;
}
@ -154,10 +171,14 @@ void Socket_outTerminate(void)
FUNC_ENTRY;
ListFree(mod_s.connect_pending);
ListFree(mod_s.write_pending);
#if defined(USE_SELECT)
ListFree(mod_s.clientsds);
#else
if (mod_s.fds)
free(mod_s.fds);
if (mod_s.saved.fds)
free(mod_s.saved.fds);
#endif
SocketBuffer_terminate();
#if defined(_WIN32) || defined(_WIN64)
WSACleanup();
@ -166,6 +187,54 @@ void Socket_outTerminate(void)
}
#if defined(USE_SELECT)
/**
* Add a socket to the list of socket to check with select
* @param newSd the new socket to add
*/
int Socket_addSocket(int newSd)
{
int rc = 0;
FUNC_ENTRY;
if (ListFindItem(mod_s.clientsds, &newSd, intcompare) == NULL) /* make sure we don't add the same socket twice */
{
if (mod_s.clientsds->count >= FD_SETSIZE)
{
Log(LOG_ERROR, -1, "addSocket: exceeded FD_SETSIZE %d", FD_SETSIZE);
rc = SOCKET_ERROR;
}
else
{
int* pnewSd = (int*)malloc(sizeof(newSd));
if (!pnewSd)
{
rc = PAHO_MEMORY_ERROR;
goto exit;
}
*pnewSd = newSd;
if (!ListAppend(mod_s.clientsds, pnewSd, sizeof(newSd)))
{
free(pnewSd);
rc = PAHO_MEMORY_ERROR;
goto exit;
}
FD_SET(newSd, &(mod_s.rset_saved));
mod_s.maxfdp1 = max(mod_s.maxfdp1, newSd + 1);
rc = Socket_setnonblocking(newSd);
if (rc == SOCKET_ERROR)
Log(LOG_ERROR, -1, "addSocket: setnonblocking");
}
}
else
Log(LOG_ERROR, -1, "addSocket: socket %d already in the list", newSd);
exit:
FUNC_EXIT_RC(rc);
return rc;
}
#else
static int cmpfds(const void *p1, const void *p2)
{
SOCKET key1 = ((struct pollfd*)p1)->fd;
@ -222,8 +291,31 @@ exit:
FUNC_EXIT_RC(rc);
return rc;
}
#endif
#if defined(USE_SELECT)
/**
* Don't accept work from a client unless it is accepting work back, i.e. its socket is writeable
* this seems like a reasonable form of flow control, and practically, seems to work.
* @param socket the socket to check
* @param read_set the socket read set (see select doc)
* @param write_set the socket write set (see select doc)
* @return boolean - is the socket ready to go?
*/
int isReady(int socket, fd_set* read_set, fd_set* write_set)
{
int rc = 1;
FUNC_ENTRY;
if (ListFindItem(mod_s.connect_pending, &socket, intcompare) && FD_ISSET(socket, write_set))
ListRemoveItem(mod_s.connect_pending, &socket, intcompare);
else
rc = FD_ISSET(socket, read_set) && FD_ISSET(socket, write_set) && Socket_noPendingWrites(socket);
FUNC_EXIT_RC(rc);
return rc;
}
#else
/**
* Don't accept work from a client unless it is accepting work back, i.e. its socket is writeable
* this seems like a reasonable form of flow control, and practically, seems to work.
@ -249,8 +341,113 @@ int isReady(int index)
FUNC_EXIT_RC(rc);
return rc;
}
#endif
#if defined(USE_SELECT)
/**
* Returns the next socket ready for communications as indicated by select
* @param more_work flag to indicate more work is waiting, and thus a timeout value of 0 should
* be used for the select
* @param timeout the timeout to be used for the select, unless overridden
* @param rc a value other than 0 indicates an error of the returned socket
* @return the socket next ready, or 0 if none is ready
*/
int Socket_getReadySocket(int more_work, int timeout, mutex_type mutex, int* rc)
{
int sock = 0;
*rc = 0;
int timeout_ms = 1000;
FUNC_ENTRY;
Thread_lock_mutex(mutex);
if (mod_s.clientsds->count == 0)
goto exit;
if (more_work)
timeout_ms = 0;
else if (timeout >= 0)
timeout_ms = timeout;
while (mod_s.cur_clientsds != NULL)
{
if (isReady(*((int*)(mod_s.cur_clientsds->content)), &(mod_s.rset), &wset))
break;
ListNextElement(mod_s.clientsds, &mod_s.cur_clientsds);
}
if (mod_s.cur_clientsds == NULL)
{
static struct timeval zero = {0L, 0L}; /* 0 seconds */
int rc1, maxfdp1_saved;
fd_set pwset;
struct timeval timeout_tv = {0L, timeout_ms*1000};
memcpy((void*)&(mod_s.rset), (void*)&(mod_s.rset_saved), sizeof(mod_s.rset));
memcpy((void*)&(pwset), (void*)&(mod_s.pending_wset), sizeof(pwset));
maxfdp1_saved = mod_s.maxfdp1;
if (maxfdp1_saved == 0)
{
sock = 0;
goto exit; /* no work to do */
}
/* Prevent performance issue by unlocking the socket_mutex while waiting for a ready socket. */
Thread_unlock_mutex(mutex);
*rc = select(maxfdp1_saved, &(mod_s.rset), &pwset, NULL, &timeout_tv);
Thread_lock_mutex(mutex);
if (*rc == SOCKET_ERROR)
{
Socket_error("read select", 0);
goto exit;
}
Log(TRACE_MAX, -1, "Return code %d from read select", *rc);
if (Socket_continueWrites(&pwset, &sock, mutex) == SOCKET_ERROR)
{
*rc = SOCKET_ERROR;
goto exit;
}
memcpy((void*)&wset, (void*)&(mod_s.rset_saved), sizeof(wset));
if ((rc1 = select(mod_s.maxfdp1, NULL, &(wset), NULL, &zero)) == SOCKET_ERROR)
{
Socket_error("write select", 0);
*rc = rc1;
goto exit;
}
Log(TRACE_MAX, -1, "Return code %d from write select", rc1);
if (*rc == 0 && rc1 == 0)
{
sock = 0;
goto exit; /* no work to do */
}
mod_s.cur_clientsds = mod_s.clientsds->first;
while (mod_s.cur_clientsds != NULL)
{
int cursock = *((int*)(mod_s.cur_clientsds->content));
if (isReady(cursock, &(mod_s.rset), &wset))
break;
ListNextElement(mod_s.clientsds, &mod_s.cur_clientsds);
}
}
*rc = 0;
if (mod_s.cur_clientsds == NULL)
sock = 0;
else
{
sock = *((int*)(mod_s.cur_clientsds->content));
ListNextElement(mod_s.clientsds, &mod_s.cur_clientsds);
}
exit:
Thread_unlock_mutex(mutex);
FUNC_EXIT_RC(sock);
return sock;
} /* end getReadySocket */
#else
/**
* Returns the next socket ready for communications as indicated by select
* @param more_work flag to indicate more work is waiting, and thus a timeout value of 0 should
@ -345,6 +542,7 @@ exit:
FUNC_EXIT_RC(sock);
return sock;
} /* end getReadySocket */
#endif
/**
@ -582,7 +780,9 @@ int Socket_putdatas(SOCKET socket, char* buf0, size_t buf0len, PacketBuffers buf
rc = PAHO_MEMORY_ERROR;
goto exit;
}
//FD_SET(socket, &(mod_s.pending_wset));
#if defined(USE_SELECT)
FD_SET(socket, &(mod_s.pending_wset));
#endif
rc = TCPSOCKET_INTERRUPTED;
}
}
@ -600,7 +800,9 @@ exit:
*/
void Socket_addPendingWrite(SOCKET socket)
{
//FD_SET(socket, &(mod_s.pending_wset));
#if defined(USE_SELECT)
FD_SET(socket, &(mod_s.pending_wset));
#endif
}
@ -610,8 +812,10 @@ void Socket_addPendingWrite(SOCKET socket)
*/
void Socket_clearPendingWrite(SOCKET socket)
{
/*if (FD_ISSET(socket, &(mod_s.pending_wset)))
FD_CLR(socket, &(mod_s.pending_wset));*/
#if defined(USE_SELECT)
if (FD_ISSET(socket, &(mod_s.pending_wset)))
FD_CLR(socket, &(mod_s.pending_wset));
#endif
}
@ -642,7 +846,52 @@ int Socket_close_only(SOCKET socket)
return rc;
}
#if defined(USE_SELECT)
/**
* Close a socket and remove it from the select list.
* @param socket the socket to close
* @return completion code
*/
int Socket_close(SOCKET socket)
{
int rc = 0;
FUNC_ENTRY;
Socket_close_only(socket);
FD_CLR(socket, &(mod_s.rset_saved));
if (FD_ISSET(socket, &(mod_s.pending_wset)))
FD_CLR(socket, &(mod_s.pending_wset));
if (mod_s.cur_clientsds != NULL && *(int*)(mod_s.cur_clientsds->content) == socket)
mod_s.cur_clientsds = mod_s.cur_clientsds->next;
Socket_abortWrite(socket);
SocketBuffer_cleanup(socket);
ListRemoveItem(mod_s.connect_pending, &socket, intcompare);
ListRemoveItem(mod_s.write_pending, &socket, intcompare);
if (ListRemoveItem(mod_s.clientsds, &socket, intcompare))
Log(TRACE_MIN, -1, "Removed socket %d", socket);
else
{
Log(LOG_ERROR, -1, "Failed to remove socket %d", socket);
rc = -1;
goto exit;
}
if (socket + 1 >= mod_s.maxfdp1)
{
/* now we have to reset mod_s.maxfdp1 */
ListElement* cur_clientsds = NULL;
mod_s.maxfdp1 = 0;
while (ListNextElement(mod_s.clientsds, &cur_clientsds))
mod_s.maxfdp1 = max(*((int*)(cur_clientsds->content)), mod_s.maxfdp1);
++(mod_s.maxfdp1);
Log(TRACE_MAX, -1, "Reset max fdp1 to %d", mod_s.maxfdp1);
}
exit:
FUNC_EXIT_RC(rc);
return rc;
}
#else
/**
* Close a socket and remove it from the select list.
* @param socket the socket to close
@ -692,6 +941,7 @@ exit:
FUNC_EXIT_RC(rc);
return rc;
}
#endif
/**
@ -1016,12 +1266,23 @@ exit:
}
#if defined(USE_SELECT)
/**
* Continue any outstanding writes for a socket set
* @param pwset the set of sockets
* @param sock in case of a socket error contains the affected socket
* @return completion code, 0 or SOCKET_ERROR
*/
int Socket_continueWrites(fd_set* pwset, int* sock, mutex_type mutex)
#else
/**
* Continue any outstanding socket writes
* @param sock in case of a socket error contains the affected socket
* @return completion code, 0 or SOCKET_ERROR
*/
int Socket_continueWrites(SOCKET* sock, mutex_type mutex)
#endif
{
int rc1 = 0;
ListElement* curpending = mod_s.write_pending->first;
@ -1031,15 +1292,23 @@ int Socket_continueWrites(SOCKET* sock, mutex_type mutex)
{
int socket = *(int*)(curpending->content);
int rc = 0;
#if defined(USE_SELECT)
if (FD_ISSET(socket, pwset) && ((rc = Socket_continueWrite(socket)) != 0))
#else
struct pollfd* fd;
/* find the socket in the fds structure */
fd = bsearch(&socket, mod_s.saved.fds, (size_t)mod_s.saved.nfds, sizeof(mod_s.saved.fds[0]), cmpsockfds);
if ((fd->revents & POLLOUT) && ((rc = Socket_continueWrite(socket)) != 0))
#endif
{
if (!SocketBuffer_writeComplete(socket))
Log(LOG_SEVERE, -1, "Failed to remove pending write from socket buffer list");
#if defined(USE_SELECT)
FD_CLR(socket, &(mod_s.pending_wset));
#endif
if (!ListRemove(mod_s.write_pending, curpending->content))
{
Log(LOG_SEVERE, -1, "Failed to remove pending write from list");

View File

@ -114,6 +114,14 @@ typedef struct
List* connect_pending; /**< list of sockets for which a connect is pending */
List* write_pending; /**< list of sockets for which a write is pending */
#if defined(USE_SELECT)
fd_set rset, /**< socket read set (see select doc) */
rset_saved; /**< saved socket read set */
int maxfdp1; /**< max descriptor used +1 (again see select doc) */
List* clientsds; /**< list of client socket descriptors */
ListElement* cur_clientsds; /**< current client socket descriptor (iterator) */
fd_set pending_wset; /**< socket pending write set for select */
#else
unsigned int nfds; /**< no of file descriptors for poll */
struct pollfd* fds; /**< poll read file descriptors */
@ -122,6 +130,7 @@ typedef struct
unsigned int nfds; /**< number of fds in the fds_saved array */
struct pollfd* fds;
} saved;
#endif
} Sockets;

View File

@ -6,7 +6,7 @@ rm -rf build.paho
mkdir build.paho
cd build.paho
echo "travis build dir $TRAVIS_BUILD_DIR pwd $PWD with OpenSSL root $OPENSSL_ROOT_DIR"
cmake -DPAHO_BUILD_STATIC=$PAHO_BUILD_STATIC -DPAHO_BUILD_SHARED=$PAHO_BUILD_SHARED -DCMAKE_BUILD_TYPE=Debug -DPAHO_WITH_SSL=TRUE -DOPENSSL_ROOT_DIR=$OPENSSL_ROOT_DIR -DPAHO_BUILD_DOCUMENTATION=FALSE -DPAHO_BUILD_SAMPLES=TRUE -DPAHO_HIGH_PERFORMANCE=$PAHO_HIGH_PERFORMANCE ..
cmake -DPAHO_BUILD_STATIC=$PAHO_BUILD_STATIC -DPAHO_BUILD_SHARED=$PAHO_BUILD_SHARED -DCMAKE_BUILD_TYPE=Debug -DPAHO_WITH_SSL=TRUE -DOPENSSL_ROOT_DIR=$OPENSSL_ROOT_DIR -DPAHO_BUILD_DOCUMENTATION=FALSE -DPAHO_BUILD_SAMPLES=TRUE -DPAHO_HIGH_PERFORMANCE=$PAHO_HIGH_PERFORMANCE -DPAHO_USE_SELECT=$PAHO_USE_SELECT ..
cmake --build .
python3 ../test/mqttsas.py &
ctest -VV --timeout 600