From 16545d573e7c85e45b21dbfcbf2c9713077e78b6 Mon Sep 17 00:00:00 2001 From: Aaron Kimbrell Date: Wed, 30 Sep 2026 02:44:16 -0500 Subject: [PATCH] feat(master): watch world zone files and replace instances that loaded old ones Master hashes each reported file on a worker, polls size and mtime every world_watch_seconds (a change must hold still for one poll), and replaces stale instances with the instance migration (Mythran shift, or seamless with world_reload_seamless=1); empty ones are stopped. Also on WORLD_RELOAD. Co-Authored-By: Claude Opus 5.5 --- dDashboardServer/routes/SettingsCatalog.cpp | 2 + dMasterServer/CMakeLists.txt | 1 + dMasterServer/MasterServer.cpp | 36 +- dMasterServer/WorldFileWatch.h | 273 +++++++++++++++ dMasterServer/WorldReloader.cpp | 349 ++++++++++++++++++++ dMasterServer/WorldReloader.h | 47 +++ resources/masterconfig.ini | 5 + tests/dCommonTests/CMakeLists.txt | 1 + tests/dCommonTests/WorldFileWatchTests.cpp | 252 ++++++++++++++ 9 files changed, 965 insertions(+), 1 deletion(-) create mode 100644 dMasterServer/WorldFileWatch.h create mode 100644 dMasterServer/WorldReloader.cpp create mode 100644 dMasterServer/WorldReloader.h create mode 100644 tests/dCommonTests/WorldFileWatchTests.cpp diff --git a/dDashboardServer/routes/SettingsCatalog.cpp b/dDashboardServer/routes/SettingsCatalog.cpp index 9d12d3dd6..f7911eb8c 100644 --- a/dDashboardServer/routes/SettingsCatalog.cpp +++ b/dDashboardServer/routes/SettingsCatalog.cpp @@ -141,6 +141,8 @@ namespace { c.Add(Unit(Int(MASTER, "live_update_warn_seconds", "Warning before moving players", "The game's Mythran maintenance warning is shown this long before players are moved. The dashboard can pick another for one update.", "10", 0, 300), "seconds")); c.Add(Int(MASTER, "live_update_parallel_worlds", "Worlds at once", "World instances replaced at the same time.", "4", 1, 64)); c.Add(Int(MASTER, "cdclient_watch_seconds", "CDClient watch interval", "Seconds between checks of the client's cdclient.fdb for changes; a change reloads it on every server. 0 turns the check off (/reloadcdclient still works).", "5", 0, 3600)); + c.Add(Int(MASTER, "world_watch_seconds", "World file watch interval", "Seconds between checks of the zone files running worlds loaded (.luz, .lvl, .lutriggers, terrain, navmesh). A change replaces the instances that loaded the old version with new ones and moves their players over. 0 turns the check off (/reloadworld and the dashboard's Reload still work).", "5", 0, 3600)); + c.Add(Bool(MASTER, "world_reload_seamless", "Seamless world reloads", "Move players to a reloaded world with the experimental seamless mode (no loading screen) instead of the Mythran shift. Untested with the real client.", false)); c.Add(Unit(Int(MASTER, "live_update_player_wait", "Wait for busy players", "Players who are dead or building are moved once they are done, or after this long.", "30", 0, 600), "seconds")); c.Add(Unit(Int(MASTER, "live_update_property_build_wait", "Wait for property builders", "A property is saved for its new instance once nobody builds there, or after this long (they leave build mode then).", "60", 0, 600), "seconds")); c.Add(Unit(Int(MASTER, "live_update_char_select_wait", "Wait at character select", "Players picking a character get this long to go in by themselves; then they are moved to the new character select.", "60", 0, 3600), "seconds")); diff --git a/dMasterServer/CMakeLists.txt b/dMasterServer/CMakeLists.txt index adc6c28be..11cfc5f2e 100644 --- a/dMasterServer/CMakeLists.txt +++ b/dMasterServer/CMakeLists.txt @@ -1,5 +1,6 @@ set(DMASTERSERVER_SOURCES "CDClientReloader.cpp" + "WorldReloader.cpp" "InstanceManager.cpp" "MigrationCoordinator.cpp" "LiveUpdateCoordinator.cpp" diff --git a/dMasterServer/MasterServer.cpp b/dMasterServer/MasterServer.cpp index d926f16a1..32a5e0120 100644 --- a/dMasterServer/MasterServer.cpp +++ b/dMasterServer/MasterServer.cpp @@ -61,6 +61,8 @@ #include "master/UgcModelsMade.h" #include "master/CDClientReload.h" #include "CDClientReloader.h" +#include "WorldReloader.h" +#include "master/WorldFiles.h" #include "BuildInfo.h" #ifdef DARKFLAME_PLATFORM_UNIX @@ -629,6 +631,7 @@ int main(int argc, char** argv) { } LiveUpdateCoordinator::Update(); CDClientReloader::Update(); + WorldReloader::Update(); CheckPlayerActionTimeouts(); // Spare instances for busy zones (zone_limits), checked every few seconds @@ -702,6 +705,7 @@ int main(int argc, char** argv) { if (instance->GetShutdownComplete()) { MigrationCoordinator::OnInstanceGone(*instance); + WorldReloader::OnInstanceGone(*instance); Game::im->RemoveInstance(instance); } } @@ -794,6 +798,7 @@ namespace { case ServiceType::DASHBOARD: dashboardServerMasterPeerSysAddr = sysAddr; g_DashboardConnects++; + WorldReloader::Republish(); // Its traffic isn't a player's; packet captures leave it out PacketCapture::IgnorePeer(sysAddr); break; @@ -1143,6 +1148,25 @@ namespace { else CDClientReloader::Request("a GM, character " + std::to_string(request.requesterId)); } + // World hot reload (docs/WorldHotReload.md): a world's zone files, and a GM's /reloadworld or the dashboard + void OnWorldFiles(const WorldFilesReport& report, const SystemAddress& sysAddr) { + WorldReloader::HandleReport(sysAddr, report); + } + + void OnWorldReload(const WorldReloadRequest& request, const SystemAddress& sysAddr) { + const bool fromDashboard = sysAddr == dashboardServerMasterPeerSysAddr && sysAddr != UNASSIGNED_SYSTEM_ADDRESS; + if (!fromDashboard && !Game::im->GetInstanceBySysAddr(sysAddr)) { + LOG("Ignoring a world reload request from a server that is neither the dashboard nor a world"); + return; + } + if (shutdownSequenceStarted) { + LOG("Shutdown sequence has been started. Not reloading zone %u.", request.zoneId); + return; + } + const std::string by = request.requestedBy.empty() ? (fromDashboard ? "the dashboard" : "a GM") : request.requestedBy; + WorldReloader::HandleRequest(request, fromDashboard ? "the dashboard, " + by : by); + } + void OnConfigReload(const ConfigReload& reload, const SystemAddress& sysAddr) { if (sysAddr != dashboardServerMasterPeerSysAddr) { LOG("Ignoring config reload from a server that is not the dashboard"); @@ -1247,6 +1271,8 @@ namespace { handlers.On(Master::DATA_CHANGED, ForwardWorldToDashboard); handlers.On(Master::MESSAGE_CAPTURE_CONTROL, OnMessageCaptureControl); handlers.On(Master::CDCLIENT_RELOAD, OnCDClientReload); + handlers.On(Master::WORLD_FILES, OnWorldFiles); + handlers.On(Master::WORLD_RELOAD, OnWorldReload); handlers.On(Master::MESSAGE_CAPTURE_DATA, OnMessageCaptureData); handlers.On(Master::REQUEST_SERVER_LIST, OnRequestServerList); handlers.On(Master::SERVER_TRAFFIC, OnServerTraffic); @@ -1281,6 +1307,7 @@ void HandlePacket(Packet* packet) { } MigrationCoordinator::OnInstanceGone(*instance); + WorldReloader::OnInstanceGone(*instance); Game::im->RemoveInstance(instance); } @@ -1343,6 +1370,7 @@ void HandlePacket(Packet* packet) { int ShutdownSequence(int32_t signal) { if (!Game::logger) return -1; CDClientReloader::Shutdown(); + WorldReloader::Shutdown(); LOG("Recieved Signal %d", signal); if (shutdownSequenceStarted) { LOG("Duplicate Shutdown Sequence"); @@ -1532,7 +1560,13 @@ void InitializeLiveUpdates() { } }; LiveUpdateCoordinator::Initialize(std::move(hooks)); - MigrationCoordinator::SetObserver(LiveUpdateCoordinator::OnMigrationStatus); + MigrationCoordinator::SetObserver([](const MigrationStatus& status) { + LiveUpdateCoordinator::OnMigrationStatus(status); + WorldReloader::OnMigrationStatus(status); + }); + WorldReloader::SetPublisher([](const WorldFilesStatus& status) { + if (dashboardServerMasterPeerSysAddr != UNASSIGNED_SYSTEM_ADDRESS) MasterPackets::SendTo(dashboardServerMasterPeerSysAddr, status); + }); // The worlds; the UGC and dashboard servers read the current files when they start (docs/CDClientFdb.md) CDClientReloader::SetBroadcast([](const CDClientReload& reload) { diff --git a/dMasterServer/WorldFileWatch.h b/dMasterServer/WorldFileWatch.h new file mode 100644 index 000000000..83e0050e1 --- /dev/null +++ b/dMasterServer/WorldFileWatch.h @@ -0,0 +1,273 @@ +#ifndef __WORLDFILEWATCH__H__ +#define __WORLDFILEWATCH__H__ + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "FdbSnapshot.h" +#include "ZoneFileLog.h" +#include "master/InstanceMigration.h" + +/** + * World hot reload without master's state (docs/WorldHotReload.md): which zone files the running instances loaded, + * which of those changed on disk, and which instances to replace. Header only, so it is unit tested without a master. + * + * Master reads (hashes) each watched file itself: once when a world first reports it, and again whenever its size or + * mtime changed and then held still for one poll (FdbSnapshot::Watcher). An instance is stale when a file it loaded + * has another hash on disk now. Files a world read from the client's packs are listed but not watched. + */ +namespace WorldFileWatch { + using Stamp = FdbSnapshot::Stamp; + + struct InstanceKey { + uint32_t zone{}; + uint32_t instance{}; + + auto operator<=>(const InstanceKey&) const = default; + }; + + struct FileState { + ZoneFileLog::eKind kind{}; + bool packed{}; + uint64_t reportedSize{}; // what the first world to report it read + uint64_t reportedHash{}; + std::optional hash; // master's read of the file on disk + uint64_t size{}; + bool missing{}; // the last poll or read found no file + bool hashing{}; + FdbSnapshot::Watcher watcher; + }; + + struct Loaded { + uint32_t clone{}; + std::map files; // path -> hash the world read + }; + + // A file to hash now, with the stamp it had when the poll saw it (accepted once the hash is in) + struct HashJob { + std::string path; + Stamp stamp; + }; + + class Tracker { + public: + // A world reported the files it loaded (replaces what that instance reported before) + void Report(uint32_t zone, uint32_t instance, uint32_t clone, const std::vector& files) { + auto& loaded = m_Instances[{ zone, instance }]; + loaded = {}; + loaded.clone = clone; + for (const auto& file : files) { + loaded.files[file.path] = file.hash; + auto [it, added] = m_Files.try_emplace(file.path); + if (!added) continue; + it->second.kind = file.kind; + it->second.packed = file.packed; + it->second.reportedSize = file.size; + it->second.reportedHash = file.hash; + } + } + + // An instance stopped; files no instance loaded any more stop being watched + void Forget(uint32_t zone, uint32_t instance) { + m_Instances.erase({ zone, instance }); + std::erase_if(m_Files, [this](const auto& entry) { return !entry.second.hashing && !IsLoaded(entry.first); }); + } + + /** + * Looks at every watched file (stampOf: its size and mtime now). Returns the ones to hash: files master hasn't + * read yet, and files whose size or mtime changed and held still since the last poll. Those are marked as + * being hashed until Hashed is called for them. + */ + std::vector Poll(const std::function& stampOf) { + std::vector due; + for (auto& [path, file] : m_Files) { + if (file.packed || file.hashing) continue; + const auto stamp = stampOf(path); + file.missing = !stamp.exists; + const bool first = !file.hash && stamp.exists; + if (!file.watcher.Poll(stamp) && !first) continue; + file.hashing = true; + due.push_back({ path, stamp }); + } + return due; + } + + /** + * A hash finished (hash nullopt: the file couldn't be read). Returns true when the file on disk is another + * version than master read before. + */ + bool Hashed(const HashJob& job, std::optional hash, uint64_t size) { + const auto it = m_Files.find(job.path); + if (it == m_Files.end()) return false; + auto& file = it->second; + file.hashing = false; + // Taken either way, so an unreadable file isn't read again every poll; the next change to it is + file.watcher.Accept(job.stamp); + if (!hash) { + file.missing = true; + if (!IsLoaded(job.path)) m_Files.erase(it); + return false; + } + file.missing = false; + const bool changed = file.hash && *file.hash != *hash; + file.hash = hash; + file.size = size; + if (!IsLoaded(job.path)) m_Files.erase(it); + return changed; + } + + // A file this instance loaded is another version on disk now + bool IsStale(const InstanceKey& key) const { + const auto it = m_Instances.find(key); + if (it == m_Instances.end()) return false; + for (const auto& [path, loadedHash] : it->second.files) { + const auto file = m_Files.find(path); + if (file != m_Files.end() && file->second.hash && *file->second.hash != loadedHash) return true; + } + return false; + } + + std::vector Stale() const { + std::vector stale; + for (const auto& [key, loaded] : m_Instances) { + if (IsStale(key)) stale.push_back(key); + } + return stale; + } + + // The versions on disk of the files this instance loaded; an automatic reload is tried once per signature + uint64_t Signature(const InstanceKey& key) const { + uint64_t signature = 14695981039346656037ULL; + const auto it = m_Instances.find(key); + if (it == m_Instances.end()) return signature; + for (const auto& [path, loadedHash] : it->second.files) { + const auto file = m_Files.find(path); + const uint64_t onDisk = file != m_Files.end() && file->second.hash ? *file->second.hash : loadedHash; + signature = FdbSnapshot::Hash(reinterpret_cast(path.data()), path.size(), signature); + signature = FdbSnapshot::Hash(reinterpret_cast(&onDisk), sizeof(onDisk), signature); + } + return signature; + } + + // The files of a stale instance that changed (for the log) + std::vector ChangedFiles(const InstanceKey& key) const { + std::vector changed; + const auto it = m_Instances.find(key); + if (it == m_Instances.end()) return changed; + for (const auto& [path, loadedHash] : it->second.files) { + const auto file = m_Files.find(path); + if (file != m_Files.end() && file->second.hash && *file->second.hash != loadedHash) changed.push_back(path); + } + return changed; + } + + // Any instance of zone loaded another version of path than is on disk + bool IsChanged(uint32_t zone, const std::string& path) const { + const auto file = m_Files.find(path); + if (file == m_Files.end() || !file->second.hash) return false; + for (const auto& [key, loaded] : m_Instances) { + if (key.zone != zone) continue; + const auto it = loaded.files.find(path); + if (it != loaded.files.end() && it->second != *file->second.hash) return true; + } + return false; + } + + // Every file some instance of zone loaded + std::set FilesOf(uint32_t zone) const { + std::set files; + for (const auto& [key, loaded] : m_Instances) { + if (key.zone != zone) continue; + for (const auto& [path, hash] : loaded.files) files.insert(path); + } + return files; + } + + std::set Zones() const { + std::set zones; + for (const auto& [key, loaded] : m_Instances) zones.insert(key.zone); + return zones; + } + + bool Knows(const InstanceKey& key) const { return m_Instances.contains(key); } + const std::map& Files() const { return m_Files; } + + private: + bool IsLoaded(const std::string& path) const { + for (const auto& [key, loaded] : m_Instances) { + if (loaded.files.contains(path)) return true; + } + return false; + } + + std::map m_Files; + std::map m_Instances; + }; + + enum class eAction : uint8_t { + REPLACE, // start a new instance (same clone, same password), move its players there, stop it + START_THEN_STOP, // nobody there, but the zone always has an instance: start a new one, stop this one + STOP, // nobody there: stop it (a new instance starts when someone goes there) + SKIP, + }; + + inline const char* ActionName(eAction action) { + switch (action) { + case eAction::REPLACE: return "replace"; + case eAction::START_THEN_STOP: return "start a new one, then stop it"; + case eAction::STOP: return "stop"; + case eAction::SKIP: return "skip"; + } + return "skip"; + } + + struct Choice { + InstanceMigration::InstanceView view; + eAction action{ eAction::SKIP }; + std::string reason; + }; + + /** + * What to do with each instance wanted() picks: replace it when players are there, stop it when it is empty (a zone + * in keepZones gets one new public instance first). Character selection loads no zone; instances still starting + * load the files on disk now (and report them); instances shutting down or being emptied already are left alone. + */ + inline std::vector Choose(const std::vector& instances, + const std::function& wanted, const std::set& keepZones) { + std::vector choices; + for (const auto& view : instances) { + if (!wanted(view)) continue; + Choice choice; + choice.view = view; + if (view.zoneId == 0) choice.reason = "Character selection loads no zone"; + else if (view.shuttingDown) choice.reason = "Already shutting down"; + else if (view.draining) choice.reason = "Its players are already being moved"; + else if (!view.ready) choice.reason = "Still starting; it loads the files on disk now"; + else if (view.players > 0) choice.action = eAction::REPLACE; + else choice.action = eAction::STOP; + choices.push_back(std::move(choice)); + } + // A zone that always has an instance keeps one: a public instance replaced for its players, or else one new + // instance started for the first empty public one + std::set keptZones; + for (const auto& choice : choices) { + if (choice.action == eAction::REPLACE && !choice.view.isPrivate && choice.view.cloneId == 0) keptZones.insert(choice.view.zoneId); + } + for (auto& choice : choices) { + const auto& view = choice.view; + if (choice.action != eAction::STOP || !keepZones.contains(view.zoneId) || view.isPrivate || view.cloneId != 0 || keptZones.contains(view.zoneId)) continue; + choice.action = eAction::START_THEN_STOP; + keptZones.insert(view.zoneId); + } + return choices; + } +} + +#endif //!__WORLDFILEWATCH__H__ diff --git a/dMasterServer/WorldReloader.cpp b/dMasterServer/WorldReloader.cpp new file mode 100644 index 000000000..a76226812 --- /dev/null +++ b/dMasterServer/WorldReloader.cpp @@ -0,0 +1,349 @@ +#include "WorldReloader.h" + +#include +#include +#include +#include +#include + +#include "Game.h" +#include "GeneralUtils.h" +#include "InstanceManager.h" +#include "Logger.h" +#include "MigrationCoordinator.h" +#include "WorldFileWatch.h" +#include "dConfig.h" +#include "master/InstanceMigration.h" +#include "master/WorldFiles.h" + +using namespace WorldFileWatch; + +namespace { + using Clock = std::chrono::steady_clock; + // World reload migrations get IDs of their own, far from the worlds' (time-based) and the live update's ones + constexpr uint32_t MIGRATION_ID_BASE = 0x80000000; + constexpr auto PUBLISH_INTERVAL = std::chrono::seconds(1); + + struct HashResult { + HashJob job; + std::optional hash; + uint64_t size{}; + }; + + Tracker g_Tracker; + std::future> g_Job; + Clock::time_point g_NextPoll{}; + Clock::time_point g_NextPublish{}; + bool g_Dirty = true; + uint32_t g_NextMigrationId = 0; + std::map g_Tried; // automatic reloads: the signature last tried per instance + std::map g_Migrations; // our migrations -> the instance being replaced + std::map g_ZoneMessages; // the last reload of each zone, for the dashboard + std::function g_Publisher; + WorldFilesStatus g_Published; + + uint32_t WatchSeconds() { + return GeneralUtils::TryParse(Game::config->GetValue("world_watch_seconds")).value_or(5); + } + + bool Seamless() { + return Game::config->GetValue("world_reload_seamless") == "1"; + } + + std::set KeepZones() { + std::set zones; + for (auto part : GeneralUtils::SplitString(Game::config->GetValue("prestart_worlds", "0,1000"), ',')) { + std::erase_if(part, [](const char c) { return std::isspace(static_cast(c)); }); + if (const auto zone = GeneralUtils::TryParse(part)) zones.insert(*zone); + } + return zones; + } + + std::string FileName(const std::string& path) { + const auto slash = path.find_last_of("/\\"); + return slash == std::string::npos ? path : path.substr(slash + 1); + } + + // The worker gets copies of the paths; it only reads the files and never logs or touches shared state + std::vector HashFiles(std::vector jobs) { + std::vector results; + results.reserve(jobs.size()); + for (auto& job : jobs) { + HashResult result; + std::error_code code; + const auto size = std::filesystem::file_size(job.path, code); + if (!code) { + result.size = static_cast(size); + result.hash = FdbSnapshot::HashFile(job.path); + } + result.job = std::move(job); + results.push_back(std::move(result)); + } + return results; + } + + std::vector Views() { + std::vector views; + for (const auto& instance : Game::im->GetInstances()) { + if (instance) views.push_back(instance->View()); + } + return views; + } + + void SetZoneMessage(uint32_t zone, const std::string& message) { + g_ZoneMessages[zone] = message; + g_Dirty = true; + } + + /** + * Replaces what Choose picked. warnSeconds: how long players are warned; requesterId: the GM told how each move + * goes. Returns how many instances are being replaced or stopped. + */ + uint32_t Apply(const std::vector& choices, uint16_t warnSeconds, LWOOBJID requesterId, const std::string& by) { + uint32_t acted = 0; + std::map> perZone; // zone -> (acted, skipped) + for (const auto& choice : choices) { + const auto& view = choice.view; + auto& counts = perZone[view.zoneId]; + switch (choice.action) { + case eAction::SKIP: + LOG("World reload (%s): zone %u instance %u left alone: %s", by.c_str(), view.zoneId, view.instanceId, choice.reason.c_str()); + counts.second++; + continue; + case eAction::STOP: + case eAction::START_THEN_STOP: { + const auto& instance = Game::im->FindInstanceWithPrivate(static_cast(view.zoneId), static_cast(view.instanceId)); + if (!instance) continue; + if (choice.action == eAction::START_THEN_STOP) { + // Started first: the stopped one takes nobody new meanwhile + instance->SetIsDraining(true); + const auto& started = Game::im->StartNewInstance(static_cast(view.zoneId), 0); + if (started) LOG("World reload (%s): started instance %u of zone %u", by.c_str(), started->GetInstanceID(), view.zoneId); + } + // Look it up again: starting one grew the instance list + const auto& stopping = Game::im->FindInstanceWithPrivate(static_cast(view.zoneId), static_cast(view.instanceId)); + if (!stopping) continue; + LOG("World reload (%s): zone %u instance %u is empty; stopping it", by.c_str(), view.zoneId, view.instanceId); + stopping->Shutdown(); + counts.first++; + acted++; + continue; + } + case eAction::REPLACE: { + InstanceMigrationRequest request; + request.requestId = MIGRATION_ID_BASE | (++g_NextMigrationId & 0x3FFFFFFF); + request.kind = InstanceMigration::eKind::REPLACE; + request.zoneId = view.zoneId; + request.sourceInstance = view.instanceId; + request.warnSeconds = warnSeconds; + request.shutdownSource = true; + request.seamless = Seamless(); + request.requesterId = requesterId; + request.requestedBy = ("world reload (" + by + ")").substr(0, InstanceMigrationRequest::MAX_BY); + MigrationCoordinator::Options options; + // Properties, private instances and activity zones are moved too; a property is saved and frozen first + options.liveUpdate = true; + options.prepare = view.cloneId != 0; + g_Migrations[request.requestId] = { view.zoneId, view.instanceId }; + const auto refusal = MigrationCoordinator::Start(request, options); + if (refusal != InstanceMigration::eRefusal::NONE) { + g_Migrations.erase(request.requestId); + LOG("World reload (%s): zone %u instance %u could not be replaced: %s", by.c_str(), view.zoneId, view.instanceId, InstanceMigration::Describe(refusal)); + counts.second++; + continue; + } + counts.first++; + acted++; + continue; + } + } + } + for (const auto& [zone, counts] : perZone) { + std::string message = "Reload (" + by + "): " + std::to_string(counts.first) + " instance(s) replaced or stopped"; + if (counts.second) message += ", " + std::to_string(counts.second) + " left alone"; + SetZoneMessage(zone, message); + } + return acted; + } + + // Instances that loaded a file that changed since, each tried once per version of its files + void ReloadStale() { + std::set wanted; + for (const auto& key : g_Tracker.Stale()) { + const auto tried = g_Tried.find(key); + if (tried != g_Tried.end() && tried->second == g_Tracker.Signature(key)) continue; + wanted.insert(key); + } + if (wanted.empty()) return; + auto choices = Choose(Views(), [&wanted](const InstanceMigration::InstanceView& view) { + return wanted.contains({ view.zoneId, view.instanceId }); + }, KeepZones()); + // Instances still starting (or already moving) are tried again at the next poll + std::erase_if(choices, [](const Choice& choice) { return choice.action == eAction::SKIP; }); + if (choices.empty()) return; + for (const auto& choice : choices) { + const InstanceKey key{ choice.view.zoneId, choice.view.instanceId }; + g_Tried[key] = g_Tracker.Signature(key); + std::string files; + for (const auto& path : g_Tracker.ChangedFiles(key)) files += (files.empty() ? "" : ", ") + FileName(path); + LOG("World reload: zone %u instance %u loaded files that changed on disk (%s): %s", choice.view.zoneId, choice.view.instanceId, + files.c_str(), ActionName(choice.action)); + } + Apply(choices, WorldReloadRequest::DEFAULT_WARN_SECONDS, LWOOBJID_EMPTY, "files changed"); + } + + WorldFilesStatus BuildStatus() { + WorldFilesStatus status; + const auto seconds = WatchSeconds(); + status.watching = seconds > 0; + status.watchSeconds = static_cast(std::min(seconds, UINT16_MAX)); + status.seamless = Seamless(); + for (const auto zoneId : g_Tracker.Zones()) { + WorldFilesStatus::Zone zone; + zone.zoneId = zoneId; + for (const auto& path : g_Tracker.FilesOf(zoneId)) { + const auto it = g_Tracker.Files().find(path); + if (it == g_Tracker.Files().end()) continue; + const auto& state = it->second; + WorldFilesStatus::File file; + file.disk.kind = state.kind; + file.disk.packed = state.packed; + file.disk.path = path; + file.hashed = state.hash.has_value(); + file.disk.size = state.hash ? state.size : state.reportedSize; + file.disk.hash = state.hash ? *state.hash : state.reportedHash; + file.missing = state.missing; + file.changed = g_Tracker.IsChanged(zoneId, path); + zone.files.push_back(std::move(file)); + } + for (const auto& instance : Game::im->GetInstances()) { + if (!instance || instance->GetMapID() != zoneId) continue; + WorldFilesStatus::Instance row; + row.instanceId = instance->GetInstanceID(); + row.cloneId = instance->GetCloneID(); + row.players = instance->GetCurrentClientCount(); + row.stale = g_Tracker.IsStale({ zoneId, row.instanceId }); + row.reloading = instance->GetIsShuttingDown() || instance->GetIsDraining(); + zone.instances.push_back(row); + } + if (const auto message = g_ZoneMessages.find(zoneId); message != g_ZoneMessages.end()) zone.message = message->second; + status.zones.push_back(std::move(zone)); + } + return status; + } + + void Publish(bool force) { + if (!g_Publisher) return; + const auto now = Clock::now(); + if (!force && (!g_Dirty || now < g_NextPublish)) return; + g_Dirty = false; + g_NextPublish = now + PUBLISH_INTERVAL; + auto status = BuildStatus(); + const bool same = status.watching == g_Published.watching && status.watchSeconds == g_Published.watchSeconds && + status.seamless == g_Published.seamless && status.zones == g_Published.zones; + if (same && !force) return; + g_Published = status; + g_Publisher(status); + } +} + +void WorldReloader::SetPublisher(std::function publisher) { + g_Publisher = std::move(publisher); +} + +void WorldReloader::HandleReport(const SystemAddress& from, const WorldFilesReport& report) { + const auto& instance = Game::im->GetInstanceBySysAddr(from); + if (!instance) { + LOG("Ignoring a world file list from a server that is not a world"); + return; + } + uint32_t watched = 0; + for (const auto& file : report.files) watched += file.packed ? 0 : 1; + LOG("Zone %u instance %u loaded %zu zone file(s) (%u loose, watched)", instance->GetMapID(), instance->GetInstanceID(), report.files.size(), watched); + g_Tracker.Report(instance->GetMapID(), instance->GetInstanceID(), instance->GetCloneID(), report.files); + g_Dirty = true; +} + +void WorldReloader::HandleRequest(const WorldReloadRequest& request, const std::string& who) { + if (request.zoneId == 0) { + LOG("World reload asked for by %s without a zone; nothing to do", who.c_str()); + return; + } + const auto zone = request.zoneId; + const auto choices = Choose(Views(), [zone](const InstanceMigration::InstanceView& view) { return view.zoneId == zone; }, KeepZones()); + if (choices.empty()) { + LOG("World reload of zone %u asked for by %s: no instance of it runs", zone, who.c_str()); + SetZoneMessage(zone, "Reload (" + who + "): no instance was running"); + return; + } + LOG("World reload of zone %u asked for by %s: %zu instance(s)", zone, who.c_str(), choices.size()); + for (const auto& choice : choices) { + // A replaced instance is not tried again automatically for the files on disk now + g_Tried[{ choice.view.zoneId, choice.view.instanceId }] = g_Tracker.Signature({ choice.view.zoneId, choice.view.instanceId }); + } + Apply(choices, std::min(request.warnSeconds, InstanceMigrationRequest::MAX_WARN_SECONDS), request.requesterId, who); +} + +void WorldReloader::OnMigrationStatus(const MigrationStatus& status) { + const auto it = g_Migrations.find(status.migrationId); + if (it == g_Migrations.end()) return; + if (!status.Finished()) { + g_Dirty = true; + return; + } + const auto key = it->second; + g_Migrations.erase(it); + std::string message = "Instance " + std::to_string(key.instance); + if (status.state == InstanceMigration::eState::DONE) { + message += " replaced by " + std::to_string(status.targetInstance) + " (" + std::to_string(status.moved) + " player(s) moved)"; + } else { + message += " was not replaced: " + (status.message.empty() ? std::string("the move failed") : status.message); + // Tried again once its files change again, or when asked + } + SetZoneMessage(key.zone, message); +} + +void WorldReloader::OnInstanceGone(const Instance& instance) { + const InstanceKey key{ instance.GetMapID(), instance.GetInstanceID() }; + g_Tracker.Forget(key.zone, key.instance); + g_Tried.erase(key); + if (g_Tracker.Zones().count(key.zone) == 0) g_ZoneMessages.erase(key.zone); + g_Dirty = true; +} + +void WorldReloader::Republish() { + g_Dirty = true; + Publish(true); +} + +void WorldReloader::Update() { + if (g_Job.valid()) { + if (g_Job.wait_for(std::chrono::seconds(0)) != std::future_status::ready) return Publish(false); + for (const auto& result : g_Job.get()) { + const bool changed = g_Tracker.Hashed(result.job, result.hash, result.size); + if (changed) LOG("World reload: %s changed on disk", result.job.path.c_str()); + else if (!result.hash) LOG("World reload: could not read %s", result.job.path.c_str()); + } + g_Dirty = true; + ReloadStale(); + } + + const auto now = Clock::now(); + if (now >= g_NextPoll) { + const auto seconds = WatchSeconds(); + g_NextPoll = now + std::chrono::seconds(seconds > 0 ? seconds : 60); + if (seconds > 0) { + // Worlds that reported after the last change are caught here too + ReloadStale(); + auto due = g_Tracker.Poll([](const std::string& path) { return FdbSnapshot::StampOf(path); }); + if (!due.empty()) g_Job = std::async(std::launch::async, HashFiles, std::move(due)); + } + // Player counts change without telling us; the status is compared before it is sent + g_Dirty = true; + } + Publish(false); +} + +void WorldReloader::Shutdown() { + if (g_Job.valid()) g_Job.wait(); +} diff --git a/dMasterServer/WorldReloader.h b/dMasterServer/WorldReloader.h new file mode 100644 index 000000000..9c59201c8 --- /dev/null +++ b/dMasterServer/WorldReloader.h @@ -0,0 +1,47 @@ +#ifndef WORLDRELOADER_H +#define WORLDRELOADER_H + +#include +#include + +#include "RakNetTypes.h" + +class Instance; +struct MigrationStatus; +struct WorldFilesReport; +struct WorldFilesStatus; +struct WorldReloadRequest; + +/** + * Master's side of world hot reload (docs/WorldHotReload.md). Worlds report the zone files they loaded (WORLD_FILES); + * master watches them (size and mtime every world_watch_seconds, then a hash on a worker thread) and, when one changes + * or a GM or the dashboard asks (WORLD_RELOAD), replaces the affected instances with new ones on the files on disk, + * moving their players with the instance migration (MigrationCoordinator). Main thread only. + */ +namespace WorldReloader { + // Where the dashboard's status goes (WORLD_FILES_STATUS) + void SetPublisher(std::function publisher); + + // A world reported its files + void HandleReport(const SystemAddress& from, const WorldFilesReport& report); + + // A GM's /reloadworld or the dashboard; who is for the log + void HandleRequest(const WorldReloadRequest& request, const std::string& who); + + // Every migration's progress (MigrationCoordinator's observer) + void OnMigrationStatus(const MigrationStatus& status); + + // A world server went away + void OnInstanceGone(const Instance& instance); + + // The dashboard connected: send it the status again + void Republish(); + + // Every master frame: polls the files, finishes hashes and reloads what changed + void Update(); + + // Waits for a running hash (at shutdown) + void Shutdown(); +}; + +#endif // WORLDRELOADER_H diff --git a/resources/masterconfig.ini b/resources/masterconfig.ini index 6fc6d901e..7d9b8ff3a 100644 --- a/resources/masterconfig.ini +++ b/resources/masterconfig.ini @@ -32,6 +32,11 @@ live_update_parallel_worlds=4 # Seconds between checks of the client's cdclient.fdb for changes (a change is reloaded on every server); 0 = off cdclient_watch_seconds=5 +# Seconds between checks of the zone files running worlds loaded; a change replaces those instances with new ones +# (players are moved over); 0 = off (/reloadworld still works) +world_watch_seconds=5 +# 1 = move players to a reloaded world without a loading screen (experimental seamless mode), 0 = Mythran shift +world_reload_seamless=0 # Seconds dead or building players may take before they are moved anyway live_update_player_wait=30 # Seconds builders on a property may take before it is saved for its new instance anyway diff --git a/tests/dCommonTests/CMakeLists.txt b/tests/dCommonTests/CMakeLists.txt index 71297f176..c9d2db8ba 100644 --- a/tests/dCommonTests/CMakeLists.txt +++ b/tests/dCommonTests/CMakeLists.txt @@ -37,6 +37,7 @@ set(DCOMMONTEST_SOURCES "Sd0Tests.cpp" "FdbReaderTests.cpp" "FdbSnapshotTests.cpp" + "WorldFileWatchTests.cpp" ) add_subdirectory(dEnumsTests) diff --git a/tests/dCommonTests/WorldFileWatchTests.cpp b/tests/dCommonTests/WorldFileWatchTests.cpp new file mode 100644 index 000000000..8f1c0c314 --- /dev/null +++ b/tests/dCommonTests/WorldFileWatchTests.cpp @@ -0,0 +1,252 @@ +#include + +#include + +#include "WorldFileWatch.h" + +using namespace WorldFileWatch; +using InstanceMigration::InstanceView; + +namespace { + using Kind = ZoneFileLog::eKind; + + ZoneFileLog::Entry File(const std::string& path, uint64_t hash, Kind kind = Kind::SCENE, bool packed = false) { + return { kind, packed, 100, hash, path }; + } + + Stamp At(uint64_t size, int64_t mtime) { + Stamp stamp; + stamp.exists = true; + stamp.size = size; + stamp.mtime = mtime; + return stamp; + } + + // The files on disk: path -> stamp + struct Disk { + std::map stamps; + std::vector Poll(Tracker& tracker) { + return tracker.Poll([this](const std::string& path) { + const auto it = stamps.find(path); + return it == stamps.end() ? Stamp{} : it->second; + }); + } + }; + + // Polls and finishes every hash with the given contents: path -> hash + std::vector PollAndHash(Tracker& tracker, Disk& disk, const std::map& contents, std::vector* changed = nullptr) { + std::vector hashed; + for (const auto& job : disk.Poll(tracker)) { + hashed.push_back(job.path); + const auto it = contents.find(job.path); + const bool didChange = tracker.Hashed(job, it == contents.end() ? std::nullopt : std::optional(it->second), 100); + if (changed && didChange) changed->push_back(job.path); + } + return hashed; + } + + InstanceView View(uint32_t zone, uint32_t instance, int32_t players, uint32_t clone = 0) { + InstanceView view; + view.zoneId = zone; + view.instanceId = instance; + view.cloneId = clone; + view.players = players; + view.ready = true; + return view; + } +} + +TEST(WorldFileWatchTests, NewFilesAreReadOnceAndMatchingInstancesAreNotStale) { + Tracker tracker; + Disk disk; + disk.stamps["/res/a.luz"] = At(10, 1); + disk.stamps["/res/a.lvl"] = At(20, 1); + tracker.Report(1100, 1, 0, { File("/res/a.luz", 1, Kind::ZONE), File("/res/a.lvl", 2) }); + // A second instance of the zone shares the files: they are read once + tracker.Report(1100, 2, 0, { File("/res/a.luz", 1, Kind::ZONE), File("/res/a.lvl", 2) }); + + EXPECT_EQ(PollAndHash(tracker, disk, { { "/res/a.luz", 1 }, { "/res/a.lvl", 2 } }).size(), 2u); + EXPECT_EQ(tracker.Files().size(), 2u); + EXPECT_TRUE(tracker.Stale().empty()); + // Nothing changed: nothing is read again + EXPECT_TRUE(PollAndHash(tracker, disk, { { "/res/a.luz", 1 }, { "/res/a.lvl", 2 } }).empty()); +} + +TEST(WorldFileWatchTests, ChangeMustHoldStillForOnePoll) { + Tracker tracker; + Disk disk; + disk.stamps["/res/a.lvl"] = At(20, 1); + tracker.Report(1100, 1, 0, { File("/res/a.lvl", 2) }); + PollAndHash(tracker, disk, { { "/res/a.lvl", 2 } }); + + // Being written: the stamp moves, so it isn't read yet + disk.stamps["/res/a.lvl"] = At(25, 2); + EXPECT_TRUE(PollAndHash(tracker, disk, { { "/res/a.lvl", 3 } }).empty()); + disk.stamps["/res/a.lvl"] = At(30, 3); + EXPECT_TRUE(PollAndHash(tracker, disk, { { "/res/a.lvl", 3 } }).empty()); + EXPECT_TRUE(tracker.Stale().empty()); + + // Held still for one poll: read, changed, and the instance is stale + std::vector changed; + EXPECT_EQ(PollAndHash(tracker, disk, { { "/res/a.lvl", 3 } }, &changed).size(), 1u); + EXPECT_EQ(changed, std::vector{ "/res/a.lvl" }); + ASSERT_EQ(tracker.Stale().size(), 1u); + EXPECT_EQ(tracker.Stale()[0], (InstanceKey{ 1100, 1 })); + EXPECT_TRUE(tracker.IsChanged(1100, "/res/a.lvl")); + EXPECT_EQ(tracker.ChangedFiles({ 1100, 1 }), std::vector{ "/res/a.lvl" }); +} + +TEST(WorldFileWatchTests, TouchedButSameContentIsNotAChange) { + Tracker tracker; + Disk disk; + disk.stamps["/res/a.lvl"] = At(20, 1); + tracker.Report(1100, 1, 0, { File("/res/a.lvl", 2) }); + PollAndHash(tracker, disk, { { "/res/a.lvl", 2 } }); + + disk.stamps["/res/a.lvl"] = At(20, 9); + PollAndHash(tracker, disk, { { "/res/a.lvl", 2 } }); + std::vector changed; + EXPECT_EQ(PollAndHash(tracker, disk, { { "/res/a.lvl", 2 } }, &changed).size(), 1u); + EXPECT_TRUE(changed.empty()); + EXPECT_TRUE(tracker.Stale().empty()); +} + +TEST(WorldFileWatchTests, NewInstanceOnTheNewFileIsNotStale) { + Tracker tracker; + Disk disk; + disk.stamps["/res/a.lvl"] = At(20, 1); + tracker.Report(1100, 1, 0, { File("/res/a.lvl", 2) }); + PollAndHash(tracker, disk, { { "/res/a.lvl", 2 } }); + disk.stamps["/res/a.lvl"] = At(30, 2); + PollAndHash(tracker, disk, { { "/res/a.lvl", 3 } }); + PollAndHash(tracker, disk, { { "/res/a.lvl", 3 } }); + + // The replacement loaded the new version; the old one is stale until it stops + tracker.Report(1100, 5, 0, { File("/res/a.lvl", 3) }); + EXPECT_EQ(tracker.Stale(), std::vector{ (InstanceKey{ 1100, 1 }) }); + tracker.Forget(1100, 1); + EXPECT_TRUE(tracker.Stale().empty()); + EXPECT_FALSE(tracker.IsChanged(1100, "/res/a.lvl")); +} + +TEST(WorldFileWatchTests, AWorldThatReadAnOlderVersionIsStaleOnceMasterReadsTheFile) { + Tracker tracker; + Disk disk; + disk.stamps["/res/a.lvl"] = At(20, 1); + // The file changed between the world's read and master's first read + tracker.Report(1100, 1, 0, { File("/res/a.lvl", 2) }); + EXPECT_TRUE(tracker.Stale().empty()); + PollAndHash(tracker, disk, { { "/res/a.lvl", 7 } }); + EXPECT_EQ(tracker.Stale().size(), 1u); +} + +TEST(WorldFileWatchTests, PackedFilesAreNotWatchedAndUnusedFilesAreDropped) { + Tracker tracker; + Disk disk; + disk.stamps["/res/b.luz"] = At(20, 1); + tracker.Report(1200, 1, 0, { File("maps/b.lvl", 2, Kind::SCENE, true), File("/res/b.luz", 1, Kind::ZONE) }); + EXPECT_EQ(PollAndHash(tracker, disk, { { "/res/b.luz", 1 } }), std::vector{ "/res/b.luz" }); + EXPECT_EQ(tracker.FilesOf(1200).size(), 2u); + + tracker.Forget(1200, 1); + EXPECT_TRUE(tracker.Files().empty()); + EXPECT_TRUE(tracker.Zones().empty()); +} + +TEST(WorldFileWatchTests, MissingFileIsNotAChangeUntilItComesBack) { + Tracker tracker; + Disk disk; + disk.stamps["/res/a.lvl"] = At(20, 1); + tracker.Report(1100, 1, 0, { File("/res/a.lvl", 2) }); + PollAndHash(tracker, disk, { { "/res/a.lvl", 2 } }); + + disk.stamps.erase("/res/a.lvl"); + EXPECT_TRUE(PollAndHash(tracker, disk, {}).empty()); + EXPECT_TRUE(PollAndHash(tracker, disk, {}).empty()); + EXPECT_TRUE(tracker.Files().at("/res/a.lvl").missing); + EXPECT_TRUE(tracker.Stale().empty()); + + disk.stamps["/res/a.lvl"] = At(22, 5); + PollAndHash(tracker, disk, { { "/res/a.lvl", 4 } }); + PollAndHash(tracker, disk, { { "/res/a.lvl", 4 } }); + EXPECT_FALSE(tracker.Files().at("/res/a.lvl").missing); + EXPECT_EQ(tracker.Stale().size(), 1u); +} + +TEST(WorldFileWatchTests, SignatureFollowsTheVersionOnDisk) { + Tracker tracker; + Disk disk; + disk.stamps["/res/a.lvl"] = At(20, 1); + tracker.Report(1100, 1, 0, { File("/res/a.lvl", 2) }); + PollAndHash(tracker, disk, { { "/res/a.lvl", 2 } }); + const auto before = tracker.Signature({ 1100, 1 }); + disk.stamps["/res/a.lvl"] = At(30, 2); + PollAndHash(tracker, disk, { { "/res/a.lvl", 3 } }); + PollAndHash(tracker, disk, { { "/res/a.lvl", 3 } }); + const auto after = tracker.Signature({ 1100, 1 }); + EXPECT_NE(before, after); + EXPECT_EQ(after, tracker.Signature({ 1100, 1 })); +} + +TEST(WorldFileWatchTests, ChooseReplacesBusyAndStopsEmptyInstances) { + std::vector views{ + View(1100, 1, 5), + View(1100, 2, 0), + View(1200, 3, 2), + View(1150, 4, 1, 77), // a property keeps its clone: replaced like any other + View(0, 5, 3), // character selection + }; + auto starting = View(1100, 6, 0); + starting.ready = false; + views.push_back(starting); + auto draining = View(1100, 7, 4); + draining.draining = true; + views.push_back(draining); + auto stopping = View(1100, 8, 0); + stopping.shuttingDown = true; + views.push_back(stopping); + + const auto choices = Choose(views, [](const InstanceView& view) { return view.zoneId != 1200; }, {}); + std::map byInstance; + for (const auto& choice : choices) byInstance[choice.view.instanceId] = choice.action; + EXPECT_EQ(byInstance.size(), 7u); // 1200 wasn't wanted + EXPECT_EQ(byInstance[1], eAction::REPLACE); + EXPECT_EQ(byInstance[2], eAction::STOP); + EXPECT_EQ(byInstance[4], eAction::REPLACE); + EXPECT_EQ(byInstance[5], eAction::SKIP); + EXPECT_EQ(byInstance[6], eAction::SKIP); + EXPECT_EQ(byInstance[7], eAction::SKIP); + EXPECT_EQ(byInstance[8], eAction::SKIP); + for (const auto& choice : choices) { + if (choice.action == eAction::SKIP) EXPECT_FALSE(choice.reason.empty()); + } +} + +TEST(WorldFileWatchTests, ChooseKeepsAnInstanceOfKeptZones) { + // Only empty public instances of a kept zone: one new instance is started for them + { + const std::vector views{ View(1000, 1, 0), View(1000, 2, 0) }; + const auto choices = Choose(views, [](const InstanceView&) { return true; }, { 1000 }); + ASSERT_EQ(choices.size(), 2u); + EXPECT_EQ(choices[0].action, eAction::START_THEN_STOP); + EXPECT_EQ(choices[1].action, eAction::STOP); + } + // A busy public instance is replaced anyway: the empty one just stops + { + const std::vector views{ View(1000, 1, 0), View(1000, 2, 3) }; + const auto choices = Choose(views, [](const InstanceView&) { return true; }, { 1000 }); + ASSERT_EQ(choices.size(), 2u); + EXPECT_EQ(choices[0].action, eAction::STOP); + EXPECT_EQ(choices[1].action, eAction::REPLACE); + } + // Not a kept zone, or a private instance: stopped + { + auto privateView = View(1000, 3, 0); + privateView.isPrivate = true; + const std::vector views{ View(1100, 1, 0), privateView }; + const auto choices = Choose(views, [](const InstanceView&) { return true; }, { 1000 }); + ASSERT_EQ(choices.size(), 2u); + EXPECT_EQ(choices[0].action, eAction::STOP); + EXPECT_EQ(choices[1].action, eAction::STOP); + } +}