diff --git a/src/runtime/CMakeLists.txt b/src/runtime/CMakeLists.txt index 4c11a5fa..d16e2b52 100644 --- a/src/runtime/CMakeLists.txt +++ b/src/runtime/CMakeLists.txt @@ -58,6 +58,7 @@ set(PATCHES_SOURCES "server/session_success_dispatch.cpp" "server/callback_unregistration.cpp" "server/session_unregister.cpp" + "server/main_thread_handoff.cpp" "server/serverdb_uri.cpp" "server/telemetry_streamer.cpp" "server/upnp.cpp" @@ -120,6 +121,7 @@ set(PATCHES_HEADERS "server/session_success_dispatch.h" "server/callback_unregistration.h" "server/session_unregister.h" + "server/main_thread_handoff.h" "server/serverdb_uri.h" "server/telemetry_streamer.h" "server/messages.h" @@ -176,7 +178,7 @@ set_source_files_properties(ext/plugin_load_plan.cpp PROPERTIES SKIP_PRECOMPILE_ # #60: plugin_manifest.cpp builds the login's nevr_plugins array with nlohmann-json; # pure and PCH-free for the same reason as plugin_load_plan.cpp. set_source_files_properties(ext/plugin_manifest.cpp PROPERTIES SKIP_PRECOMPILE_HEADERS ON) -set_source_files_properties(log/url_diagnostics.cpp server/websocket_client.cpp server/websocket_frame.cpp server/protobuf_transport.cpp server/session_success_dispatch.cpp server/callback_unregistration.cpp server/session_unregister.cpp server/serverdb_uri.cpp server/telemetry_streamer.cpp server/upnp.cpp server/gameserver.cpp server/messages.cpp PROPERTIES SKIP_PRECOMPILE_HEADERS ON) +set_source_files_properties(log/url_diagnostics.cpp server/websocket_client.cpp server/websocket_frame.cpp server/protobuf_transport.cpp server/session_success_dispatch.cpp server/callback_unregistration.cpp server/session_unregister.cpp server/main_thread_handoff.cpp server/serverdb_uri.cpp server/telemetry_streamer.cpp server/upnp.cpp server/gameserver.cpp server/messages.cpp PROPERTIES SKIP_PRECOMPILE_HEADERS ON) # Scenario-test control endpoint (docs/design/2026-10-01-social-scenario-harness.md). It can inject # messages into a live session, so it is OFF by default and only the mingw-scenario preset turns it @@ -335,6 +337,9 @@ if(BUILD_TESTING) target_link_libraries(test_url_diagnostics PRIVATE GTest::gtest GTest::gtest_main CURL::libcurl) gtest_discover_tests(test_url_diagnostics DISCOVERY_MODE PRE_TEST) + # test_main_thread_handoff.cpp (GH #44) rides in this executable, so the + # existing `just test-auth-unit` loop runs it with no recipe change. + # Issue #41: the ServerDB URI builder (percent-encoded query) — pure, curl only. add_executable(test_serverdb_uri tests/test_serverdb_uri.cpp @@ -361,9 +366,12 @@ if(BUILD_TESTING) add_executable(test_callback_unregistration tests/test_callback_unregistration.cpp + tests/test_main_thread_handoff.cpp server/callback_unregistration.cpp + server/main_thread_handoff.cpp server/server_context.cpp) - set_source_files_properties(tests/test_callback_unregistration.cpp server/callback_unregistration.cpp server/server_context.cpp + set_source_files_properties(tests/test_callback_unregistration.cpp tests/test_main_thread_handoff.cpp + server/callback_unregistration.cpp server/main_thread_handoff.cpp server/server_context.cpp PROPERTIES SKIP_PRECOMPILE_HEADERS ON) target_include_directories(test_callback_unregistration PRIVATE ${CMAKE_SOURCE_DIR}/src) target_link_libraries(test_callback_unregistration PRIVATE GTest::gtest GTest::gtest_main) diff --git a/src/runtime/server/gameserver.cpp b/src/runtime/server/gameserver.cpp index 7189d539..f1da982d 100644 --- a/src/runtime/server/gameserver.cpp +++ b/src/runtime/server/gameserver.cpp @@ -1,6 +1,8 @@ #include "runtime/server/gameserver.h" #include "core/curl_global.h" +#include +#include #include #include #include @@ -882,6 +884,10 @@ GameServerLib::~GameServerLib() { // this destructor never finishes — which is correct. When the destructor runs // without a prior BeginGracefulShutdown call the thread is not joinable and // the join is a no-op. + // GH #44: a shutdown thread still waiting for Update() to service its + // unregister would hold this join for the whole hand-off timeout; withdraw + // the request so it takes its off-game-thread path now. + m_gameThreadHandoff.Cancel(); if (m_shutdownThread.joinable()) { m_shutdownThread.join(); } @@ -911,7 +917,8 @@ VOID* GameServerLib::Initialize(EchoVR::Lobby* lobby, EchoVR::Broadcaster* broad RegisterBroadcasterCallbacks(); RegisterTcpCallbacks(); - Log(EchoVR::LogLevel::Info, "[NEVR.GAMESERVER] Initialized game server"); + Log(EchoVR::LogLevel::Info, "[NEVR.GAMESERVER] Initialized game server (game thread %lu)", + static_cast(GetCurrentThreadId())); // N87: the game has installed its own console ctrl handler by now, which sits // in front of the one InstallConsoleCtrlHandler() registered during @@ -929,6 +936,7 @@ VOID* GameServerLib::Initialize(EchoVR::Lobby* lobby, EchoVR::Broadcaster* broad } void GameServerLib::RegisterBroadcasterCallbacks() { + m_registryThreadId.store(GetCurrentThreadId()); auto& cb = m_context->GetCallbackRegistry(); cb.sessionStart = @@ -1090,6 +1098,17 @@ void GameServerLib::RegisterTcpCallbacks() { } void GameServerLib::UnregisterAllCallbacks() { + // GH #44: the registry and EchoVR::BroadcasterUnlisten are game-thread-only + // (server_context.h). Make any future off-thread caller visible in the log. + const DWORD registryThread = m_registryThreadId.load(); + const DWORD thisThread = GetCurrentThreadId(); + if (registryThread != 0 && thisThread != registryThread) { + Log(EchoVR::LogLevel::Warning, + "[NEVR.GAMESERVER] callback registry reached from thread %lu; callbacks were registered on game thread %lu " + "(registry is game-thread-only, GH #44)", + static_cast(thisThread), static_cast(registryThread)); + } + auto* lobby = m_context->GetLobby(); auto& cb = m_context->GetCallbackRegistry(); EchoVR::Broadcaster* liveOwner = lobby != nullptr ? lobby->broadcaster : nullptr; @@ -1130,6 +1149,12 @@ static bool s_wasConnectedToServerDb = false; static bool s_exitPending = false; VOID GameServerLib::Update() { + // GH #44: run the graceful-shutdown thread's EndSession + Unregister here, on + // the game thread that owns the callback registry. Once it has run the server + // is unregistered and about to exit; skip the rest of the frame rather than + // process ServerDB traffic for a server that no longer exists. + if (m_gameThreadHandoff.Service()) return; + // Dispatch incoming ServerDB messages on the main thread if (m_wsClient) m_wsClient->ProcessReceivedMessages(); @@ -1212,8 +1237,48 @@ void GameServerLib::BeginGracefulShutdown(bool registrationFailed) { } } - self->EndSession(); - self->Unregister(); + // GH #44: EndSession + Unregister reach UnregisterAllCallbacks, i.e. the + // game-thread-only callback registry and EchoVR::BroadcasterUnlisten, which + // takes no lock (echovr.exe 0x140f8df20). Hand the work to Update() on the + // game thread and wait. The game can stop calling Update() (level + // transitions, teardown), so the wait is bounded; past it, do the + // ServerDB-facing half here and leave the registry alone — this process + // ends in ForceFatalExit below, which takes the listeners with it. + // 30 s is a chosen bound, not a measured one: the round-end wait above + // already allowed the post-round level transition 10 s of grace. + constexpr std::chrono::milliseconds kGameThreadHandoffTimeout{30 * 1000}; + const unsigned long shutdownThreadId = static_cast(GetCurrentThreadId()); + std::atomic gameThreadId{0}; + const auto outcome = self->m_gameThreadHandoff.RunOnServicingThread( + [self, &gameThreadId]() { + gameThreadId.store(static_cast(GetCurrentThreadId())); + self->ShutdownUnregisterOnGameThread(); + }, + kGameThreadHandoffTimeout); + const char* outcomeName = GameServer::MainThreadHandoffOutcomeName(outcome); + switch (outcome) { + case GameServer::MainThreadHandoff::Outcome::kRan: + Log(EchoVR::LogLevel::Info, + "[NEVR.GAMESERVER] shutdown unregister handoff=%s game_thread=%lu shutdown_thread=%lu", outcomeName, + gameThreadId.load(), shutdownThreadId); + break; + case GameServer::MainThreadHandoff::Outcome::kTaskThrew: + Log(EchoVR::LogLevel::Error, + "[NEVR.GAMESERVER] shutdown unregister handoff=%s game_thread=%lu shutdown_thread=%lu — " + "unregister threw on the game thread; exiting anyway", + outcomeName, gameThreadId.load(), shutdownThreadId); + break; + case GameServer::MainThreadHandoff::Outcome::kTimedOut: + case GameServer::MainThreadHandoff::Outcome::kCancelled: + case GameServer::MainThreadHandoff::Outcome::kBusy: + Log(EchoVR::LogLevel::Warning, + "[NEVR.GAMESERVER] shutdown unregister handoff=%s timeout_ms=%lld shutdown_thread=%lu — game thread did " + "not run it; ending session and unregistering from ServerDB on the shutdown thread, broadcaster " + "callbacks left registered (process is exiting)", + outcomeName, static_cast(kGameThreadHandoffTimeout.count()), shutdownThreadId); + self->ShutdownUnregisterOffGameThread(); + break; + } self->m_shutdownComplete.store(true); @@ -1599,12 +1664,32 @@ VOID GameServerLib::RequestRegistration(INT64 serverId, CHAR*, EchoVR::SymbolId } VOID GameServerLib::Unregister() { + // IServerLib entry point: the game calls this on its own thread. + UnregisterFromServerDb(true); +} + +void GameServerLib::ShutdownUnregisterOnGameThread() { + EndSession(); + UnregisterFromServerDb(true); +} + +void GameServerLib::ShutdownUnregisterOffGameThread() { + // Everything here is safe off the game thread: ServerContext state is + // mutex-guarded and WebSocketClient is documented thread-safe. The callback + // registry is not, so it is skipped (GH #44). + EndSession(); + UnregisterFromServerDb(false); +} + +void GameServerLib::UnregisterFromServerDb(bool touchCallbackRegistry) { const auto sendEnvelope = [this](const gameservice::v1::Envelope& envelope) { return GameServer::SendProtobufEnvelope(*m_wsClient, envelope); }; + GameServer::ServerLifecycleAction unregisterCallbacks; + if (touchCallbackRegistry) unregisterCallbacks = [this]() { UnregisterAllCallbacks(); }; const auto endResult = GameServer::UnregisterRegisteredServer( - *m_context, sendEnvelope, [this]() { m_wsClient->DiscardPendingMessages(); }, - [this]() { UnregisterAllCallbacks(); }, [this]() { m_wsClient->Disconnect(); }); + *m_context, sendEnvelope, [this]() { m_wsClient->DiscardPendingMessages(); }, unregisterCallbacks, + [this]() { m_wsClient->Disconnect(); }); if (endResult.attempted && endResult.sendResult == GameServer::ProtobufSendResult::TransportRejected) { Log(EchoVR::LogLevel::Warning, "[NEVR.SERVER] CODE_ENDED transport rejected during unregister"); } else if (endResult.sendResult == GameServer::ProtobufSendResult::AcceptedQueued) { diff --git a/src/runtime/server/gameserver.h b/src/runtime/server/gameserver.h index 749f5003..a91d230e 100644 --- a/src/runtime/server/gameserver.h +++ b/src/runtime/server/gameserver.h @@ -7,6 +7,7 @@ #include "runtime/server/constants.h" #include "abi/echovr.h" #include "core/pch.h" +#include "runtime/server/main_thread_handoff.h" #include "runtime/server/server_context.h" #include "runtime/server/telemetry_streamer.h" #include "runtime/server/websocket_client.h" @@ -45,7 +46,8 @@ class GameServerLib : public EchoVR::IServerLib { TelemetryStreamer& GetTelemetry() { return *m_telemetry; } /// Initiate graceful shutdown: disable reconnection, wait for round end (if active), - /// send EndSession, call Unregister, then ExitProcess(0). + /// have the game thread run EndSession + Unregister (GH #44), then exit via + /// ForceFatalExit(0). /// @param registrationFailed true if we're shutting down because registration was rejected. void BeginGracefulShutdown(bool registrationFailed); @@ -58,6 +60,16 @@ class GameServerLib : public EchoVR::IServerLib { // Destructor waits on this before tearing down members. std::atomic m_shutdownComplete{false}; + // GH #44: the shutdown thread hands EndSession + Unregister to the game thread + // through this; Update() services it. Declared before m_shutdownThread so it + // outlives the thread that waits on it. + GameServer::MainThreadHandoff m_gameThreadHandoff; + + // Thread that registered the broadcaster callbacks (the game thread). The + // callback registry is only safe on that thread; UnregisterAllCallbacks logs + // a warning when it is reached from any other. + std::atomic m_registryThreadId{0}; + // Joinable handle for the shutdown thread — replaces detached thread. // Joined in ~GameServerLib to prevent use-after-free on member destruction. std::thread m_shutdownThread; @@ -66,6 +78,14 @@ class GameServerLib : public EchoVR::IServerLib { void RegisterBroadcasterCallbacks(); void RegisterTcpCallbacks(); void UnregisterAllCallbacks(); + + // Unregister() body. touchCallbackRegistry=false skips UnregisterAllCallbacks + // and is the only form allowed off the game thread. + void UnregisterFromServerDb(bool touchCallbackRegistry); + + // Shutdown-thread work, split by the thread it may run on (GH #44). + void ShutdownUnregisterOnGameThread(); // EndSession + full Unregister + void ShutdownUnregisterOffGameThread(); // EndSession + ServerDB unregister, registry untouched }; // Logging helper (uses game's logging system) diff --git a/src/runtime/server/main_thread_handoff.cpp b/src/runtime/server/main_thread_handoff.cpp new file mode 100644 index 00000000..dbe5eb0a --- /dev/null +++ b/src/runtime/server/main_thread_handoff.cpp @@ -0,0 +1,100 @@ +/* SYNTHESIS -- custom tool code, not from binary */ + +#include "runtime/server/main_thread_handoff.h" + +#include +#include + +namespace GameServer { + +bool MainThreadHandoff::IsFinishedLocked() const { + return state_ == State::kDone || state_ == State::kThrew || state_ == State::kCancelled; +} + +MainThreadHandoff::Outcome MainThreadHandoff::RunOnServicingThread(Task task, + std::chrono::milliseconds timeout) { + std::unique_lock lock(mutex_); + if (state_ != State::kIdle) return Outcome::kBusy; + + task_ = std::move(task); + state_ = State::kPending; + pending_.store(true, std::memory_order_release); + + const auto finished = [this]() { return IsFinishedLocked(); }; + + if (!finished_.wait_for(lock, timeout, finished)) { + if (state_ == State::kPending) { + // Nobody picked it up. Withdraw it so a late Service() cannot run it + // after the caller has moved on to its fallback. + task_ = nullptr; + pending_.store(false, std::memory_order_release); + state_ = State::kIdle; + return Outcome::kTimedOut; + } + // kRunning: the servicing thread owns the task now. Wait for it. + finished_.wait(lock, finished); + } + + Outcome outcome = Outcome::kCancelled; + if (state_ == State::kDone) outcome = Outcome::kRan; + if (state_ == State::kThrew) outcome = Outcome::kTaskThrew; + state_ = State::kIdle; + return outcome; +} + +bool MainThreadHandoff::Service() { + if (!pending_.load(std::memory_order_acquire)) return false; + + Task task; + { + std::lock_guard lock(mutex_); + if (state_ != State::kPending) return false; + task = std::move(task_); + task_ = nullptr; + pending_.store(false, std::memory_order_release); + state_ = State::kRunning; + } + + bool threw = false; + try { + if (task) task(); + } catch (const std::exception&) { + threw = true; + } + + { + std::lock_guard lock(mutex_); + state_ = threw ? State::kThrew : State::kDone; + } + finished_.notify_all(); + return true; +} + +void MainThreadHandoff::Cancel() { + { + std::lock_guard lock(mutex_); + if (state_ != State::kPending) return; + task_ = nullptr; + pending_.store(false, std::memory_order_release); + state_ = State::kCancelled; + } + finished_.notify_all(); +} + +const char* MainThreadHandoffOutcomeName(MainThreadHandoff::Outcome outcome) { + switch (outcome) { + case MainThreadHandoff::Outcome::kRan: + return "ran"; + case MainThreadHandoff::Outcome::kTaskThrew: + return "task_threw"; + case MainThreadHandoff::Outcome::kTimedOut: + return "timed_out"; + case MainThreadHandoff::Outcome::kCancelled: + return "cancelled"; + case MainThreadHandoff::Outcome::kBusy: + return "busy"; + } + return "unknown"; +} + +} // namespace GameServer diff --git a/src/runtime/server/main_thread_handoff.h b/src/runtime/server/main_thread_handoff.h new file mode 100644 index 00000000..7f62382a --- /dev/null +++ b/src/runtime/server/main_thread_handoff.h @@ -0,0 +1,81 @@ +/* SYNTHESIS -- custom tool code, not from binary */ +#pragma once + +#include +#include +#include +#include +#include + +namespace GameServer { + +/// Hands one task from a background thread to the thread that calls Service() +/// — in production GameServerLib::Update(), which the game calls on its main +/// thread (GH #44). +/// +/// Why it exists: the callback registry (ServerContext::GetCallbackRegistry) +/// and the game's broadcaster listener table are main-thread-only. +/// EchoVR::BroadcasterUnlisten (echovr.exe 0x140f8df20) takes no lock: it +/// unlinks the handle from the broadcaster's listener hash chain and only +/// defers the delete when a plain "dispatch in progress" bit is set — a +/// same-thread reentrancy guard for SBroadcasterData's dispatch loop, not a +/// cross-thread one. The graceful-shutdown thread must therefore not +/// unregister by itself; it asks the game thread to do it and waits. +/// +/// The game thread is not ours: it can stop calling Update() (level +/// transitions, teardown). The wait is therefore bounded, and the caller +/// decides what to do on kTimedOut / kCancelled. +/// +/// One request at a time. A request the servicing thread has already started +/// is never abandoned: the requester keeps waiting past its timeout until the +/// task finishes, so it never runs the fallback while the task is still +/// running on the other thread. +class MainThreadHandoff { + public: + using Task = std::function; + + enum class Outcome { + kRan, // the servicing thread ran the task to completion + kTaskThrew, // the servicing thread ran the task and it threw a std::exception + kTimedOut, // nobody called Service() in time; the task did not run and never will + kCancelled, // Cancel() withdrew the request before it started; the task did not run + kBusy, // another request was already in flight; this task did not run + }; + + MainThreadHandoff() = default; + MainThreadHandoff(const MainThreadHandoff&) = delete; + MainThreadHandoff& operator=(const MainThreadHandoff&) = delete; + + /// Requester side (any thread except the servicing one). Queues the task and + /// blocks until it has run, the timeout elapses before it started, or + /// Cancel() withdraws it. + Outcome RunOnServicingThread(Task task, std::chrono::milliseconds timeout); + + /// Servicing side. Runs the pending task, if any, on the calling thread. + /// Returns true when a task ran (successfully or by throwing). Costs one + /// atomic load when nothing is pending. Never lets an exception escape: this + /// is called from a game vtable entry point. + bool Service(); + + /// Withdraws a request that has not started yet; its requester returns + /// kCancelled. A request already running is left to finish. + void Cancel(); + + /// True while a request is waiting to be serviced. + bool HasPending() const { return pending_.load(std::memory_order_acquire); } + + private: + enum class State { kIdle, kPending, kRunning, kDone, kThrew, kCancelled }; + + bool IsFinishedLocked() const; // caller holds mutex_ + + std::mutex mutex_; + std::condition_variable finished_; + std::atomic pending_{false}; + State state_ = State::kIdle; + Task task_; +}; + +const char* MainThreadHandoffOutcomeName(MainThreadHandoff::Outcome outcome); + +} // namespace GameServer diff --git a/src/runtime/server/server_context.cpp b/src/runtime/server/server_context.cpp index 5d920a05..75e3bddf 100644 --- a/src/runtime/server/server_context.cpp +++ b/src/runtime/server/server_context.cpp @@ -228,7 +228,9 @@ void ServerContext::UpdateSessionState(const SessionState& state) { CallbackRegistry& ServerContext::GetCallbackRegistry() { // Not synchronized — only safe from the game's main thread. // All current call sites (RegisterBroadcasterCallbacks, UnregisterAllCallbacks, - // Initialize, Terminate) run on the main thread. + // Initialize, Terminate) run on the main thread. The graceful-shutdown thread + // reaches UnregisterAllCallbacks only through GameServerLib::Update() via + // MainThreadHandoff (GH #44). return m_callbacks; } diff --git a/src/runtime/server/server_context.h b/src/runtime/server/server_context.h index 671d3737..6bd010df 100644 --- a/src/runtime/server/server_context.h +++ b/src/runtime/server/server_context.h @@ -118,7 +118,9 @@ class ServerContext { // Callback registry — NOT internally synchronized. // Safe to call without locking when all access is from the game's main thread // (RegisterBroadcasterCallbacks, UnregisterAllCallbacks, Initialize, Terminate). - // Must not be called from ixwebsocket or other background threads. + // Must not be called from ixwebsocket or other background threads. A background + // thread that needs registry work hands it to the game thread instead + // (GameServerLib::m_gameThreadHandoff, serviced in Update(); GH #44). CallbackRegistry& GetCallbackRegistry(); const CallbackRegistry& GetCallbackRegistry() const; diff --git a/src/runtime/tests/test_main_thread_handoff.cpp b/src/runtime/tests/test_main_thread_handoff.cpp new file mode 100644 index 00000000..0d8dab87 --- /dev/null +++ b/src/runtime/tests/test_main_thread_handoff.cpp @@ -0,0 +1,187 @@ +// GH #44: the graceful-shutdown thread hands EndSession + Unregister to the game +// thread through MainThreadHandoff instead of touching the main-thread-only +// callback registry itself. These tests pin the hand-off contract. +// +// The "game thread" below is the test's own thread calling Service() once per +// simulated frame, the way GameServerLib::Update() does. + +#include "runtime/server/main_thread_handoff.h" + +#include + +#include +#include +#include +#include +#include + +using GameServer::MainThreadHandoff; +using namespace std::chrono_literals; + +namespace { + +// Calls Service() once per simulated 1 ms frame until it runs a task or the +// deadline passes. Returns whether a task ran. +bool RunFramesUntilServiced(MainThreadHandoff& handoff, std::chrono::milliseconds deadline) { + const auto end = std::chrono::steady_clock::now() + deadline; + while (std::chrono::steady_clock::now() < end) { + if (handoff.Service()) return true; + std::this_thread::sleep_for(1ms); + } + return false; +} + +// Simulated frames until a request is pending (or the deadline passes). +bool RunFramesUntilPending(const MainThreadHandoff& handoff, std::chrono::milliseconds deadline) { + const auto end = std::chrono::steady_clock::now() + deadline; + while (std::chrono::steady_clock::now() < end) { + if (handoff.HasPending()) return true; + std::this_thread::sleep_for(1ms); + } + return false; +} + +} // namespace + +// The defect itself: before #44 the shutdown thread ran Unregister on its own +// thread. The task must run on the thread that calls Service(), not the requester. +TEST(MainThreadHandoff, TaskRunsOnTheServicingThreadNotTheRequester) { + MainThreadHandoff handoff; + const std::thread::id gameThread = std::this_thread::get_id(); + std::thread::id ranOn; + std::thread::id requesterThread; + + std::future outcome = std::async(std::launch::async, [&]() { + requesterThread = std::this_thread::get_id(); + return handoff.RunOnServicingThread([&]() { ranOn = std::this_thread::get_id(); }, 10s); + }); + + ASSERT_TRUE(RunFramesUntilServiced(handoff, 10s)) << "the request never reached the game thread"; + EXPECT_EQ(outcome.get(), MainThreadHandoff::Outcome::kRan); + EXPECT_EQ(ranOn, gameThread) << "the task ran off the servicing (game) thread"; + EXPECT_NE(ranOn, requesterThread) << "the task ran on the requesting (shutdown) thread"; +} + +TEST(MainThreadHandoff, ServiceIsANoOpWhenNothingIsPending) { + MainThreadHandoff handoff; + EXPECT_FALSE(handoff.HasPending()); + EXPECT_FALSE(handoff.Service()); +} + +// The game thread can stop calling Update(). The requester must get its thread +// back, and the withdrawn task must never run later behind the fallback's back. +TEST(MainThreadHandoff, UnservicedRequestTimesOutAndIsNeverRunLater) { + MainThreadHandoff handoff; + std::atomic runs{0}; + + EXPECT_EQ(handoff.RunOnServicingThread([&]() { ++runs; }, 20ms), MainThreadHandoff::Outcome::kTimedOut); + EXPECT_FALSE(handoff.HasPending()); + EXPECT_FALSE(handoff.Service()) << "a timed-out request was still serviceable"; + EXPECT_EQ(runs.load(), 0); +} + +// Once the game thread has started the task the requester must not give up on +// it — giving up would run the fallback concurrently with the task. +TEST(MainThreadHandoff, RequesterWaitsPastItsTimeoutForATaskAlreadyRunning) { + MainThreadHandoff handoff; + std::promise release; + std::shared_future released = release.get_future().share(); + std::promise started; + + std::future outcome = std::async(std::launch::async, [&]() { + return handoff.RunOnServicingThread( + [&]() { + started.set_value(); + released.wait(); + }, + 20ms); + }); + + // Game thread: wait for the request, then run it on a helper "frame" so this + // thread can hold the task open past the requester's 20 ms timeout. + if (!RunFramesUntilPending(handoff, 10s)) { + // Release before failing: a task that somehow already started (e.g. run + // inline on the requester) would otherwise block the future's destructor + // and turn this failure into a hang. + release.set_value(); + FAIL() << "the request never became pending for the game thread"; + } + std::thread frame([&]() { handoff.Service(); }); + started.get_future().wait(); + std::this_thread::sleep_for(100ms); // well past the requester's timeout + EXPECT_EQ(outcome.wait_for(0ms), std::future_status::timeout) + << "the requester returned while its task was still running"; + release.set_value(); + frame.join(); + + EXPECT_EQ(outcome.get(), MainThreadHandoff::Outcome::kRan); +} + +TEST(MainThreadHandoff, CancelWithdrawsAPendingRequest) { + MainThreadHandoff handoff; + std::atomic runs{0}; + + std::future outcome = std::async(std::launch::async, [&]() { + return handoff.RunOnServicingThread([&]() { ++runs; }, 10s); + }); + + ASSERT_TRUE(RunFramesUntilPending(handoff, 10s)); + handoff.Cancel(); + EXPECT_EQ(outcome.get(), MainThreadHandoff::Outcome::kCancelled); + EXPECT_FALSE(handoff.Service()); + EXPECT_EQ(runs.load(), 0); +} + +// Service() is called from a game vtable entry point; an exception must not +// escape into the game, and the requester must learn the task failed. +TEST(MainThreadHandoff, ThrowingTaskIsContainedAndReported) { + MainThreadHandoff handoff; + + std::future outcome = std::async(std::launch::async, [&]() { + return handoff.RunOnServicingThread([]() { throw std::runtime_error("unregister failed"); }, 10s); + }); + + bool serviced = false; + EXPECT_NO_THROW(serviced = RunFramesUntilServiced(handoff, 10s)); + EXPECT_TRUE(serviced); + EXPECT_EQ(outcome.get(), MainThreadHandoff::Outcome::kTaskThrew); +} + +TEST(MainThreadHandoff, SecondConcurrentRequestIsRefusedAndTheFirstStillRuns) { + MainThreadHandoff handoff; + std::atomic firstRuns{0}; + std::atomic secondRuns{0}; + + std::future first = std::async(std::launch::async, [&]() { + return handoff.RunOnServicingThread([&]() { ++firstRuns; }, 10s); + }); + ASSERT_TRUE(RunFramesUntilPending(handoff, 10s)); + + EXPECT_EQ(handoff.RunOnServicingThread([&]() { ++secondRuns; }, 10s), MainThreadHandoff::Outcome::kBusy); + + ASSERT_TRUE(RunFramesUntilServiced(handoff, 10s)); + EXPECT_EQ(first.get(), MainThreadHandoff::Outcome::kRan); + EXPECT_EQ(firstRuns.load(), 1); + EXPECT_EQ(secondRuns.load(), 0); +} + +TEST(MainThreadHandoff, IsReusableAfterARequestCompletes) { + MainThreadHandoff handoff; + for (int round = 0; round < 3; ++round) { + std::atomic runs{0}; + std::future outcome = std::async(std::launch::async, [&]() { + return handoff.RunOnServicingThread([&]() { ++runs; }, 10s); + }); + ASSERT_TRUE(RunFramesUntilServiced(handoff, 10s)) << "round " << round; + EXPECT_EQ(outcome.get(), MainThreadHandoff::Outcome::kRan) << "round " << round; + EXPECT_EQ(runs.load(), 1) << "round " << round; + } +} + +TEST(MainThreadHandoff, OutcomeNamesAreStableLogTokens) { + EXPECT_STREQ(GameServer::MainThreadHandoffOutcomeName(MainThreadHandoff::Outcome::kRan), "ran"); + EXPECT_STREQ(GameServer::MainThreadHandoffOutcomeName(MainThreadHandoff::Outcome::kTaskThrew), "task_threw"); + EXPECT_STREQ(GameServer::MainThreadHandoffOutcomeName(MainThreadHandoff::Outcome::kTimedOut), "timed_out"); + EXPECT_STREQ(GameServer::MainThreadHandoffOutcomeName(MainThreadHandoff::Outcome::kCancelled), "cancelled"); + EXPECT_STREQ(GameServer::MainThreadHandoffOutcomeName(MainThreadHandoff::Outcome::kBusy), "busy"); +} diff --git a/tools/tests/test_runtime_lifecycle_invariants.py b/tools/tests/test_runtime_lifecycle_invariants.py index 176ab728..2c3919b0 100644 --- a/tools/tests/test_runtime_lifecycle_invariants.py +++ b/tools/tests/test_runtime_lifecycle_invariants.py @@ -101,6 +101,37 @@ def test_only_reviewed_boot_hooks_are_optional(self): "EchoVR::JsonValueAsString", }) + def test_shutdown_thread_never_touches_the_callback_registry(self): + # Issue #44: the graceful-shutdown thread called self->Unregister(), which reaches + # UnregisterAllCallbacks -> GetCallbackRegistry() and EchoVR::BroadcasterUnlisten. The + # registry is documented game-thread-only (server_context.h) and BroadcasterUnlisten takes + # no lock (echovr.exe 0x140f8df20). The shutdown thread now hands that work to Update() + # through MainThreadHandoff; its own fallback must skip the registry. + source = (ROOT / "src/runtime/server/gameserver.cpp").read_text() + shutdown = extract_braced_function(source, "void GameServerLib::BeginGracefulShutdown(") + + for forbidden in (r"\bUnregister\s*\(\s*\)", r"\bUnregisterAllCallbacks\s*\(", + r"\bGetCallbackRegistry\s*\(", r"\bUnregisterFromServerDb\s*\(\s*true"): + self.assertNotRegex(shutdown, forbidden, + "the shutdown thread reaches the game-thread-only callback registry") + self.assertRegex(shutdown, r"m_gameThreadHandoff\.RunOnServicingThread\(", + "subject vanished: the shutdown thread no longer hands off to the game thread") + self.assertRegex(shutdown, r"ShutdownUnregisterOnGameThread\s*\(") + + off_thread = extract_braced_function(source, "void GameServerLib::ShutdownUnregisterOffGameThread(") + self.assertRegex(off_thread, r"UnregisterFromServerDb\s*\(\s*false\s*\)") + for forbidden in (r"\bUnregister\s*\(\s*\)", r"\bUnregisterAllCallbacks\s*\(", + r"\bGetCallbackRegistry\s*\(", r"\bUnregisterFromServerDb\s*\(\s*true"): + self.assertNotRegex(off_thread, forbidden) + + # The skip must be real: the registry action is only built when asked for. + impl = extract_braced_function(source, "void GameServerLib::UnregisterFromServerDb(") + self.assertRegex(impl, r"if\s*\(\s*touchCallbackRegistry\s*\)\s*unregisterCallbacks\s*=") + + # And the game thread must actually service the hand-off, or every shutdown times out. + update = extract_braced_function(source, "VOID GameServerLib::Update(") + self.assertRegex(update, r"m_gameThreadHandoff\.Service\(\)") + def test_bridge_never_closes_a_remote_on_unrequire(self): # d0190c4/dd1e9e7 (2026-09-14) closed the remote websocket whenever STcpConnectionUnrequireEvent # arrived in server mode, as an experiment (refuted in 79e27d5). Nakama sends that event on