From: Lucian Petrut Date: Wed, 22 Jan 2020 14:00:14 +0000 (+0000) Subject: common,msg: Fix socket handling X-Git-Tag: v16.1.0~652^2~12 X-Git-Url: http://git-server-git.apps.pok.os.sepia.ceph.com/?a=commitdiff_plain;h=7d6a2064464d90ff26a6d3839e4a5488089a3c59;p=ceph.git common,msg: Fix socket handling The Windows socket API is slightly different. send/recv/closesocket must be used instead of read/write/close. At the same time, WSALastError() must be used instead of errno, covered by ceph_sock_errno(). The Windows socket errors are similar to the errno ones, except for the "WSA" prefix and different values. Signed-off-by: Lucian Petrut --- diff --git a/src/common/CMakeLists.txt b/src/common/CMakeLists.txt index 39bc50d41fd..7a68bda23e1 100644 --- a/src/common/CMakeLists.txt +++ b/src/common/CMakeLists.txt @@ -135,6 +135,8 @@ elseif(SUN) list(APPEND common_srcs solaris_errno.cc) elseif(AIX) list(APPEND common_srcs aix_errno.cc) +elseif(WIN32) + list(APPEND common_srcs win32_errno.c) endif() if(WITH_EVENTTRACE) diff --git a/src/common/admin_socket.cc b/src/common/admin_socket.cc index 1b5254fef86..328a9f1ae8a 100644 --- a/src/common/admin_socket.cc +++ b/src/common/admin_socket.cc @@ -131,7 +131,7 @@ std::string AdminSocket::create_wakeup_pipe(int *pipe_rd, int *pipe_wr) #else if (pipe_cloexec(pipefd, O_NONBLOCK) < 0) { #endif - int e = errno; + int e = ceph_sock_errno(); ostringstream oss; oss << "AdminSocket::create_wakeup_pipe error: " << cpp_strerror(e); return oss.str(); @@ -146,14 +146,10 @@ std::string AdminSocket::destroy_wakeup_pipe() { // Send a byte to the wakeup pipe that the thread is listening to char buf[1] = { 0x0 }; - #ifdef _WIN32 - int ret = send(m_wakeup_wr_fd, buf, sizeof(buf), 0) != 1; - #else - int ret = safe_write(m_wakeup_wr_fd, buf, sizeof(buf)); - #endif + int ret = safe_send(m_wakeup_wr_fd, buf, sizeof(buf)); // Close write end - retry_sys_call(::close, m_wakeup_wr_fd); + retry_sys_call(::compat_closesocket, m_wakeup_wr_fd); m_wakeup_wr_fd = -1; if (ret != 0) { @@ -167,7 +163,7 @@ std::string AdminSocket::destroy_wakeup_pipe() // Close read end. Doing this before join() blocks the listenter and prevents // joining. - retry_sys_call(::close, m_wakeup_rd_fd); + retry_sys_call(::compat_closesocket, m_wakeup_rd_fd); m_wakeup_rd_fd = -1; return ""; @@ -188,7 +184,7 @@ std::string AdminSocket::bind_and_listen(const std::string &sock_path, int *fd) } int sock_fd = socket_cloexec(PF_UNIX, SOCK_STREAM, 0); if (sock_fd < 0) { - int err = errno; + int err = ceph_sock_errno(); ostringstream oss; oss << "AdminSocket::bind_and_listen: " << "failed to create socket: " << cpp_strerror(err); @@ -201,7 +197,7 @@ std::string AdminSocket::bind_and_listen(const std::string &sock_path, int *fd) "%s", sock_path.c_str()); if (::bind(sock_fd, (struct sockaddr*)&address, sizeof(struct sockaddr_un)) != 0) { - int err = errno; + int err = ceph_sock_errno(); if (err == EADDRINUSE) { AdminSocketClient client(sock_path); bool ok; @@ -216,7 +212,7 @@ std::string AdminSocket::bind_and_listen(const std::string &sock_path, int *fd) sizeof(struct sockaddr_un)) == 0) { err = 0; } else { - err = errno; + err = ceph_sock_errno(); } } } @@ -230,7 +226,7 @@ std::string AdminSocket::bind_and_listen(const std::string &sock_path, int *fd) } } if (listen(sock_fd, 5) != 0) { - int err = errno; + int err = ceph_sock_errno(); ostringstream oss; oss << "AdminSocket::bind_and_listen: " << "failed to listen to socket: " << cpp_strerror(err); @@ -257,7 +253,7 @@ void AdminSocket::entry() noexcept ldout(m_cct,20) << __func__ << " waiting" << dendl; int ret = poll(fds, 2, -1); if (ret < 0) { - int err = errno; + int err = ceph_sock_errno(); if (err == EINTR) { continue; } @@ -274,9 +270,9 @@ void AdminSocket::entry() noexcept if (fds[1].revents & POLLIN) { // read off one byte char buf; - auto s = ::read(m_wakeup_rd_fd, &buf, 1); + auto s = safe_recv(m_wakeup_rd_fd, &buf, 1); if (s == -1) { - int e = errno; + int e = ceph_sock_errno(); ldout(m_cct, 5) << "AdminSocket: (ignoring) read(2) error: '" << cpp_strerror(e) << dendl; } @@ -322,7 +318,7 @@ void AdminSocket::do_accept() int connection_fd = accept_cloexec(m_sock_fd, (struct sockaddr*) &address, &address_length); if (connection_fd < 0) { - int err = errno; + int err = ceph_sock_errno(); lderr(m_cct) << "AdminSocket: do_accept error: '" << cpp_strerror(err) << dendl; return; @@ -333,13 +329,13 @@ void AdminSocket::do_accept() unsigned pos = 0; string c; while (1) { - int ret = safe_read(connection_fd, &cmd[pos], 1); + int ret = safe_recv(connection_fd, &cmd[pos], 1); if (ret <= 0) { if (ret < 0) { lderr(m_cct) << "AdminSocket: error reading request code: " << cpp_strerror(ret) << dendl; } - retry_sys_call(::close, connection_fd); + retry_sys_call(::compat_closesocket, connection_fd); return; } if (cmd[0] == '\0') { @@ -373,7 +369,7 @@ void AdminSocket::do_accept() } if (++pos >= sizeof(cmd)) { lderr(m_cct) << "AdminSocket: error reading request too long" << dendl; - retry_sys_call(::close, connection_fd); + retry_sys_call(::compat_closesocket, connection_fd); return; } } @@ -396,17 +392,18 @@ void AdminSocket::do_accept() out.claim_append(o); } uint32_t len = htonl(out.length()); - int ret = safe_write(connection_fd, &len, sizeof(len)); + int ret = safe_send(connection_fd, &len, sizeof(len)); if (ret < 0) { lderr(m_cct) << "AdminSocket: error writing response length " << cpp_strerror(ret) << dendl; } else { - if (out.write_fd(connection_fd) < 0) { + int r = out.send_fd(connection_fd); + if (r < 0) { lderr(m_cct) << "AdminSocket: error writing response payload " << cpp_strerror(ret) << dendl; } } - retry_sys_call(::close, connection_fd); + retry_sys_call(::compat_closesocket, connection_fd); } void AdminSocket::do_tell_queue() @@ -753,7 +750,7 @@ void AdminSocket::shutdown() lderr(m_cct) << "AdminSocket::shutdown: error: " << err << dendl; } - retry_sys_call(::close, m_sock_fd); + retry_sys_call(::compat_closesocket, m_sock_fd); unregister_commands(version_hook.get()); version_hook.reset(); @@ -772,10 +769,6 @@ void AdminSocket::wakeup() { // Send a byte to the wakeup pipe that the thread is listening to char buf[1] = { 0x0 }; - #ifdef _WIN32 - int r = send(m_wakeup_wr_fd, buf, sizeof(buf), 0); - #else - int r = safe_write(m_wakeup_wr_fd, buf, sizeof(buf)); - #endif + int r = safe_send(m_wakeup_wr_fd, buf, sizeof(buf)); (void)r; } diff --git a/src/common/admin_socket_client.cc b/src/common/admin_socket_client.cc index f2dd6d0e927..aeddc100fc0 100644 --- a/src/common/admin_socket_client.cc +++ b/src/common/admin_socket_client.cc @@ -35,6 +35,11 @@ const char* get_rand_socket_path() if (g_socket_path == NULL) { char buf[512]; const char *tdir = getenv("TMPDIR"); + #ifdef _WIN32 + if (tdir == NULL) { + tdir = getenv("TEMP"); + } + #endif /* _WIN32 */ if (tdir == NULL) { tdir = "/tmp"; } @@ -50,7 +55,7 @@ static std::string asok_connect(const std::string &path, int *fd) { int socket_fd = socket_cloexec(PF_UNIX, SOCK_STREAM, 0); if(socket_fd < 0) { - int err = errno; + int err = ceph_sock_errno(); ostringstream oss; oss << "socket(PF_UNIX, SOCK_STREAM, 0) failed: " << cpp_strerror(err); return oss.str(); @@ -64,10 +69,10 @@ static std::string asok_connect(const std::string &path, int *fd) if (::connect(socket_fd, (struct sockaddr *) &address, sizeof(struct sockaddr_un)) != 0) { - int err = errno; + int err = ceph_sock_errno(); ostringstream oss; oss << "connect(" << socket_fd << ") failed: " << cpp_strerror(err); - close(socket_fd); + compat_closesocket(socket_fd); return oss.str(); } @@ -75,21 +80,21 @@ static std::string asok_connect(const std::string &path, int *fd) timer.tv_sec = 10; timer.tv_usec = 0; if (::setsockopt(socket_fd, SOL_SOCKET, SO_RCVTIMEO, (SOCKOPT_VAL_TYPE)&timer, sizeof(timer))) { - int err = errno; + int err = ceph_sock_errno(); ostringstream oss; oss << "setsockopt(" << socket_fd << ", SO_RCVTIMEO) failed: " << cpp_strerror(err); - close(socket_fd); + compat_closesocket(socket_fd); return oss.str(); } timer.tv_sec = 10; timer.tv_usec = 0; if (::setsockopt(socket_fd, SOL_SOCKET, SO_SNDTIMEO, (SOCKOPT_VAL_TYPE)&timer, sizeof(timer))) { - int err = errno; + int err = ceph_sock_errno(); ostringstream oss; oss << "setsockopt(" << socket_fd << ", SO_SNDTIMEO) failed: " << cpp_strerror(err); - close(socket_fd); + compat_closesocket(socket_fd); return oss.str(); } @@ -99,7 +104,7 @@ static std::string asok_connect(const std::string &path, int *fd) static std::string asok_request(int socket_fd, std::string request) { - ssize_t res = safe_write(socket_fd, request.c_str(), request.length() + 1); + ssize_t res = safe_send(socket_fd, request.c_str(), request.length() + 1); if (res < 0) { int err = res; ostringstream oss; @@ -138,30 +143,30 @@ std::string AdminSocketClient::do_request(std::string request, std::string *resu if (!err.empty()) { goto done; } - res = safe_read_exact(socket_fd, &message_size_raw, + res = safe_recv_exact(socket_fd, &message_size_raw, sizeof(message_size_raw)); if (res < 0) { int e = res; ostringstream oss; - oss << "safe_read(" << socket_fd << ") failed to read message size: " + oss << "safe_recv(" << socket_fd << ") failed to read message size: " << cpp_strerror(e); err = oss.str(); goto done; } message_size = ntohl(message_size_raw); buffer.resize(message_size, 0); - res = safe_read_exact(socket_fd, &buffer[0], message_size); + res = safe_recv_exact(socket_fd, &buffer[0], message_size); if (res < 0) { int e = res; ostringstream oss; - oss << "safe_read(" << socket_fd << ") failed: " << cpp_strerror(e); + oss << "safe_recv(" << socket_fd << ") failed: " << cpp_strerror(e); err = oss.str(); goto done; } //printf("MESSAGE FROM SERVER: %s\n", buffer.c_str()); std::swap(*result, buffer); done: - close(socket_fd); + compat_closesocket(socket_fd); out: return err; } diff --git a/src/common/buffer.cc b/src/common/buffer.cc index abd7415f017..6dbf9c00671 100644 --- a/src/common/buffer.cc +++ b/src/common/buffer.cc @@ -1799,6 +1799,17 @@ ssize_t buffer::list::read_fd(int fd, size_t len) return ret; } +ssize_t buffer::list::recv_fd(int fd, size_t len) +{ + auto bp = ptr_node::create(buffer::create(len)); + ssize_t ret = safe_recv(fd, (void*)bp->c_str(), len); + if (ret >= 0) { + bp->set_length(ret); + push_back(std::move(bp)); + } + return ret; +} + int buffer::list::write_file(const char *fn, int mode) { int fd = TEMP_FAILURE_RETRY(::open(fn, O_WRONLY|O_CREAT|O_TRUNC|O_CLOEXEC, mode)); @@ -1916,6 +1927,10 @@ int buffer::list::write_fd(int fd) const return 0; } +int buffer::list::send_fd(int fd) const { + return buffer::list::write_fd(fd); +} + int buffer::list::write_fd(int fd, uint64_t offset) const { iovec iov[IOV_MAX]; @@ -1959,11 +1974,37 @@ int buffer::list::write_fd(int fd) const written += r; } + + left_pbrs--; + p++; } return 0; +} + +int buffer::list::send_fd(int fd) const +{ + // There's no writev on Windows. WriteFileGather may be an option, + // but it has strict requirements in terms of buffer size and alignment. + auto p = std::cbegin(_buffers); + uint64_t left_pbrs = get_num_buffers(); + while (left_pbrs) { + int written = 0; + while (written < p->length()) { + int r = ::send(fd, p->c_str(), p->length() - written, 0); + if (r < 0) + return -ceph_sock_errno(); + written += r; + } + + left_pbrs--; + p++; + } + + return 0; } + int buffer::list::write_fd(int fd, uint64_t offset) const { int r = ::lseek64(fd, offset, SEEK_SET); diff --git a/src/common/safe_io.c b/src/common/safe_io.c index 22369af38fb..9b7dbb77924 100644 --- a/src/common/safe_io.c +++ b/src/common/safe_io.c @@ -21,6 +21,7 @@ #include #include #include +#include ssize_t safe_read(int fd, void *buf, size_t count) { @@ -43,6 +44,39 @@ ssize_t safe_read(int fd, void *buf, size_t count) return cnt; } +#ifdef _WIN32 +// "read" doesn't work with Windows sockets. +ssize_t safe_recv(int fd, void *buf, size_t count) +{ + size_t cnt = 0; + + while (cnt < count) { + ssize_t r = recv(fd, (SOCKOPT_VAL_TYPE)buf, count - cnt, 0); + if (r <= 0) { + if (r == 0) { + // EOF + return cnt; + } + int err = ceph_sock_errno(); + if (err == EAGAIN || err == EINTR) { + continue; + } + return -err; + } + cnt += r; + buf = (char *)buf + r; + } + return cnt; +} +#else +ssize_t safe_recv(int fd, void *buf, size_t count) +{ + // We'll use "safe_read" so that this can work with any type of + // file descriptor. + return safe_read(fd, buf, count); +} +#endif /* _WIN32 */ + ssize_t safe_read_exact(int fd, void *buf, size_t count) { ssize_t ret = safe_read(fd, buf, count); @@ -52,7 +86,17 @@ ssize_t safe_read_exact(int fd, void *buf, size_t count) return -EDOM; return 0; } - + +ssize_t safe_recv_exact(int fd, void *buf, size_t count) +{ + ssize_t ret = safe_recv(fd, buf, count); + if (ret < 0) + return ret; + if ((size_t)ret != count) + return -EDOM; + return 0; +} + ssize_t safe_write(int fd, const void *buf, size_t count) { while (count > 0) { @@ -68,6 +112,30 @@ ssize_t safe_write(int fd, const void *buf, size_t count) return 0; } +#ifdef _WIN32 +ssize_t safe_send(int fd, const void *buf, size_t count) +{ + while (count > 0) { + ssize_t r = send(fd, (SOCKOPT_VAL_TYPE)buf, count, 0); + if (r < 0) { + int err = ceph_sock_errno(); + if (err == EINTR || err == EAGAIN) { + continue; + } + return -err; + } + count -= r; + buf = (char *)buf + r; + } + return 0; +} +#else +ssize_t safe_send(int fd, const void *buf, size_t count) +{ + return safe_write(fd, buf, count); +} +#endif /* _WIN32 */ + ssize_t safe_pread(int fd, void *buf, size_t count, off_t offset) { size_t cnt = 0; diff --git a/src/common/safe_io.h b/src/common/safe_io.h index 7ccbf37b604..6eb25d52c91 100644 --- a/src/common/safe_io.h +++ b/src/common/safe_io.h @@ -26,11 +26,17 @@ extern "C" { * Safe functions wrapping the raw read() and write() libc functions. * These retry on EINTR, and on error return -errno instead of returning * -1 and setting errno). + * + * On Windows, only recv/send work with sockets. */ ssize_t safe_read(int fd, void *buf, size_t count) WARN_UNUSED_RESULT; ssize_t safe_write(int fd, const void *buf, size_t count) WARN_UNUSED_RESULT; + ssize_t safe_recv(int fd, void *buf, size_t count) + WARN_UNUSED_RESULT; + ssize_t safe_send(int fd, const void *buf, size_t count) + WARN_UNUSED_RESULT; ssize_t safe_pread(int fd, void *buf, size_t count, off_t offset) WARN_UNUSED_RESULT; ssize_t safe_pwrite(int fd, const void *buf, size_t count, off_t offset) @@ -54,6 +60,8 @@ extern "C" { */ ssize_t safe_read_exact(int fd, void *buf, size_t count) WARN_UNUSED_RESULT; + ssize_t safe_recv_exact(int fd, void *buf, size_t count) + WARN_UNUSED_RESULT; ssize_t safe_pread_exact(int fd, void *buf, size_t count, off_t offset) WARN_UNUSED_RESULT; diff --git a/src/include/buffer.h b/src/include/buffer.h index 7e3c3d34233..67f3c67efab 100644 --- a/src/include/buffer.h +++ b/src/include/buffer.h @@ -1165,9 +1165,11 @@ struct error_code; ssize_t pread_file(const char *fn, uint64_t off, uint64_t len, std::string *error); int read_file(const char *fn, std::string *error); ssize_t read_fd(int fd, size_t len); + ssize_t recv_fd(int fd, size_t len); int write_file(const char *fn, int mode=0644); int write_fd(int fd) const; int write_fd(int fd, uint64_t offset) const; + int send_fd(int fd) const; template void prepare_iov(VectorT *piov) const { #ifdef __CEPH__ diff --git a/src/include/compat.h b/src/include/compat.h index f1eb0701dae..748fafee983 100644 --- a/src/include/compat.h +++ b/src/include/compat.h @@ -15,6 +15,7 @@ #include "acconfig.h" #include #include +#include #if defined(__linux__) #define PROCPREFIX @@ -212,6 +213,8 @@ unsigned get_page_size(); #include +#include "include/win32/win32_errno.h" + // There are a few name collisions between Windows headers and Ceph. // Updating Ceph definitions would be the prefferable fix in order to avoid // confussion, unless it requires too many changes, in which case we're going @@ -262,19 +265,6 @@ struct iovec { #define SIGKILL 9 #endif -#ifndef ENODATA -// mingw doesn't define this, the Windows SDK does. -#define ENODATA 120 -#endif - -#ifndef EDQUOT -#define EDQUOT ENOSPC -#endif - -#define ESHUTDOWN ECONNABORTED -#define ESTALE 256 -#define EREMOTEIO 257 - #define IOV_MAX 1024 #ifdef __cplusplus @@ -313,6 +303,7 @@ extern _CRTIMP errno_t __cdecl _putenv_s(const char *_Name,const char *_Value); } #endif +#define compat_closesocket closesocket // Use "aligned_free" when freeing memory allocated using posix_memalign or // _aligned_malloc. Using "free" will crash. #define aligned_free(ptr) _aligned_free(ptr) @@ -328,6 +319,9 @@ extern _CRTIMP errno_t __cdecl _putenv_s(const char *_Name,const char *_Value); #define SOCKOPT_VAL_TYPE void* #define aligned_free(ptr) free(ptr) +static inline int compat_closesocket(int fildes) { + return close(fildes); +} #endif /* WIN32 */ diff --git a/src/include/types.h b/src/include/types.h index 9723564b5eb..9eae90fdc08 100644 --- a/src/include/types.h +++ b/src/include/types.h @@ -524,9 +524,12 @@ WRITE_EQ_OPERATORS_1(shard_id_t, id) WRITE_CMP_OPERATORS_1(shard_id_t, id) std::ostream &operator<<(std::ostream &lhs, const shard_id_t &rhs); -#if defined(__sun) || defined(_AIX) || defined(__APPLE__) || defined(__FreeBSD__) +#if defined(__sun) || defined(_AIX) || defined(__APPLE__) || \ + defined(__FreeBSD__) || defined(_WIN32) +extern "C" { __s32 ceph_to_hostos_errno(__s32 e); __s32 hostos_to_ceph_errno(__s32 e); +} #else #define ceph_to_hostos_errno(e) (e) #define hostos_to_ceph_errno(e) (e) diff --git a/src/msg/async/Event.cc b/src/msg/async/Event.cc index 25bf6dde213..1e109230282 100644 --- a/src/msg/async/Event.cc +++ b/src/msg/async/Event.cc @@ -189,9 +189,9 @@ EventCenter::~EventCenter() //assert(time_events.empty()); if (notify_receive_fd >= 0) - ::close(notify_receive_fd); + compat_closesocket(notify_receive_fd); if (notify_send_fd >= 0) - ::close(notify_send_fd); + compat_closesocket(notify_send_fd); delete driver; if (notify_handler) diff --git a/src/msg/async/PosixStack.cc b/src/msg/async/PosixStack.cc index 93a79da84b4..c971546a10f 100644 --- a/src/msg/async/PosixStack.cc +++ b/src/msg/async/PosixStack.cc @@ -204,7 +204,7 @@ class PosixConnectedSocketImpl final : public ConnectedSocketImpl { ::shutdown(_fd, SHUT_RDWR); } void close() override { - ::close(_fd); + compat_closesocket(_fd); } int fd() const override { return _fd; diff --git a/src/msg/async/net_handler.cc b/src/msg/async/net_handler.cc index 7d9b81a4b32..59e641511fc 100644 --- a/src/msg/async/net_handler.cc +++ b/src/msg/async/net_handler.cc @@ -53,7 +53,7 @@ int NetHandler::create_socket(int domain, bool reuse_addr) r = ceph_sock_errno(); lderr(cct) << __func__ << " setsockopt SO_REUSEADDR failed: " << strerror(r) << dendl; - close(s); + compat_closesocket(s); return -r; } } @@ -178,7 +178,7 @@ int NetHandler::generic_connect(const entity_addr_t& addr, const entity_addr_t & if (nonblock) { ret = set_nonblock(s); if (ret < 0) { - close(s); + compat_closesocket(s); return ret; } } @@ -193,7 +193,7 @@ int NetHandler::generic_connect(const entity_addr_t& addr, const entity_addr_t & if (ret < 0) { ret = ceph_sock_errno(); ldout(cct, 2) << __func__ << " client bind error " << ", " << cpp_strerror(ret) << dendl; - close(s); + compat_closesocket(s); return -ret; } } @@ -207,7 +207,7 @@ int NetHandler::generic_connect(const entity_addr_t& addr, const entity_addr_t & return s; ldout(cct, 10) << __func__ << " connect: " << cpp_strerror(ret) << dendl; - close(s); + compat_closesocket(s); return -ret; }