diff --git a/.github/workflows/fork-ci.yml b/.github/workflows/fork-ci.yml index f8b37998f6..7e71515d7e 100644 --- a/.github/workflows/fork-ci.yml +++ b/.github/workflows/fork-ci.yml @@ -13,6 +13,11 @@ # The jobs only build and run the tests; they upload nothing. The workflow # uses GitHub's own actions only, so it runs under a repository setting that # allows no other actions. +# +# Each job gets 90 minutes (timeout-minutes; the slowest, Windows, takes +# about 35), and each test case its own limit from tests/CMakeLists.txt. A +# step that hangs then fails in time instead of holding the job until GitHub +# cancels it after six hours. name: Fork CI on: @@ -32,6 +37,7 @@ jobs: name: Windows x64 Release, editor ${{ matrix.editor }} if: github.repository != 'sven-n/MuMain' || github.event_name == 'workflow_dispatch' runs-on: windows-latest + timeout-minutes: 90 strategy: fail-fast: false matrix: @@ -120,6 +126,7 @@ jobs: name: Linux x64 Release, editor ON if: github.repository != 'sven-n/MuMain' || github.event_name == 'workflow_dispatch' runs-on: ubuntu-latest + timeout-minutes: 90 steps: - name: Checkout repository @@ -175,6 +182,7 @@ jobs: name: macOS arm64 Release, editor ON if: github.repository != 'sven-n/MuMain' || github.event_name == 'workflow_dispatch' runs-on: macos-latest + timeout-minutes: 90 steps: - name: Checkout repository diff --git a/src/source/Core/Platform/LocalSocket.cpp b/src/source/Core/Platform/LocalSocket.cpp index a5474384ca..6bb0bc2829 100644 --- a/src/source/Core/Platform/LocalSocket.cpp +++ b/src/source/Core/Platform/LocalSocket.cpp @@ -1,5 +1,7 @@ #include "Core/Platform/LocalSocket.h" +#include "Core/Platform/NonBlockingSocket.h" + #include #ifdef _WIN32 @@ -28,12 +30,6 @@ constexpr int SendFlags = 0; constexpr int SendFlags = MSG_NOSIGNAL; #endif -bool SetNonBlocking(SOCKET handle) -{ - u_long nonBlocking = 1; - return ioctlsocket(handle, FIONBIO, &nonBlocking) != SOCKET_ERROR; -} - // Whether a non-blocking connect has not finished yet, rather than failed. bool ConnectPending(int error) { @@ -44,15 +40,6 @@ bool ConnectPending(int error) #endif } -bool WouldBlock(int error) -{ -#ifdef _WIN32 - return error == WSAEWOULDBLOCK; -#else - return error == EWOULDBLOCK || error == EAGAIN || error == EINTR; -#endif -} - void ApplyOwnerOnlyMode(SOCKET handle, const std::string& path) { #ifdef _WIN32 @@ -171,7 +158,7 @@ bool LocalSocketConnection::ReadAvailable() return true; } - if (WouldBlock(WSAGetLastError())) + if (NonBlockingSocket::WouldBlock(WSAGetLastError())) { return true; } @@ -262,7 +249,7 @@ bool LocalSocketConnection::Flush() continue; } - if (sent < 0 && WouldBlock(WSAGetLastError())) + if (sent < 0 && NonBlockingSocket::WouldBlock(WSAGetLastError())) { // Peer is not reading yet; the rest goes out on a later poll. return true; @@ -395,7 +382,7 @@ bool SomethingIsListening(const std::string& path) // very case this probe detects. If the mode cannot be changed the probe // is abandoned rather than run blocking: not detecting a second client is // better than refusing to start. - if (!SetNonBlocking(probe)) + if (!NonBlockingSocket::Enable(probe)) { closesocket(probe); return true; @@ -529,7 +516,7 @@ bool LocalSocketListener::Listen(const std::string& path, std::string& error) ApplyOwnerOnlyMode(handle, path); - if (!SetNonBlocking(handle)) + if (!NonBlockingSocket::Enable(handle)) { error = "setting the control socket non-blocking failed: " + DescribeLastSocketError(); closesocket(handle); @@ -572,7 +559,7 @@ std::unique_ptr LocalSocketListener::Accept() return nullptr; } - if (!SetNonBlocking(accepted)) + if (!NonBlockingSocket::Enable(accepted)) { closesocket(accepted); return nullptr; diff --git a/src/source/Core/Platform/NonBlockingSocket.h b/src/source/Core/Platform/NonBlockingSocket.h new file mode 100644 index 0000000000..f96cd524af --- /dev/null +++ b/src/source/Core/Platform/NonBlockingSocket.h @@ -0,0 +1,28 @@ +// Non-blocking socket calls, for Windows and POSIX alike: the local socket +// transport and its tests share them, so both treat "try again later" the +// same way. +#pragma once + +#include "Core/Platform/WinSock.h" + +namespace Core::Platform::NonBlockingSocket +{ +// Switches a socket to non-blocking mode. Returns false when the mode could +// not be changed. +[[nodiscard]] inline bool Enable(SOCKET handle) +{ + u_long nonBlocking = 1; + return ioctlsocket(handle, FIONBIO, &nonBlocking) != SOCKET_ERROR; +} + +// Whether a call on a non-blocking socket failed only because it would have +// had to wait, and can simply be tried again later. +[[nodiscard]] inline bool WouldBlock(int error) +{ +#ifdef _WIN32 + return error == WSAEWOULDBLOCK; +#else + return error == EWOULDBLOCK || error == EAGAIN || error == EINTR; +#endif +} +} // namespace Core::Platform::NonBlockingSocket diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index e14e56ecdd..9eea42541a 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -32,6 +32,11 @@ if(MSVC) target_compile_options(mu_test_main PUBLIC /utf-8) endif() +# Longest one TEST_CASE may run before ctest stops it and reports it as +# failed. The slowest case takes about a second; without a limit, a case that +# hangs holds its CI job until the runner gives up hours later. +set(MU_TEST_CASE_TIMEOUT_SECONDS 120) + # Helper for module CMakeLists. Usage: # mu_add_test(NAME SOURCES a.cpp b.cpp LINK_LIBS lib1 lib2) # Registers each TEST_CASE inside the binary as its own CTest entry, so @@ -47,7 +52,7 @@ function(mu_add_test) if(MSVC) target_compile_options(${MAT_NAME} PRIVATE /utf-8) endif() - doctest_discover_tests(${MAT_NAME}) + doctest_discover_tests(${MAT_NAME} PROPERTIES TIMEOUT ${MU_TEST_CASE_TIMEOUT_SECONDS}) endfunction() find_package(Python3 COMPONENTS Interpreter REQUIRED) diff --git a/tests/core/test_local_socket.cpp b/tests/core/test_local_socket.cpp index 61b12a0854..f704f82ef3 100644 --- a/tests/core/test_local_socket.cpp +++ b/tests/core/test_local_socket.cpp @@ -2,8 +2,8 @@ // Linux alike. // // The listener is exercised against a real socket file in a temporary -// directory; the "client" is a plain blocking socket created by the test, so -// nothing here needs a window, a renderer or the game's globals. +// directory; the "client" is a plain socket created by the test, so nothing +// here needs a window, a renderer or the game's globals. // // Run: ctest --test-dir --build-config Release -R "\[core\]\[local-socket\]" @@ -19,6 +19,7 @@ #endif #include "Core/Platform/LocalSocket.h" +#include "Core/Platform/NonBlockingSocket.h" #include "Core/Platform/WinSock.h" // SOCKET, closesocket, WSAStartup (no-ops on POSIX) @@ -44,6 +45,26 @@ namespace { +// Bytes the send loops below hand the socket in one call: twice macOS's +// default AF_UNIX buffer, so a partial send is the normal case there. +constexpr std::size_t ChunkBytes = 16 * 1024; + +// Send buffer asked for on a LargeWriteClient: less than one chunk on every +// platform (Linux doubles what it is given). +constexpr int ClientSendBufferBytes = 4 * 1024; + +// How long a send loop waits for a client that cannot get a byte out before +// it gives up, instead of spinning until ctest stops the case. +constexpr std::chrono::milliseconds StallTimeout{5000}; + +// Pause between two polls of a socket that had nothing to offer. +constexpr std::chrono::milliseconds PollInterval{2}; + +// How long a writer has to stay held off before the pacing case takes the +// connection as paused: a connection that is still reading lets it send +// again on the very next pass. +constexpr std::chrono::milliseconds PauseSettleTimeout{200}; + // Winsock has to be started before the first socket call; the shim makes both // calls no-ops on POSIX. void EnsureSocketLibrary() @@ -88,7 +109,10 @@ std::filesystem::path MakeSocketDirectory() return directory; } -// Blocking client side, standing in for a test script. +// Client side, standing in for a test script. It blocks, which suits only +// payloads the kernel buffers whole: the connection is read on this same +// thread, so a send() that has to wait for it never returns. Larger payloads +// need ConnectForLargeWrites. SOCKET ConnectTo(const std::string& path) { EnsureSocketLibrary(); @@ -114,12 +138,48 @@ SOCKET ConnectTo(const std::string& path) return handle; } -// send()/recv() take an int length on Winsock and a size_t on POSIX. -int SendAll(SOCKET handle, std::string_view payload) +// A client that may send more than the socket buffer holds. Only +// ConnectForLargeWrites makes one, and only SendWhileServing takes one, so +// a large payload cannot go out through a blocking client by mistake. +struct LargeWriteClient +{ + SOCKET handle = INVALID_SOCKET; +}; + +// Non-blocking, because a send() that had to wait for the connection would +// wait for a read that can only run on this thread. With a send buffer +// smaller than one chunk, so every platform takes the partial sends that +// macOS's 8 KiB default forces, not only macOS. +LargeWriteClient ConnectForLargeWrites(const std::string& path) +{ + const SOCKET handle = ConnectTo(path); + REQUIRE(handle != INVALID_SOCKET); + const int bufferBytes = ClientSendBufferBytes; + REQUIRE(::setsockopt(handle, SOL_SOCKET, SO_SNDBUF, reinterpret_cast(&bufferBytes), + static_cast(sizeof(bufferBytes))) == 0); + REQUIRE(Core::Platform::NonBlockingSocket::Enable(handle)); + return LargeWriteClient{handle}; +} + +// One send() call, which may take only part of the payload. send()/recv() +// take an int length on Winsock and a size_t on POSIX. +int SendOnce(SOCKET handle, std::string_view payload) { return ::send(handle, payload.data(), static_cast(payload.size()), 0); } +// Bytes a non-blocking client's send() took: 0 while the socket buffer is +// full, until the connection reads from it; -1 when the socket failed. +int SendWhatFits(const LargeWriteClient& client, std::string_view payload) +{ + const int sent = SendOnce(client.handle, payload); + if (sent < 0 && Core::Platform::NonBlockingSocket::WouldBlock(WSAGetLastError())) + { + return 0; + } + return sent; +} + // The listener is non-blocking, so a connection may not be queued yet when // Accept() is first called; poll briefly instead of sleeping a fixed time. std::unique_ptr AcceptWithin(Core::Platform::LocalSocketListener& listener, @@ -132,33 +192,98 @@ std::unique_ptr AcceptWithin(Core::Platfo { return connection; } - std::this_thread::sleep_for(std::chrono::milliseconds(2)); + std::this_thread::sleep_for(PollInterval); } return nullptr; } -// Pushes a payload in through chunks, buffering each one on the connection -// without draining lines: how a pipelining script and the frame loop that -// only serves a few requests per frame interleave. -bool BufferInto(SOCKET client, Core::Platform::LocalSocketConnection& connection, std::string_view payload) +// Takes every complete line already buffered on the connection. +std::size_t TakeBufferedLines(Core::Platform::LocalSocketConnection& connection) { - constexpr std::size_t ChunkBytes = 16 * 1024; - std::size_t offset = 0; - while (offset < payload.size()) + std::size_t taken = 0; + std::string line; + while (connection.TakeLine(line)) { - const std::size_t size = std::min(ChunkBytes, payload.size() - offset); - const int sent = SendAll(client, payload.substr(offset, size)); - if (sent <= 0) + ++taken; + } + return taken; +} + +// Whether SendWhileServing takes the lines it buffers as they arrive. +enum class Drain +{ + Lines, // a frame loop that keeps up + Nothing, // a frame loop that has fallen behind +}; + +// How far SendWhileServing got. +struct Delivery +{ + std::size_t bytesSent = 0; + std::size_t linesTaken = 0; +}; + +// Sends a payload in chunks with the connection reading between them: how a +// pipelining script and the frame loop interleave. Stops once the payload is +// out, when the socket fails or the connection closes, or when the client +// has not got a byte out for `stallTimeout`: a connection that stopped +// reading, which would otherwise keep this loop going forever. +Delivery SendWhileServing(const LargeWriteClient& client, Core::Platform::LocalSocketConnection& connection, + std::string_view payload, Drain drain, std::chrono::milliseconds stallTimeout = StallTimeout) +{ + Delivery delivery; + auto lastProgress = std::chrono::steady_clock::now(); + while (delivery.bytesSent < payload.size()) + { + const std::size_t size = std::min(ChunkBytes, payload.size() - delivery.bytesSent); + const int sent = SendWhatFits(client, payload.substr(delivery.bytesSent, size)); + if (sent < 0) { - return false; + break; } - offset += static_cast(sent); + delivery.bytesSent += static_cast(sent); if (!connection.ReadAvailable()) { - return false; + break; + } + if (drain == Drain::Lines) + { + delivery.linesTaken += TakeBufferedLines(connection); + } + + const auto now = std::chrono::steady_clock::now(); + if (sent > 0) + { + lastProgress = now; + continue; } + if (now - lastProgress >= stallTimeout) + { + break; + } + std::this_thread::sleep_for(PollInterval); + } + return delivery; +} + +// Takes lines until `expected` have arrived, the connection closes or the +// timeout passes: the end of a stream may still be on its way when the last +// send() returns. +std::size_t TakeLinesWithin(Core::Platform::LocalSocketConnection& connection, std::size_t expected, + std::chrono::milliseconds timeout) +{ + const auto deadline = std::chrono::steady_clock::now() + timeout; + std::size_t taken = 0; + while (true) + { + const bool open = connection.ReadAvailable(); + taken += TakeBufferedLines(connection); + if (taken >= expected || !open || std::chrono::steady_clock::now() >= deadline) + { + return taken; + } + std::this_thread::sleep_for(PollInterval); } - return true; } bool ReadLineWithin(Core::Platform::LocalSocketConnection& connection, std::string& line, @@ -178,7 +303,7 @@ bool ReadLineWithin(Core::Platform::LocalSocketConnection& connection, std::stri // line yet" for a connection that can never produce one. return false; } - std::this_thread::sleep_for(std::chrono::milliseconds(2)); + std::this_thread::sleep_for(PollInterval); } return false; } @@ -205,7 +330,7 @@ TEST_CASE("Local socket serves a line round-trip [core][local-socket]") const std::string request = R"({"cmd":"ping"})" "\n"; - REQUIRE(SendAll(client, request) == static_cast(request.size())); + REQUIRE(SendOnce(client, request) == static_cast(request.size())); std::string line; REQUIRE(ReadLineWithin(*connection, line, std::chrono::milliseconds(500))); @@ -248,7 +373,7 @@ TEST_CASE("Local socket splits and preserves partial lines [core][local-socket]" // Two complete lines (one with a CRLF terminator) plus an unterminated tail. const std::string payload = "first\r\nsecond\nthi"; - REQUIRE(SendAll(client, payload) == static_cast(payload.size())); + REQUIRE(SendOnce(client, payload) == static_cast(payload.size())); std::string line; REQUIRE(ReadLineWithin(*connection, line, std::chrono::milliseconds(500))); @@ -261,7 +386,7 @@ TEST_CASE("Local socket splits and preserves partial lines [core][local-socket]" CHECK_FALSE(connection->TakeLine(line)); const std::string rest = "rd\n"; - REQUIRE(SendAll(client, rest) == static_cast(rest.size())); + REQUIRE(SendOnce(client, rest) == static_cast(rest.size())); REQUIRE(ReadLineWithin(*connection, line, std::chrono::milliseconds(500))); CHECK(line == "third"); @@ -392,7 +517,7 @@ TEST_CASE("Local socket refuses a path another listener is serving [core][local- const SOCKET client = ConnectTo(path); REQUIRE(client != INVALID_SOCKET); const std::string probe = "{\"cmd\":\"ping\"}\n"; - REQUIRE(SendAll(client, probe) == static_cast(probe.size())); + REQUIRE(SendOnce(client, probe) == static_cast(probe.size())); std::string line; std::unique_ptr served; @@ -445,7 +570,7 @@ TEST_CASE("Local socket reports a closed peer [core][local-socket]") while (!connection->PeerClosed() && std::chrono::steady_clock::now() < deadline) { connection->ReadAvailable(); - std::this_thread::sleep_for(std::chrono::milliseconds(2)); + std::this_thread::sleep_for(PollInterval); } CHECK(connection->PeerClosed()); CHECK_FALSE(connection->HasLine()); @@ -467,8 +592,7 @@ TEST_CASE("Local socket bounds the unterminated tail, not a pipelined batch [cor std::string error; REQUIRE(listener.Listen(path, error)); - const SOCKET client = ConnectTo(path); - REQUIRE(client != INVALID_SOCKET); + const LargeWriteClient client = ConnectForLargeWrites(path); auto connection = AcceptWithin(listener, std::chrono::milliseconds(500)); REQUIRE(connection != nullptr); @@ -485,25 +609,13 @@ TEST_CASE("Local socket bounds the unterminated tail, not a pipelined batch [cor } REQUIRE(batch.size() > Core::Platform::LocalSocketConnection::MaxPendingInputBytes); - // Served as it arrives, the way the frame loop does: reading pauses - // once a frame's worth of complete lines is waiting, so the batch is - // taken in several passes rather than all at once. - constexpr std::size_t ChunkBytes = 16 * 1024; - std::size_t offset = 0; - std::size_t taken = 0; - std::string line; - while (offset < batch.size()) - { - const std::size_t size = std::min(ChunkBytes, batch.size() - offset); - const int sent = SendAll(client, batch.substr(offset, size)); - REQUIRE(sent > 0); - offset += static_cast(sent); - REQUIRE(connection->ReadAvailable()); - while (connection->TakeLine(line)) - { - ++taken; - } - } + // Served as it arrives, the way a frame loop that keeps up does: the + // batch is taken in many passes, and its complete lines never count + // against the cap on an unterminated one. + const Delivery delivery = SendWhileServing(client, *connection, batch, Drain::Lines); + REQUIRE(delivery.bytesSent == batch.size()); + const std::size_t taken = + delivery.linesTaken + TakeLinesWithin(*connection, lines - delivery.linesTaken, std::chrono::milliseconds(500)); CHECK(connection->IsOpen()); CHECK(taken == lines); @@ -511,10 +623,11 @@ TEST_CASE("Local socket bounds the unterminated tail, not a pipelined batch [cor // Comfortably past the cap: crossing it on the last byte of the payload // would depend on that byte having arrived before ReadAvailable() runs. const std::string blob(Core::Platform::LocalSocketConnection::MaxPendingInputBytes + (64 * 1024), 'x'); - CHECK_FALSE(BufferInto(client, *connection, blob)); + const Delivery cutOff = SendWhileServing(client, *connection, blob, Drain::Nothing); + CHECK(cutOff.bytesSent < blob.size()); CHECK_FALSE(connection->IsOpen()); - closesocket(client); + closesocket(client.handle); listener.Close(); std::filesystem::remove_all(directory); } @@ -528,8 +641,7 @@ TEST_CASE("Local socket paces a peer that outruns the drain rate [core][local-so std::string error; REQUIRE(listener.Listen(path, error)); - const SOCKET client = ConnectTo(path); - REQUIRE(client != INVALID_SOCKET); + const LargeWriteClient client = ConnectForLargeWrites(path); auto connection = AcceptWithin(listener, std::chrono::milliseconds(500)); REQUIRE(connection != nullptr); @@ -544,38 +656,34 @@ TEST_CASE("Local socket paces a peer that outruns the drain rate [core][local-so batch += command; } - // Written in chunks with the reader and a consumer interleaved, the way - // the frame loop does it: the writer is never blocked out and nothing - // is dropped, while the buffer stays at the pause mark. - constexpr std::size_t ChunkBytes = 16 * 1024; - std::size_t offset = 0; - std::size_t taken = 0; - std::string line; - while (offset < batch.size()) - { - const std::size_t size = std::min(ChunkBytes, batch.size() - offset); - const int sent = SendAll(client, batch.substr(offset, size)); - REQUIRE(sent > 0); - offset += static_cast(sent); - - REQUIRE(connection->ReadAvailable()); - REQUIRE(connection->IsOpen()); - while (connection->TakeLine(line)) - { - ++taken; - } - } + const std::size_t lines = batch.size() / command.size(); - // Everything sent arrived, in order, and the connection is still open: - // a fast writer is paced rather than dropped. - CHECK(taken == batch.size() / command.size()); + // Nothing is served at first, like a frame loop that has fallen behind. + // The connection reads up to the pause mark and stops there; the kernel + // buffer fills behind it and the writer is held off, not dropped. + const Delivery behind = SendWhileServing(client, *connection, batch, Drain::Nothing, PauseSettleTimeout); CHECK(connection->IsOpen()); - - // Draining makes room, and the connection is still there to read more. - REQUIRE(connection->ReadAvailable()); + REQUIRE(behind.bytesSent < batch.size()); + + // What is buffered reached the pause mark, and stopped within one read + // of it. + const std::size_t pausedLines = TakeBufferedLines(*connection); + const std::size_t pausedBytes = pausedLines * command.size(); + CHECK(pausedBytes >= Core::Platform::LocalSocketConnection::ReadPauseBytes); + CHECK(pausedBytes < Core::Platform::LocalSocketConnection::ReadPauseBytes + + Core::Platform::LocalSocketConnection::ReadChunkBytes); + + // Served again, it picks up where it paused: the rest of the batch goes + // out, every line arrives and the connection is still open. + const std::string_view rest = std::string_view(batch).substr(behind.bytesSent); + const Delivery caughtUp = SendWhileServing(client, *connection, rest, Drain::Lines); + REQUIRE(caughtUp.bytesSent == rest.size()); + const std::size_t served = pausedLines + caughtUp.linesTaken; + const std::size_t taken = served + TakeLinesWithin(*connection, lines - served, std::chrono::milliseconds(500)); + CHECK(taken == lines); CHECK(connection->IsOpen()); - closesocket(client); + closesocket(client.handle); listener.Close(); std::filesystem::remove_all(directory); }