diff --git a/CHANGELOG.md b/CHANGELOG.md index 819a741b73..9bdadb5049 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,10 @@ Increment the: * [CODE HEALTH] Move remaining API test helpers into anonymous namespaces [#4301](https://github.com/open-telemetry/opentelemetry-cpp/pull/4301) +* [BUG] Fix transient would-block send/recv handling in the ext/http embedded + server + [#4283](https://github.com/open-telemetry/opentelemetry-cpp/pull/4283) + * [CODE HEALTH] Move metrics storage test fixtures into anonymous namespace [#4286](https://github.com/open-telemetry/opentelemetry-cpp/pull/4286) diff --git a/ext/include/opentelemetry/ext/http/server/http_server.h b/ext/include/opentelemetry/ext/http/server/http_server.h index fed8518e94..43e9a1b7f8 100644 --- a/ext/include/opentelemetry/ext/http/server/http_server.h +++ b/ext/include/opentelemetry/ext/http/server/http_server.h @@ -282,15 +282,19 @@ class HttpServer : private SocketTools::Reactor::SocketCallback char buffer[2048] = {0}; int received = socket.recv(buffer, sizeof(buffer)); + int recvError = (received < 0) ? socket.error() : 0; LOG_TRACE("HttpServer: [%s] received %d", conn.request.client.c_str(), received); - if (received <= 0) + if (received > 0) { + conn.receiveBuffer.append(buffer, buffer + received); + handleConnection(conn); + } + else if (received == 0 || recvError != SocketTools::Socket::ErrorWouldBlock) + { + // Orderly shutdown (received == 0) or a real error; a transient would-block is + // retried on the next readable event. handleConnectionClosed(conn); - return; } - conn.receiveBuffer.append(buffer, buffer + received); - - handleConnection(conn); } void onSocketWritable(SocketTools::Socket socket) override @@ -340,12 +344,20 @@ class HttpServer : private SocketTools::Reactor::SocketCallback } int sent = conn.socket.send(conn.sendBuffer.data(), static_cast(conn.sendBuffer.size())); + int sendError = (sent < 0) ? conn.socket.error() : 0; LOG_TRACE("HttpServer: [%s] sent %d", conn.request.client.c_str(), sent); - if (sent < 0 && conn.socket.error() != SocketTools::Socket::ErrorWouldBlock) + if (sent < 0) { + if (sendError == SocketTools::Socket::ErrorWouldBlock) + { + // Backpressure: keep the unsent data and wait for the socket to become writable. + m_reactor.addSocket(conn.socket, + SocketTools::Reactor::Writable | SocketTools::Reactor::Closed); + } + // Never erase the buffer on a failed send: a negative count would wipe it. return true; } - conn.sendBuffer.erase(0, sent); + conn.sendBuffer.erase(0, static_cast(sent)); if (!conn.sendBuffer.empty()) { @@ -494,6 +506,13 @@ class HttpServer : private SocketTools::Reactor::SocketCallback return; } + // addSocket() replaces the interest set rather than adding to it, so a would-block + // while sending the interim response leaves the socket registered for writability + // only. Restore read interest before going back to waiting for the request body: + // this is the one transition from a sending state to a receiving one, and the + // keep-alive path at the end of SendingBody already does the same thing. + m_reactor.addSocket(conn.socket, + SocketTools::Reactor::Readable | SocketTools::Reactor::Closed); conn.state = Connection::ReceivingBody; LOG_TRACE("HttpServer: [%s] receiving body", conn.request.client.c_str()); } diff --git a/ext/test/http/BUILD b/ext/test/http/BUILD index 2e6a07e688..3db53373b5 100644 --- a/ext/test/http/BUILD +++ b/ext/test/http/BUILD @@ -3,6 +3,19 @@ load("@rules_cc//cc:cc_test.bzl", "cc_test") +cc_test( + name = "http_server_test", + srcs = [ + "http_server_test.cc", + ], + # The would-block test uses a POSIX socketpair; the body is compiled out on Windows. + tags = ["test"], + deps = [ + "//ext:headers", + "@com_google_googletest//:gtest_main", + ], +) + cc_test( name = "curl_http_test", srcs = [ diff --git a/ext/test/http/CMakeLists.txt b/ext/test/http/CMakeLists.txt index b7707bd8c2..e6f4597ae0 100644 --- a/ext/test/http/CMakeLists.txt +++ b/ext/test/http/CMakeLists.txt @@ -22,3 +22,18 @@ gtest_add_tests( TARGET ${URL_PARSER_FILENAME} TEST_PREFIX ext.http.urlparser. TEST_LIST ${URL_PARSER_FILENAME}) + +# http_server_test drives the would-block path with a POSIX socketpair; the body +# is compiled out on Windows, so build the target only where it has tests +# (gtest_add_tests parses the source for TEST() names and would otherwise +# register tests the Windows binary does not have). +if(NOT WIN32) + set(HTTP_SERVER_FILENAME http_server_test) + add_executable(${HTTP_SERVER_FILENAME} ${HTTP_SERVER_FILENAME}.cc) + target_link_libraries(${HTTP_SERVER_FILENAME} opentelemetry_ext ${GMOCK_LIB} + ${GTEST_BOTH_LIBRARIES} ${CMAKE_THREAD_LIBS_INIT}) + gtest_add_tests( + TARGET ${HTTP_SERVER_FILENAME} + TEST_PREFIX ext.http.httpserver. + TEST_LIST ${HTTP_SERVER_FILENAME}) +endif() diff --git a/ext/test/http/http_server_test.cc b/ext/test/http/http_server_test.cc new file mode 100644 index 0000000000..e5455e5b7c --- /dev/null +++ b/ext/test/http/http_server_test.cc @@ -0,0 +1,195 @@ +// Copyright The OpenTelemetry Authors +// SPDX-License-Identifier: Apache-2.0 + +#include + +#include "opentelemetry/ext/http/server/http_server.h" +#include "opentelemetry/ext/http/server/socket_tools.h" + +// The would-block path is exercised with a POSIX socketpair whose kernel send buffer is filled. +// Windows has no socketpair(); the shared HttpServer would-block decision logic these tests drive +// is platform-independent, so the coverage is provided on the POSIX runners. Windows-specific +// reactor notification behavior is outside this target. +#ifndef _WIN32 + +# include +# include +# include +# include +# include +# include +# include +# include + +namespace +{ + +// sendMore() is protected because it is a reactor internal. A test subclass reaches it so the +// would-block branch can be driven directly: under the level-triggered reactor a single send() +// after a writable readiness event returns a positive count on the normal Linux path, so this +// branch is otherwise reached from the readable path on a backlogged keep-alive connection, which +// is hard to stage deterministically. +class SendMoreProbe : public HTTP_SERVER_NS::HttpServer +{ +public: + using HttpServer::Connection; + using HttpServer::handleConnection; + using HttpServer::m_connections; + using HttpServer::m_reactor; + using HttpServer::sendMore; + + // After a would-block, sendMore() must re-arm the socket for Writable | Closed so the reactor + // calls back once the buffer drains. m_reactor is protected and Reactor::m_sockets is public, so + // the test can read back the software interest flags addSocket() recorded. + int reactorFlags(const SocketTools::Socket &socket) + { + for (const auto &data : m_reactor.m_sockets) + { + if (data.socket == socket) + { + return data.flags; + } + } + return 0; + } +}; + +// A nonblocking send() to a socket whose kernel send buffer is already full returns -1 with +// EWOULDBLOCK. sendMore() must treat that as transient: keep the unsent response intact and report +// there is more to send. Erasing on a negative count would compute erase(0, SIZE_MAX) and wipe the +// entire response. +TEST(HttpServerSendMoreTest, WouldBlockPreservesTheSendBuffer) +{ + int fds[2]; + ASSERT_EQ(::socketpair(AF_UNIX, SOCK_STREAM, 0, fds), 0); + + // A small send buffer plus a nonblocking sender fills quickly; the peer never reads. + int sndbuf = 1024; + ASSERT_EQ(::setsockopt(fds[0], SOL_SOCKET, SO_SNDBUF, &sndbuf, sizeof(sndbuf)), 0); + const int fl = ::fcntl(fds[0], F_GETFL, 0); + ASSERT_NE(fl, -1); + ASSERT_EQ(::fcntl(fds[0], F_SETFL, fl | O_NONBLOCK), 0); + + // Fill the kernel send buffer so the next send() would block. Bounded so a misbehaving kernel + // cannot spin here forever. + const std::string filler(4096, 'x'); + bool blocked = false; + for (int i = 0; i < 100000; ++i) + { + errno = 0; + const ssize_t sent = ::send(fds[0], filler.data(), filler.size(), 0); + if (sent < 0 && errno == EINTR) + { + continue; + } + if (sent < 0 && errno == SocketTools::Socket::ErrorWouldBlock) + { + blocked = true; + break; + } + ASSERT_GE(sent, 0); + } + ASSERT_TRUE(blocked) << "could not fill the socket send buffer"; + + SendMoreProbe server; + SendMoreProbe::Connection conn; + conn.socket = SocketTools::Socket(fds[0]); + const std::string response("HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nhello"); + conn.sendBuffer = response; + + const bool more = server.sendMore(conn); + + // sendMore() hit the would-block branch: the response is kept (not wiped) and there is more to + // send. + EXPECT_TRUE(more); + EXPECT_EQ(conn.sendBuffer, response); + // The socket is also re-armed for Writable | Closed so the reactor resumes the send once the + // buffer drains; without this the response is kept but the connection would stall forever. + EXPECT_EQ(server.reactorFlags(conn.socket), + SocketTools::Reactor::Writable | SocketTools::Reactor::Closed); + + conn.socket.close(); // closes fds[0] + ::close(fds[1]); +} + +// A readiness notification can go stale before the callback runs, so a nonblocking recv() in +// onSocketReadable() can return -1 with EWOULDBLOCK even though the readable callback fired. That +// transient must be ignored (wait for the next readable event) rather than treated as end-of-stream +// that tears the connection down. +TEST(HttpServerReadableTest, WouldBlockKeepsTheConnection) +{ + int fds[2]; + ASSERT_EQ(::socketpair(AF_UNIX, SOCK_STREAM, 0, fds), 0); + + // The server side is nonblocking; the peer never writes, so recv() reports would-block. + const int fl = ::fcntl(fds[0], F_GETFL, 0); + ASSERT_NE(fl, -1); + ASSERT_EQ(::fcntl(fds[0], F_SETFL, fl | O_NONBLOCK), 0); + + SendMoreProbe server; + const SocketTools::Socket sock(fds[0]); + + // Stage a mid-request connection and arm it exactly as a freshly accepted connection would be. + SendMoreProbe::Connection &conn = server.m_connections[sock]; + conn.socket = sock; + conn.state = SendMoreProbe::Connection::ReceivingHeaders; + server.m_reactor.addSocket(sock, SocketTools::Reactor::Readable | SocketTools::Reactor::Closed); + + // Drive the readable path through the reactor's public callback reference. Virtual dispatch + // reaches the private onSocketReadable() override, so no production surface has to be widened. + server.m_reactor.m_callback.onSocketReadable(sock); + + // The would-block recv() is transient: the connection stays in the map, the socket stays open, + // and its reactor registration is untouched (still Readable | Closed) so a later readable event + // can retry. + EXPECT_TRUE(server.m_connections.find(sock) != server.m_connections.end()); + EXPECT_NE(::fcntl(fds[0], F_GETFD, 0), -1); + EXPECT_EQ(server.reactorFlags(sock), + SocketTools::Reactor::Readable | SocketTools::Reactor::Closed); + + ::close(fds[0]); + ::close(fds[1]); +} + +// After the interim "100 Continue" response has drained, the server must return to reading so it +// can receive the request body. addSocket() replaces (does not merge) the interest set, so a +// would-block during the interim send leaves the socket armed for Writable only; the +// Sending100Continue -> ReceivingBody transition must therefore restore Readable | Closed, +// otherwise the body is never read and the connection stalls. +TEST(HttpServer100ContinueTest, CompletedInterimResponseRestoresReadInterest) +{ + int fds[2]; + ASSERT_EQ(::socketpair(AF_UNIX, SOCK_STREAM, 0, fds), 0); + const int fl = ::fcntl(fds[0], F_GETFL, 0); + ASSERT_NE(fl, -1); + ASSERT_EQ(::fcntl(fds[0], F_SETFL, fl | O_NONBLOCK), 0); + + SendMoreProbe server; + SendMoreProbe::Connection conn; + conn.socket = SocketTools::Socket(fds[0]); + // The interim response has finished sending (empty sendBuffer) and a one-byte body is still + // outstanding, so handleConnection() parks in ReceivingBody instead of running the request. + conn.state = SendMoreProbe::Connection::Sending100Continue; + conn.sendBuffer.clear(); + conn.receiveBuffer.clear(); + conn.contentLength = 1; + + // Simulate the Writable-only interest left behind when the interim send blocked. + server.m_reactor.addSocket(conn.socket, + SocketTools::Reactor::Writable | SocketTools::Reactor::Closed); + + server.handleConnection(conn); + + // The interim send is complete, so the server advances to reading the body and restores read + // interest. + EXPECT_EQ(conn.state, SendMoreProbe::Connection::ReceivingBody); + EXPECT_EQ(server.reactorFlags(conn.socket), + SocketTools::Reactor::Readable | SocketTools::Reactor::Closed); + + ::close(fds[0]); + ::close(fds[1]); +} + +} // namespace + +#endif // _WIN32