Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 10 additions & 2 deletions src/runtime/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand Down
95 changes: 90 additions & 5 deletions src/runtime/server/gameserver.cpp
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
#include "runtime/server/gameserver.h"
#include "core/curl_global.h"

#include <atomic>
#include <chrono>
#include <cstdio>
#include <cstring>
#include <exception>
Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -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<unsigned long>(GetCurrentThreadId()));

// N87: the game has installed its own console ctrl handler by now, which sits
// in front of the one InstallConsoleCtrlHandler() registered during
Expand All @@ -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 =
Expand Down Expand Up @@ -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<unsigned long>(thisThread), static_cast<unsigned long>(registryThread));
}

auto* lobby = m_context->GetLobby();
auto& cb = m_context->GetCallbackRegistry();
EchoVR::Broadcaster* liveOwner = lobby != nullptr ? lobby->broadcaster : nullptr;
Expand Down Expand Up @@ -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();

Expand Down Expand Up @@ -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<unsigned long>(GetCurrentThreadId());
std::atomic<unsigned long> gameThreadId{0};
const auto outcome = self->m_gameThreadHandoff.RunOnServicingThread(
[self, &gameThreadId]() {
gameThreadId.store(static_cast<unsigned long>(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<long long>(kGameThreadHandoffTimeout.count()), shutdownThreadId);
self->ShutdownUnregisterOffGameThread();
break;
}

self->m_shutdownComplete.store(true);

Expand Down Expand Up @@ -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) {
Expand Down
22 changes: 21 additions & 1 deletion src/runtime/server/gameserver.h
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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);

Expand All @@ -58,6 +60,16 @@ class GameServerLib : public EchoVR::IServerLib {
// Destructor waits on this before tearing down members.
std::atomic<bool> 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<DWORD> 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;
Expand All @@ -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)
Expand Down
100 changes: 100 additions & 0 deletions src/runtime/server/main_thread_handoff.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
/* SYNTHESIS -- custom tool code, not from binary */

#include "runtime/server/main_thread_handoff.h"

#include <exception>
#include <utility>

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<std::mutex> 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<std::mutex> 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<std::mutex> lock(mutex_);
state_ = threw ? State::kThrew : State::kDone;
}
finished_.notify_all();
return true;
}

void MainThreadHandoff::Cancel() {
{
std::lock_guard<std::mutex> 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
Loading
Loading