/* Copyright (c) 2013 Anton Titov. Copyright (c) 2013 pCloud Ltd. All rights reserved. Redistribution and use in source and binary forms, with or without modification, are permitted provided that the following conditions are met: Redistributions of source code must retain the above copyright notice, this list of conditions and the following disclaimer. Redistributions in binary form must reproduce the above copyright notice, this list of conditions and the following disclaimer in the documentation and/or materials provided with the distribution. Neither the name of pCloud Ltd nor the names of its contributors may be used to endorse or promote products derived from this software without specific prior written permission. THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL pCloud Ltd BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. */ #include #include #include #include #include #include #include #include "plibs.h" #include "pmem.h" #include "psettings.h" #include "psock.h" #include "psql.h" #include "ptask.h" #include "ptimer.h" #define PROXY_NONE 0 #define PROXY_CONNECT 1 typedef struct { const char *host; const char *port; } resolve_host_port; static pthread_mutex_t mutex = PTHREAD_MUTEX_INITIALIZER; __attribute__((weak)) void psock_debug_log_wait_latency(const struct timespec *start) {} static int wait_readable(int sock, long sec, long usec) { fd_set rfds; struct timeval tv; struct timespec start; int res; tv.tv_sec = sec; tv.tv_usec = usec; clock_gettime(CLOCK_REALTIME, &start); do { FD_ZERO(&rfds); FD_SET(sock, &rfds); res = select(sock + 1, &rfds, NULL, NULL, &tv); } while (res == -1 && errno == EINTR); if (res == 1) { psock_debug_log_wait_latency(&start); return 0; } if (res == 0) { if (sec) pdbg_logf(D_WARNING, "socket read timeouted on %ld seconds", sec); errno = (ETIMEDOUT); } else pdbg_logf(D_WARNING, "select returned %d", res); return SOCKET_ERROR; } static int wait_writable(int sock, long sec, long usec) { fd_set wfds; struct timeval tv; int res; tv.tv_sec = sec; tv.tv_usec = usec; do { FD_ZERO(&wfds); FD_SET(sock, &wfds); res = select(sock + 1, NULL, &wfds, NULL, &tv); } while (res == -1 && errno == EINTR); if (res == 1) return 0; if (res == 0) errno = (ETIMEDOUT); return SOCKET_ERROR; } static int wait_ssl_ready(int sock) { fd_set fds; struct timeval tv; int res, want_read; if (psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ) { want_read = 1; tv.tv_sec = PSYNC_SOCK_READ_TIMEOUT; } else if (psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE) { want_read = 0; tv.tv_sec = PSYNC_SOCK_WRITE_TIMEOUT; } else { pdbg_logf(D_BUG, "this functions should only be called when SSL returns " "WANT_READ/WANT_WRITE"); errno = (EINVAL); return SOCKET_ERROR; } tv.tv_usec = 0; do { FD_ZERO(&fds); FD_SET(sock, &fds); res = select(sock + 1, want_read ? &fds : NULL, want_read ? NULL : &fds, NULL, &tv); } while (res == -1 && errno == EINTR); if (res == 1) return 0; if (res == 0) { pdbg_logf(D_WARNING, "socket timeouted"); errno = (ETIMEDOUT); } return pdbg_return_const(SOCKET_ERROR); } static int addr_valid(struct addrinfo *olda, struct addrinfo *newa) { struct addrinfo *a; do { a = newa; while (1) { if (a->ai_addrlen == olda->ai_addrlen && !memcmp(a->ai_addr, olda->ai_addr, a->ai_addrlen)) break; a = a->ai_next; if (!a) return 0; } olda = olda->ai_next; } while (olda); return 1; } static struct addrinfo *addr_load(const char *host, const char *port) { psync_sql_res *res; psync_uint_row row; psync_variant_row vrow; struct addrinfo *ret; char *data; const char *str; uint64_t i; size_t len; psql_rdlock(); res = psql_query_nolock("SELECT COUNT(*), SUM(LENGTH(data)) FROM " "resolver WHERE hostname=? AND port=?"); psql_bind_str(res, 1, host); psql_bind_str(res, 2, port); if (!(row = psql_fetch_int(res)) || row[0] == 0) { psql_free(res); psql_rdunlock(); return NULL; } size_t struct_size, total_size; if (row[0] != 0 && sizeof(struct addrinfo) > SIZE_MAX / row[0]) { psql_free(res); psql_rdunlock(); return NULL; } struct_size = sizeof(struct addrinfo) * row[0]; if (struct_size > SIZE_MAX - row[1]) { psql_free(res); psql_rdunlock(); return NULL; } total_size = struct_size + row[1]; ret = (struct addrinfo *)pmem_malloc(PMEM_SUBSYS_OTHER, total_size); data = (char *)(ret + row[0]); for (i = 0; i < row[0] - 1; i++) ret[i].ai_next = &ret[i + 1]; ret[i].ai_next = NULL; psql_free(res); res = psql_query_nolock( "SELECT family, socktype, protocol, data FROM resolver WHERE hostname=? " "AND port=? ORDER BY prio"); psql_bind_str(res, 1, host); psql_bind_str(res, 2, port); i = 0; while ((vrow = psql_fetch(res))) { ret[i].ai_family = psync_get_snumber(vrow[0]); ret[i].ai_socktype = psync_get_snumber(vrow[1]); ret[i].ai_protocol = psync_get_snumber(vrow[2]); str = psync_get_lstring(vrow[3], &len); ret[i].ai_addr = (struct sockaddr *)data; ret[i].ai_addrlen = len; i++; memcpy(data, str, len); data += len; } psql_free(res); psql_rdunlock(); return ret; } static void addr_save(const char *host, const char *port, struct addrinfo *addr) { psync_sql_res *res; uint64_t id; if (psql_rdlocked()) { if (psql_tryupgradeLock()) return; else pdbg_logf(D_NOTICE, "upgraded read to write lock to save data to DB"); } psql_start(); res = psql_prepare( "DELETE FROM resolver WHERE hostname=? AND port=?"); psql_bind_str(res, 1, host); psql_bind_str(res, 2, port); psql_run_free(res); res = psql_prepare( "INSERT INTO resolver (hostname, port, prio, created, family, socktype, " "protocol, data) VALUES (?, ?, ?, ?, ?, ?, ?, ?)"); psql_bind_str(res, 1, host); psql_bind_str(res, 2, port); psql_bind_uint(res, 4, ptimer_time()); id = 0; do { psql_bind_uint(res, 3, id++); psql_bind_int(res, 5, addr->ai_family); psql_bind_int(res, 6, addr->ai_socktype); psql_bind_int(res, 7, addr->ai_protocol); psql_bind_blob(res, 8, (char *)addr->ai_addr, addr->ai_addrlen); psql_run(res); addr = addr->ai_next; } while (addr); psql_free(res); psql_commit(); } static int psock_check_so_error(int sock) { int err; socklen_t len = sizeof(err); return (getsockopt(sock, SOL_SOCKET, SO_ERROR, &err, &len) == 0 && err == 0) ? 0 : -1; } static int connect_res(struct addrinfo *res) { int sock; #if defined(SOCK_NONBLOCK) #if defined(SOCK_CLOEXEC) #define PSOCK_TYPE_OR (SOCK_NONBLOCK | SOCK_CLOEXEC) #else #define PSOCK_TYPE_OR SOCK_NONBLOCK #endif #else #define PSOCK_TYPE_OR 0 #define PSOCK_NEED_NOBLOCK #endif while (res) { sock = socket(res->ai_family, res->ai_socktype | PSOCK_TYPE_OR, res->ai_protocol); if (pdbg_likely(sock != INVALID_SOCKET)) { #if defined(PSOCK_NEED_NOBLOCK) fcntl(sock, F_SETFD, FD_CLOEXEC); fcntl(sock, F_SETFL, fcntl(sock, F_GETFL) | O_NONBLOCK); #endif if ((connect(sock, res->ai_addr, res->ai_addrlen) != SOCKET_ERROR) || (errno == EINPROGRESS && !wait_writable(sock, PSYNC_SOCK_CONNECT_TIMEOUT, 0) && !psock_check_so_error(sock))) return sock; close(sock); } res = res->ai_next; } return INVALID_SOCKET; } static void cb_connect_res(void *h, void *ptr) { struct addrinfo *res; int sock; int r; res = (struct addrinfo *)ptr; sock = connect_res(res); r = psync_task_complete(h, (void *)(uintptr_t)sock); pmem_free(PMEM_SUBSYS_OTHER, res); if (r && sock != INVALID_SOCKET) close(sock); } static void cb_resolve(void *h, void *ptr) { resolve_host_port *hp; struct addrinfo *res; struct addrinfo hints; int rc; hp = (resolve_host_port *)ptr; memset(&hints, 0, sizeof(hints)); hints.ai_family = AF_UNSPEC; hints.ai_socktype = SOCK_STREAM; res = NULL; rc = getaddrinfo(hp->host, hp->port, &hints, &res); if (unlikely(rc != 0)) res = NULL; psync_task_complete(h, res); } static int connect_socket(const char *host, const char *port) { struct addrinfo *res, *dbres; struct addrinfo hints; int sock; int rc; pdbg_logf(D_NOTICE, "connecting to %s:%s", host, port); dbres = addr_load(host, port); if (dbres) { resolve_host_port resolv; void *params[2]; psync_task_callback_t callbacks[2]; psync_task_manager_t tasks; resolv.host = host; resolv.port = port; params[0] = dbres; params[1] = &resolv; callbacks[0] = cb_connect_res; callbacks[1] = cb_resolve; tasks = psync_task_run_tasks(callbacks, params, 2); res = (struct addrinfo *)psync_task_papi_result(tasks, 1); if (unlikely(!res)) { psync_task_free(tasks); pdbg_logf(D_WARNING, "failed to resolve %s", host); return INVALID_SOCKET; } addr_save(host, port, res); if (addr_valid(dbres, res)) { pdbg_logf(D_NOTICE, "successfully reused cached IP for %s:%s", host, port); sock = (int)(uintptr_t)psync_task_papi_result(tasks, 0); } else { pdbg_logf(D_NOTICE, "cached IP not valid for %s:%s", host, port); sock = connect_res(res); } freeaddrinfo(res); psync_task_free(tasks); } else { memset(&hints, 0, sizeof(hints)); hints.ai_family = AF_UNSPEC; hints.ai_socktype = SOCK_STREAM; res = NULL; rc = getaddrinfo(host, port, &hints, &res); if (unlikely(rc != 0)) { pdbg_logf(D_WARNING, "failed to resolve %s", host); return INVALID_SOCKET; } addr_save(host, port, res); sock = connect_res(res); freeaddrinfo(res); } if (likely(sock != INVALID_SOCKET)) { int sock_opt = 1; setsockopt(sock, SOL_TCP, TCP_NODELAY, (char *)&sock_opt, sizeof(sock_opt)); setsockopt(sock, SOL_SOCKET, SO_KEEPALIVE, (char *)&sock_opt, sizeof(sock_opt)); #if defined(SOL_TCP) #if defined(TCP_KEEPCNT) sock_opt = 3; setsockopt(sock, SOL_TCP, TCP_KEEPCNT, (char *)&sock_opt, sizeof(sock_opt)); #endif #if defined(TCP_KEEPIDLE) sock_opt = 60; setsockopt(sock, SOL_TCP, TCP_KEEPIDLE, (char *)&sock_opt, sizeof(sock_opt)); #endif #if defined(TCP_KEEPINTVL) sock_opt = 20; setsockopt(sock, SOL_TCP, TCP_KEEPINTVL, (char *)&sock_opt, sizeof(sock_opt)); #endif #endif } else { pdbg_logf(D_WARNING, "failed to connect to %s:%s", host, port); } return sock; } int psock_try_write_buffer(psock_t *sock) { if (sock->buffer) { psock_buf_t *b; int wrt, cw; wrt = 0; while ((b = sock->buffer)) { if (b->roffset == b->woffset) { sock->buffer = b->next; pmem_free(PMEM_SUBSYS_OTHER, b); continue; } if (sock->ssl) { cw = pssl_write(sock->ssl, b->buff + b->roffset, b->woffset - b->roffset); if (cw == PSYNC_SSL_FAIL) { if (pdbg_likely(psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ || psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE)) break; else { if (!wrt) wrt = -1; break; } } } else { cw = write(sock->sock, b->buff + b->roffset, b->woffset - b->roffset); if (cw == SOCKET_ERROR) { if (pdbg_likely(errno == EWOULDBLOCK || errno == EAGAIN || errno == EINTR)) break; else { if (!wrt) wrt = -1; break; } } } wrt += cw; b->roffset += cw; if (b->roffset != b->woffset) break; } if (wrt > 0) pdbg_logf(D_NOTICE, "wrote %d bytes to socket from buffers", wrt); return wrt; } else return 0; } int psock_try_write_buffer_thread(psock_t *sock) { int ret; pthread_mutex_lock(&mutex); ret = psock_try_write_buffer(sock); pthread_mutex_unlock(&mutex); return ret; } int psock_readable(psock_t *sock) { psock_try_write_buffer(sock); if (sock->ssl && pssl_pendingdata(sock->ssl)) return 1; else if (wait_readable(sock->sock, 0, 0)) return 0; else { sock->pending = 1; return 1; } } int psock_writable(psock_t *sock) { if (sock->buffer) return 1; return !wait_writable(sock->sock, 0, 0); } int psock_create(int domain, int type, int protocol) { int ret; ret = socket(domain, type, protocol); return ret; } psock_t *psock_connect(const char *host, int unsigned port, int ssl) { psock_t *ret; void *sslc; int sock; char sport[8]; putil_slprintf(sport, sizeof(sport), "%d", port); sock = connect_socket(host, sport); if (pdbg_unlikely(sock == INVALID_SOCKET)) { return NULL; } if (ssl) { ssl = pssl_connect(sock, &sslc, host); while (ssl == PSYNC_SSL_NEED_FINISH) { if (wait_ssl_ready(sock)) { pssl_free(sslc); break; } ssl = pssl_connect_finish(sslc, host); } if (pdbg_unlikely(ssl != PSYNC_SSL_SUCCESS)) { close(sock); return NULL; } } else { sslc = NULL; } ret = pmem_malloc(PMEM_SUBSYS_OTHER, sizeof(psock_t)); if (!ret) { if (sslc) pssl_free(sslc); close(sock); return NULL; } ret->ssl = sslc; ret->buffer = NULL; ret->sock = sock; ret->pending = 0; return ret; } int psock_wait_write_timeout(int sock) { return wait_writable(sock, PSYNC_SOCK_WRITE_TIMEOUT, 0); } void psock_close(psock_t *sock) { if (sock->ssl) while (pssl_shutdown(sock->ssl) == PSYNC_SSL_NEED_FINISH) if (wait_ssl_ready(sock->sock)) { pssl_free(sock->ssl); break; } psock_clear_write_buffered(sock); close(sock->sock); pmem_free(PMEM_SUBSYS_OTHER, sock); } void psock_close_bad(psock_t *sock) { if (sock->ssl) pssl_free(sock->ssl); psock_clear_write_buffered(sock); close(sock->sock); pmem_free(PMEM_SUBSYS_OTHER, sock); } void psock_set_write_buffered(psock_t *sock) { psock_buf_t *sb; if (sock->buffer) return; sb = (psock_buf_t *)pmem_malloc(PMEM_SUBSYS_OTHER, offsetof(psock_buf_t, buff) + PSYNC_FIRST_SOCK_WRITE_BUFF_SIZE); sb->next = NULL; sb->size = PSYNC_FIRST_SOCK_WRITE_BUFF_SIZE; sb->woffset = 0; sb->roffset = 0; sock->buffer = sb; } void psock_set_write_buffered_thread(psock_t *sock) { pthread_mutex_lock(&mutex); psock_set_write_buffered(sock); pthread_mutex_unlock(&mutex); } void psock_clear_write_buffered(psock_t *sock) { psock_buf_t *nb; while (sock->buffer) { nb = sock->buffer->next; pmem_free(PMEM_SUBSYS_OTHER, sock->buffer); sock->buffer = nb; } } void psock_clear_write_buffered_thread(psock_t *sock) { pthread_mutex_lock(&mutex); psock_clear_write_buffered(sock); pthread_mutex_unlock(&mutex); } int psock_set_recvbuf(psock_t *sock, int bufsize) { #if defined(SO_RCVBUF) && defined(SOL_SOCKET) return setsockopt(sock->sock, SOL_SOCKET, SO_RCVBUF, (const char *)&bufsize, sizeof(bufsize)); #else return -1; #endif } int psock_set_sendbuf(psock_t *sock, int bufsize) { #if defined(SO_SNDBUF) && defined(SOL_SOCKET) return setsockopt(sock->sock, SOL_SOCKET, SO_SNDBUF, (const char *)&bufsize, sizeof(bufsize)); #else return -1; #endif } int psock_is_ssl(psock_t *sock) { if (sock->ssl) return 1; else return 0; } int psock_pendingdata(psock_t *sock) { if (sock->pending) return 1; if (sock->ssl) return pssl_pendingdata(sock->ssl); else return 0; } int psock_pendingdata_buf(psock_t *sock) { int ret; #if defined(FIONREAD) if (ioctl(sock->sock, FIONREAD, &ret)) return -1; #else return -1; #endif if (sock->ssl) ret += pssl_pendingdata(sock->ssl); return ret; } int psock_pendingdata_buf_thread(psock_t *sock) { int ret; pthread_mutex_lock(&mutex); ret = psock_pendingdata_buf(sock); pthread_mutex_unlock(&mutex); return ret; } static int psync_socket_read_ssl(psock_t *sock, void *buff, int num) { int r; psock_try_write_buffer(sock); if (!pssl_pendingdata(sock->ssl) && !sock->pending && psock_wait_read_timeout(sock->sock)) return -1; sock->pending = 0; while (1) { psock_try_write_buffer(sock); r = pssl_read(sock->ssl, buff, num); if (r == PSYNC_SSL_FAIL) { if (pdbg_likely(psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ || psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE)) { if (wait_ssl_ready(sock->sock)) { if (sock->buffer) pdbg_logf(D_WARNING, "timeouted on socket with pending buffers"); return -1; } else continue; } else { errno = (ECONNRESET); return -1; } } else return r; } } static int psync_socket_read_plain(psock_t *sock, void *buff, int num) { int r; while (1) { psock_try_write_buffer(sock); if (sock->pending) sock->pending = 0; else if (psock_wait_read_timeout(sock->sock)) { pdbg_logf(D_WARNING, "timeouted on socket with pending buffers"); return -1; } else psock_try_write_buffer(sock); r = read(sock->sock, buff, num); if (r == SOCKET_ERROR) { if (pdbg_likely(errno == EWOULDBLOCK || errno == EAGAIN || errno == EINTR)) continue; else return -1; } else return r; } } int psock_read(psock_t *sock, void *buff, int num) { if (sock->ssl) return psync_socket_read_ssl(sock, buff, num); else return psync_socket_read_plain(sock, buff, num); } static int psync_socket_read_noblock_ssl(psock_t *sock, void *buff, int num) { int r; r = pssl_read(sock->ssl, buff, num); if (r == PSYNC_SSL_FAIL) { sock->pending = 0; if (pdbg_likely(psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ || psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE)) return PSYNC_SOCKET_WOULDBLOCK; else { errno = (ECONNRESET); return -1; } } else return r; } static int psync_socket_read_noblock_plain(psock_t *sock, void *buff, int num) { int r; r = read(sock->sock, buff, num); if (r == SOCKET_ERROR) { sock->pending = 0; if (pdbg_likely(errno == EWOULDBLOCK || errno == EAGAIN || errno == EINTR)) return PSYNC_SOCKET_WOULDBLOCK; else return -1; } else return r; } int psock_read_noblock(psock_t *sock, void *buff, int num) { psock_try_write_buffer(sock); if (sock->ssl) return psync_socket_read_noblock_ssl(sock, buff, num); else return psync_socket_read_noblock_plain(sock, buff, num); } static int psync_socket_read_ssl_thread(psock_t *sock, void *buff, int num) { int r; pthread_mutex_lock(&mutex); psock_try_write_buffer(sock); pthread_mutex_unlock(&mutex); if (!pssl_pendingdata(sock->ssl) && !sock->pending && psock_wait_read_timeout(sock->sock)) return -1; sock->pending = 0; while (1) { pthread_mutex_lock(&mutex); psock_try_write_buffer(sock); r = pssl_read(sock->ssl, buff, num); pthread_mutex_unlock(&mutex); if (r == PSYNC_SSL_FAIL) { if (pdbg_likely(psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ || psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE)) { if (wait_ssl_ready(sock->sock)) return -1; else continue; } else { errno = (ECONNRESET); return -1; } } else return r; } } static int psync_socket_read_plain_thread(psock_t *sock, void *buff, int num) { int r; pthread_mutex_lock(&mutex); psock_try_write_buffer(sock); pthread_mutex_unlock(&mutex); while (1) { if (sock->pending) sock->pending = 0; else if (psock_wait_read_timeout(sock->sock)) return -1; pthread_mutex_lock(&mutex); psock_try_write_buffer(sock); r = read(sock->sock, buff, num); pthread_mutex_unlock(&mutex); if (r == SOCKET_ERROR) { if (pdbg_likely(errno == EWOULDBLOCK || errno == EAGAIN)) continue; else return -1; } else return r; } } int psock_read_thread(psock_t *sock, void *buff, int num) { if (sock->ssl) return psync_socket_read_ssl_thread(sock, buff, num); else return psync_socket_read_plain_thread(sock, buff, num); } static int psync_socket_write_to_buf(psock_t *sock, const void *buff, int num) { psock_buf_t *b; pdbg_assert(sock->buffer); b = sock->buffer; while (b->next) b = b->next; if (likely(b->size - b->woffset >= num)) { memcpy(b->buff + b->woffset, buff, num); b->woffset += num; return num; } else { uint32_t rnum, wr; rnum = num; do { wr = b->size - b->woffset; if (!wr) { b->next = (psock_buf_t *)pmem_malloc(PMEM_SUBSYS_OTHER, offsetof(psock_buf_t, buff) + PSYNC_SECOND_SOCK_WRITE_BUFF_SIZE); b = b->next; b->next = NULL; b->size = PSYNC_SECOND_SOCK_WRITE_BUFF_SIZE; b->woffset = 0; b->roffset = 0; wr = PSYNC_SECOND_SOCK_WRITE_BUFF_SIZE; } if (wr > rnum) wr = rnum; memcpy(b->buff + b->woffset, buff, wr); b->woffset += wr; buff = (const char *)buff + wr; rnum -= wr; } while (rnum); return num; } } int psock_write(psock_t *sock, const void *buff, int num) { int r; if (sock->buffer) return psync_socket_write_to_buf(sock, buff, num); if (psock_wait_write_timeout(sock->sock)) return -1; if (sock->ssl) { r = pssl_write(sock->ssl, buff, num); if (r == PSYNC_SSL_FAIL) { if (pdbg_likely(psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ || psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE)) return 0; else return -1; } } else { r = write(sock->sock, buff, num); if (r == SOCKET_ERROR) { if (pdbg_likely(errno == EWOULDBLOCK || errno == EAGAIN || errno == EINTR)) return 0; else return -1; } } return r; } static int psync_socket_readall_ssl(psock_t *sock, void *buff, int num) { int br, r; br = 0; psock_try_write_buffer(sock); if (!pssl_pendingdata(sock->ssl) && !sock->pending && psock_wait_read_timeout(sock->sock)) { return -1; } sock->pending = 0; while (br < num) { psock_try_write_buffer(sock); r = pssl_read(sock->ssl, (char *)buff + br, num - br); if (r == PSYNC_SSL_FAIL) { if (pdbg_likely(psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ || psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE)) { if (wait_ssl_ready(sock->sock)) return -1; else continue; } else { errno = (ECONNRESET); return -1; } } if (r == 0) return br; br += r; } return br; } static int psync_socket_readall_plain(psock_t *sock, void *buff, int num) { int br, r; br = 0; while (br < num) { psock_try_write_buffer(sock); if (sock->pending) sock->pending = 0; else if (psock_wait_read_timeout(sock->sock)) return -1; else psock_try_write_buffer(sock); r = read(sock->sock, (char *)buff + br, num - br); if (r == SOCKET_ERROR) { if (pdbg_likely(errno == EWOULDBLOCK || errno == EAGAIN || errno == EINTR)) continue; else return -1; } if (r == 0) return br; br += r; } return br; } int psock_readall(psock_t *sock, void *buff, int num) { if (sock->ssl) { return psync_socket_readall_ssl(sock, buff, num); } else { return psync_socket_readall_plain(sock, buff, num); } } static int psync_socket_writeall_ssl(psock_t *sock, const void *buff, int num) { int br, r; br = 0; while (br < num) { r = pssl_write(sock->ssl, (char *)buff + br, num - br); if (r == PSYNC_SSL_FAIL) { if (psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ || psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE) { if (wait_ssl_ready(sock->sock)) return -1; else continue; } else { errno = (ECONNRESET); return -1; } } if (r == 0) return br; br += r; } return br; } static int psync_socket_writeall_plain(int sock, const void *buff, int num) { int br, r; br = 0; while (br < num) { r = write(sock, (const char *)buff + br, num - br); if (r == SOCKET_ERROR) { if (errno == EWOULDBLOCK || errno == EAGAIN) { if (psock_wait_write_timeout(sock)) return -1; else continue; } else if (errno == EINTR) continue; else return -1; } br += r; } return br; } int psock_writeall(psock_t *sock, const void *buff, int num) { if (sock->buffer) return psync_socket_write_to_buf(sock, buff, num); if (sock->ssl) return psync_socket_writeall_ssl(sock, buff, num); else return psync_socket_writeall_plain(sock->sock, buff, num); } static int psync_socket_readall_ssl_thread(psock_t *sock, void *buff, int num) { int br, r; br = 0; pthread_mutex_lock(&mutex); psock_try_write_buffer(sock); r = pssl_pendingdata(sock->ssl); pthread_mutex_unlock(&mutex); if (!r && !sock->pending && psock_wait_read_timeout(sock->sock)) return -1; sock->pending = 0; while (br < num) { pthread_mutex_lock(&mutex); psock_try_write_buffer(sock); r = pssl_read(sock->ssl, (char *)buff + br, num - br); pthread_mutex_unlock(&mutex); if (r == PSYNC_SSL_FAIL) { if (pdbg_likely(psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ || psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE)) { if (wait_ssl_ready(sock->sock)) return -1; else continue; } else { errno = (ECONNRESET); return -1; } } if (r == 0) return br; br += r; } return br; } static int psync_socket_readall_plain_thread(psock_t *sock, void *buff, int num) { int br, r; br = 0; pthread_mutex_lock(&mutex); psock_try_write_buffer(sock); pthread_mutex_unlock(&mutex); while (br < num) { if (sock->pending) sock->pending = 0; else if (psock_wait_read_timeout(sock->sock)) return -1; pthread_mutex_lock(&mutex); psock_try_write_buffer(sock); r = read(sock->sock, (char *)buff + br, num - br); pthread_mutex_unlock(&mutex); if (r == SOCKET_ERROR) { if (pdbg_likely(errno == EWOULDBLOCK || errno == EAGAIN || errno == EINTR)) continue; else return -1; } if (r == 0) return br; br += r; } return br; } int psock_readall_thread(psock_t *sock, void *buff, int num) { if (sock->ssl) return psync_socket_readall_ssl_thread(sock, buff, num); else return psync_socket_readall_plain_thread(sock, buff, num); } static int psync_socket_writeall_ssl_thread(psock_t *sock, const void *buff, int num) { int br, r; br = 0; while (br < num) { pthread_mutex_lock(&mutex); if (sock->buffer) r = psync_socket_write_to_buf(sock, buff, num); else r = pssl_write(sock->ssl, (char *)buff + br, num - br); pthread_mutex_unlock(&mutex); if (r == PSYNC_SSL_FAIL) { if (psync_ssl_errno == PSYNC_SSL_ERR_WANT_READ || psync_ssl_errno == PSYNC_SSL_ERR_WANT_WRITE) { if (wait_ssl_ready(sock->sock)) return -1; else continue; } else { errno = (ECONNRESET); return -1; } } if (r == 0) return br; br += r; } return br; } static int psync_socket_writeall_plain_thread(psock_t *sock, const void *buff, int num) { int br, r; br = 0; while (br < num) { pthread_mutex_lock(&mutex); if (sock->buffer) r = psync_socket_write_to_buf(sock, buff, num); else r = write(sock->sock, (const char *)buff + br, num - br); pthread_mutex_unlock(&mutex); if (r == SOCKET_ERROR) { if (errno == EWOULDBLOCK || errno == EAGAIN) { if (psock_wait_write_timeout(sock->sock)) return -1; else continue; } else if (errno == EINTR) continue; else return -1; } br += r; } return br; } int psock_writeall_thread(psock_t *sock, const void *buff, int num) { if (sock->ssl) return psync_socket_writeall_ssl_thread(sock, buff, num); else return psync_socket_writeall_plain_thread(sock, buff, num); } static void copy_address(struct sockaddr_storage *dst, const struct sockaddr *src) { dst->ss_family = src->sa_family; if (src->sa_family == AF_INET) memcpy(&((struct sockaddr_in *)dst)->sin_addr, &((const struct sockaddr_in *)src)->sin_addr, sizeof(((struct sockaddr_in *)dst)->sin_addr)); else memcpy(&((struct sockaddr_in6 *)dst)->sin6_addr, &((const struct sockaddr_in6 *)src)->sin6_addr, sizeof(((struct sockaddr_in6 *)dst)->sin6_addr)); } psock_ifaces_t *psock_list_adapters() { psock_ifaces_t *ret; size_t cnt; struct ifaddrs *addrs, *addr; sa_family_t family; size_t sz; if (pdbg_unlikely(getifaddrs(&addrs))) goto empty; cnt = 0; addr = addrs; while (addr) { if (addr->ifa_addr) { family = addr->ifa_addr->sa_family; if ((family == AF_INET || family == AF_INET6) && addr->ifa_broadaddr && addr->ifa_netmask) cnt++; } addr = addr->ifa_next; } ret = pmem_malloc(PMEM_SUBSYS_OTHER, offsetof(psock_ifaces_t, interfaces) + sizeof(psock_iface_t) * cnt); memset(ret, 0, offsetof(psock_ifaces_t, interfaces) + sizeof(psock_iface_t) * cnt); ret->interfacecnt = cnt; addr = addrs; cnt = 0; while (addr) { if (addr->ifa_addr) { family = addr->ifa_addr->sa_family; if ((family == AF_INET || family == AF_INET6) && addr->ifa_broadaddr && addr->ifa_netmask) { if (family == AF_INET) sz = sizeof(struct sockaddr_in); else sz = sizeof(struct sockaddr_in6); copy_address(&ret->interfaces[cnt].address, addr->ifa_addr); copy_address(&ret->interfaces[cnt].broadcast, addr->ifa_broadaddr); copy_address(&ret->interfaces[cnt].netmask, addr->ifa_netmask); ret->interfaces[cnt].addrsize = sz; cnt++; } } addr = addr->ifa_next; } freeifaddrs(addrs); return ret; empty: ret = pmem_malloc(PMEM_SUBSYS_OTHER, offsetof(psock_ifaces_t, interfaces)); ret->interfacecnt = 0; return ret; } int psock_is_broken(int sock) { fd_set rfds; struct timeval tv; memset(&tv, 0, sizeof(tv)); FD_ZERO(&rfds); FD_SET(sock, &rfds); return select(sock + 1, NULL, NULL, &rfds, &tv) == 1; } int psock_wait_read_timeout(int sock) { return wait_readable(sock, PSYNC_SOCK_READ_TIMEOUT, 0); } int psock_select_in(int *sockets, int cnt, int64_t timeoutmillisec) { fd_set rfds; struct timeval tv, *ptv; int max; int i; if (timeoutmillisec < 0) ptv = NULL; else { tv.tv_sec = timeoutmillisec / 1000; tv.tv_usec = (timeoutmillisec % 1000) * 1000; ptv = &tv; } FD_ZERO(&rfds); max = 0; for (i = 0; i < cnt; i++) { FD_SET(sockets[i], &rfds); if (sockets[i] >= max) max = sockets[i] + 1; } i = select(max, &rfds, NULL, NULL, ptv); if (i > 0) { for (i = 0; i < cnt; i++) if (FD_ISSET(sockets[i], &rfds)) return i; } else if (i == 0) errno = (ETIMEDOUT); return SOCKET_ERROR; }