diff --git a/dMasterServer/CMakeLists.txt b/dMasterServer/CMakeLists.txt index ec313ab0f..02e9ec12f 100644 --- a/dMasterServer/CMakeLists.txt +++ b/dMasterServer/CMakeLists.txt @@ -1,5 +1,6 @@ set(DMASTERSERVER_SOURCES "InstanceManager.cpp" + "MigrationCoordinator.cpp" "Start.cpp" ) diff --git a/dMasterServer/InstanceManager.cpp b/dMasterServer/InstanceManager.cpp index 6ddd32d5f..fb5747a62 100644 --- a/dMasterServer/InstanceManager.cpp +++ b/dMasterServer/InstanceManager.cpp @@ -31,6 +31,10 @@ const InstancePtr& InstanceManager::GetInstance(LWOMAPID mapID, bool isFriendTra auto& instance = FindInstance(mapID, isFriendTransfer, cloneID); if (instance) return instance; + return CreateInstance(mapID, cloneID); +} + +const InstancePtr& InstanceManager::CreateInstance(LWOMAPID mapID, LWOCLONEID cloneID) { // If we are shutting down, return a nullptr so a new instance is not created. if (m_IsShuttingDown) { LOG("Tried to create a new instance map/instance/clone %i/%i/%i, but Master is shutting down.", @@ -241,7 +245,7 @@ const InstancePtr& InstanceManager::GetInstanceBySysAddr(SystemAddress& sysAddr) const InstancePtr& InstanceManager::FindInstance(LWOMAPID mapID, bool isFriendTransfer, LWOCLONEID cloneId) { for (const auto& i : m_Instances) { - if (i && i->GetMapID() == mapID && i->GetCloneID() == cloneId && !i->IsFull(isFriendTransfer) && !i->GetIsPrivate() && !i->GetShutdownComplete() && !i->GetIsShuttingDown()) { + if (i && i->GetMapID() == mapID && i->GetCloneID() == cloneId && !i->IsFull(isFriendTransfer) && !i->GetIsPrivate() && !i->GetShutdownComplete() && !i->GetIsShuttingDown() && !i->GetIsDraining()) { return i; } } @@ -358,9 +362,11 @@ void Instance::Shutdown() { bool Instance::IsFull(bool isFriendTransfer) const { - if (!isFriendTransfer && GetSoftCap() > GetCurrentClientCount()) + // Seats held for players being moved in count as taken + const int load = GetCurrentClientCount() + GetReserved(); + if (!isFriendTransfer && GetSoftCap() > load) return false; - else if (isFriendTransfer && GetHardCap() > GetCurrentClientCount()) + else if (isFriendTransfer && GetHardCap() > load) return false; return true; diff --git a/dMasterServer/InstanceManager.h b/dMasterServer/InstanceManager.h index 5481efbca..80d6a8de5 100644 --- a/dMasterServer/InstanceManager.h +++ b/dMasterServer/InstanceManager.h @@ -1,4 +1,5 @@ #pragma once +#include #include #include "dCommonVars.h" #include "RakNetTypes.h" @@ -49,6 +50,12 @@ public: void SetIsReady(bool value) { m_Ready = value; } bool GetIsShuttingDown() const { return m_IsShuttingDown; } void SetIsShuttingDown(bool value) { m_IsShuttingDown = value; } + // Its players are being moved to another instance (InstanceMigration.h): nobody new is sent here + bool GetIsDraining() const { return m_IsDraining; } + void SetIsDraining(bool value) { m_IsDraining = value; } + // Seats held for players being moved in; they count towards the caps until the move is over + int GetReserved() const { return m_Reserved; } + void SetReserved(int value) { m_Reserved = std::max(0, value); } std::vector& GetPendingRequests() { return m_PendingRequests; } std::vector& GetPendingAffirmations() { return m_PendingAffirmations; } @@ -88,6 +95,8 @@ private: SystemAddress m_SysAddr{}; bool m_Ready{}; bool m_IsShuttingDown{}; + bool m_IsDraining{}; + int m_Reserved{}; std::vector m_PendingRequests{}; std::vector m_PendingAffirmations{}; @@ -135,6 +144,9 @@ public: void SetIsShuttingDown(bool value) { this->m_IsShuttingDown = value; }; void PruneUnreadyInstances(); + // Start a new public instance of a zone even when one with room is running (instance migrations) + const InstancePtr& StartNewInstance(LWOMAPID mapID, LWOCLONEID cloneID) { return CreateInstance(mapID, cloneID); } + private: std::string mExternalIP; std::vector> m_Instances; @@ -149,4 +161,5 @@ private: //Private functions: int GetSoftCap(LWOMAPID mapID); int GetHardCap(LWOMAPID mapID); + const InstancePtr& CreateInstance(LWOMAPID mapID, LWOCLONEID cloneID); }; diff --git a/dMasterServer/MasterServer.cpp b/dMasterServer/MasterServer.cpp index 240decdfe..a1e82e096 100644 --- a/dMasterServer/MasterServer.cpp +++ b/dMasterServer/MasterServer.cpp @@ -34,6 +34,7 @@ #include "AuthPackets.h" #include "Game.h" #include "InstanceManager.h" +#include "MigrationCoordinator.h" #include "MasterPackets.h" #include "FdbToSqlite.h" #include "BitStreamUtils.h" @@ -372,6 +373,16 @@ int main(int argc, char** argv) { return EXIT_FAILURE; } + // Instance migration progress goes to every world, where the GM who asked for it hears about it + MigrationCoordinator::SetReporter([](const MigrationStatus& status) { + CBITSTREAM; + BitStreamUtils::WriteHeader(bitStream, ServiceType::MASTER, MessageType::Master::MIGRATE_STATUS); + status.Serialize(bitStream); + for (const auto& instance : Game::im->GetInstances()) { + if (instance && instance->GetIsReady() && !instance->GetShutdownComplete()) Game::server->Send(bitStream, instance->GetSysAddr(), false); + } + }); + //Depending on the config, start up servers: if (Game::config->GetValue("prestart_servers") != "0") { StartChatServer(); @@ -403,6 +414,8 @@ int main(int argc, char** argv) { packet = nullptr; } + MigrationCoordinator::Update(); + //Push our log every 15s: if (framesSinceLastFlush >= logFlushTime) { Game::logger->Flush(); @@ -493,6 +506,7 @@ void HandlePacket(Packet* packet) { Game::im->GetInstanceBySysAddr(packet->systemAddress); if (instance) { LOG("Actually disconnected from zone %i clone %i instance %i port %i", instance->GetMapID(), instance->GetCloneID(), instance->GetInstanceID(), instance->GetPort()); + MigrationCoordinator::OnInstanceGone(*instance); Game::im->RemoveInstance(instance); //Delete the old } @@ -514,6 +528,7 @@ void HandlePacket(Packet* packet) { Game::im->GetInstanceBySysAddr(packet->systemAddress); if (instance) { LWOZONEID zoneID = instance->GetZoneID(); //Get the zoneID so we can recreate a server + MigrationCoordinator::OnInstanceGone(*instance); Game::im->RemoveInstance(instance); //Delete the old } @@ -834,6 +849,33 @@ void HandlePacket(Packet* packet) { break; } + case MessageType::Master::INSTANCE_MIGRATE: { + // Only servers connected to master can send this (a world, for a GM's /replaceinstance or /mergeinstance) + CINSTREAM_SKIP_HEADER; + InstanceMigrationRequest request; + if (!request.Deserialize(inStream)) break; + if (shutdownSequenceStarted) { + LOG("Shutdown sequence has been started. Not starting instance migration %u.", request.requestId); + break; + } + MigrationCoordinator::Start(request); + break; + } + + case MessageType::Master::MIGRATE_STATUS: { + CINSTREAM_SKIP_HEADER; + MigrationStatus status; + if (status.Deserialize(inStream)) MigrationCoordinator::HandleStatus(packet->systemAddress, status); + break; + } + + case MessageType::Master::MIGRATE_PLAYER_STATE: { + CINSTREAM_SKIP_HEADER; + CarriedPlayerState state; + if (state.Deserialize(inStream)) MigrationCoordinator::HandleCarriedState(packet->systemAddress, state, packet->data, packet->length); + break; + } + default: LOG("Unknown master packet ID from server: %i", packet->data[3]); } diff --git a/dMasterServer/MigrationCoordinator.cpp b/dMasterServer/MigrationCoordinator.cpp new file mode 100644 index 000000000..f7e5ca8c1 --- /dev/null +++ b/dMasterServer/MigrationCoordinator.cpp @@ -0,0 +1,313 @@ +#include "MigrationCoordinator.h" + +#include +#include + +#include "BitStreamUtils.h" +#include "CDActivitiesTable.h" +#include "CDClientManager.h" +#include "Game.h" +#include "InstanceManager.h" +#include "Logger.h" +#include "MessageType/Master.h" +#include "ServiceType.h" +#include "dServer.h" + +using namespace InstanceMigration; + +namespace { + using Clock = std::chrono::steady_clock; + // A new world server takes a few seconds to load its zone; big zones on slow disks take longer + constexpr auto TARGET_START_TIMEOUT = std::chrono::seconds(120); + // After the warning, moving everyone must be over by then (the source world gives up on stuck players sooner) + constexpr auto MOVE_TIMEOUT = std::chrono::seconds(180); + constexpr uint16_t PLAYERS_PER_SECOND = 10; + + struct Migration { + uint32_t id{}; + eKind kind{}; + eState state{}; + uint32_t zone{}; + uint32_t clone{}; + uint32_t source{}; + uint32_t target{}; + uint16_t warnSeconds{}; + bool shutdownSource{}; + bool startedTarget{}; // we started the target for this migration + bool seamless{}; + LWOOBJID requester{}; + int reserved{}; + uint16_t moved{}; + uint16_t failed{}; + Clock::time_point deadline{}; + std::string by; + }; + + std::map g_Active; + std::function g_Reporter; + + const InstancePtr& FindInstance(uint32_t zone, uint32_t instance) { + return Game::im->FindInstanceWithPrivate(static_cast(zone), static_cast(instance)); + } + + InstanceView View(const Instance& instance) { + InstanceView view; + view.zoneId = instance.GetMapID(); + view.instanceId = instance.GetInstanceID(); + view.cloneId = instance.GetCloneID(); + view.players = instance.GetCurrentClientCount(); + view.softCap = instance.GetSoftCap(); + view.hardCap = instance.GetHardCap(); + view.reserved = instance.GetReserved(); + view.ready = instance.GetIsReady(); + view.isPrivate = instance.GetIsPrivate(); + view.shuttingDown = instance.GetIsShuttingDown() || instance.GetShutdownComplete(); + view.draining = instance.GetIsDraining(); + return view; + } + + // Races, minigames and other activities run in their own zones with state that only lives in that world + bool IsActivityZone(uint32_t zone) { + auto* activities = CDClientManager::GetTable(); + if (!activities) return false; + return !activities->Query([zone](const CDActivities& activity) { return activity.instanceMapID == zone; }).empty(); + } + + bool IsTakingPart(uint32_t zone, uint32_t instance) { + for (const auto& [id, migration] : g_Active) { + if (migration.zone == zone && (migration.source == instance || migration.target == instance)) return true; + } + return false; + } + + void Report(const Migration& migration, eState state, const std::string& message, uint16_t remaining = 0) { + LOG("Migration %u (%s zone %u instance %u -> %u): %s%s%s", migration.id, KindName(migration.kind), migration.zone, migration.source, + migration.target, StateName(state), message.empty() ? "" : ": ", message.c_str()); + if (!g_Reporter) return; + MigrationStatus status; + status.migrationId = migration.id; + status.state = state; + status.kind = migration.kind; + status.zoneId = migration.zone; + status.sourceInstance = migration.source; + status.targetInstance = migration.target; + status.moved = migration.moved; + status.failed = migration.failed; + status.requesterId = migration.requester; + status.remaining = remaining; + status.message = message; + g_Reporter(status); + } + + void ReleaseSeats(Migration& migration) { + if (migration.reserved == 0) return; + if (const auto& target = FindInstance(migration.zone, migration.target)) target->SetReserved(target->GetReserved() - migration.reserved); + migration.reserved = 0; + } + + // Tell the source world to stop sending players (a target port of 0 cancels) + void SendCancel(const Migration& migration) { + const auto& source = FindInstance(migration.zone, migration.source); + if (!source) return; + MigratePlayersOrder order; + order.migrationId = migration.id; + order.targetZone = migration.zone; + order.targetInstance = migration.target; + order.targetPort = 0; + CBITSTREAM; + BitStreamUtils::WriteHeader(bitStream, ServiceType::MASTER, MessageType::Master::MIGRATE_PLAYERS); + order.Serialize(bitStream); + Game::server->Send(bitStream, source->GetSysAddr(), false); + } + + // Ends a migration. The source goes back to taking players unless it is being shut down. + void Finish(std::map::iterator it, bool success, const std::string& message) { + auto& migration = it->second; + ReleaseSeats(migration); + const auto& source = FindInstance(migration.zone, migration.source); + if (success && migration.shutdownSource && source) { + // Draining stays set: nothing is sent to it while it shuts down + source->Shutdown(); + } else if (source) { + source->SetIsDraining(false); + } + // A fresh instance nobody went to is shut down again (it would idle for half an hour otherwise) + if (!success && migration.startedTarget) { + const auto& target = FindInstance(migration.zone, migration.target); + if (target && target->GetCurrentClientCount() == 0) target->Shutdown(); + } + Report(migration, success ? eState::DONE : eState::FAILED, message); + g_Active.erase(it); + } + + void SendOrder(Migration& migration, const Instance& source, const Instance& target) { + MigratePlayersOrder order; + order.migrationId = migration.id; + order.targetZone = target.GetMapID(); + order.targetInstance = target.GetInstanceID(); + order.targetClone = target.GetCloneID(); + order.targetIp = target.GetIP(); + order.targetPort = static_cast(target.GetPort()); + order.warnSeconds = migration.warnSeconds; + order.playersPerSecond = PLAYERS_PER_SECOND; + // Without a loading screen the "dimensional shift" notice would be the only sign; leave it out then + order.seamless = migration.seamless; + order.mythranShift = !migration.seamless; + CBITSTREAM; + BitStreamUtils::WriteHeader(bitStream, ServiceType::MASTER, MessageType::Master::MIGRATE_PLAYERS); + order.Serialize(bitStream); + Game::server->Send(bitStream, source.GetSysAddr(), false); + + migration.state = migration.warnSeconds > 0 ? eState::WARNING : eState::MOVING; + migration.deadline = Clock::now() + std::chrono::seconds(migration.warnSeconds) + MOVE_TIMEOUT; + Report(migration, migration.state, migration.warnSeconds > 0 ? "Players were warned" : "Moving players", + static_cast(source.GetCurrentClientCount())); + } +} + +void MigrationCoordinator::SetReporter(std::function reporter) { + g_Reporter = std::move(reporter); +} + +eRefusal MigrationCoordinator::Start(const InstanceMigrationRequest& request) { + Migration migration; + migration.id = request.requestId; + migration.kind = request.kind; + migration.zone = request.zoneId; + migration.source = request.sourceInstance; + migration.warnSeconds = std::min(request.warnSeconds, InstanceMigrationRequest::MAX_WARN_SECONDS); + migration.shutdownSource = request.shutdownSource; + migration.seamless = request.seamless; + migration.requester = request.requesterId; + migration.by = request.requestedBy; + + const auto refuse = [&migration](eRefusal refusal) { + Report(migration, eState::FAILED, Describe(refusal)); + return refusal; + }; + + if (g_Active.contains(migration.id)) return refuse(eRefusal::ALREADY_MIGRATING); + // A raw pointer: starting an instance below grows the instance list, which moves the InstancePtrs + Instance* source = FindInstance(request.zoneId, request.sourceInstance).get(); + if (!source) return refuse(eRefusal::NOT_RUNNING); + auto sourceView = View(*source); + if (const auto refusal = CheckSource(sourceView, IsActivityZone(request.zoneId)); refusal != eRefusal::NONE) return refuse(refusal); + if (IsTakingPart(request.zoneId, request.sourceInstance)) return refuse(eRefusal::ALREADY_MIGRATING); + migration.clone = source->GetCloneID(); + + Instance* target = nullptr; + if (request.kind == eKind::MERGE) { + std::vector views; + for (const auto& instance : Game::im->GetInstances()) { + // Instances taking part in another migration don't count as targets + if (instance && !IsTakingPart(instance->GetMapID(), instance->GetInstanceID())) views.push_back(View(*instance)); + } + uint32_t targetId = request.targetInstance; + if (targetId == 0) { + const auto picked = PickMergeTarget(views, sourceView); + if (!picked) return refuse(eRefusal::NO_TARGET); + targetId = *picked; + } + const auto& found = FindInstance(request.zoneId, targetId); + if (!found) return refuse(eRefusal::NOT_RUNNING); + if (IsTakingPart(request.zoneId, targetId)) return refuse(eRefusal::ALREADY_MIGRATING); + if (const auto refusal = CheckMergeTarget(View(*found), sourceView); refusal != eRefusal::NONE) return refuse(refusal); + target = found.get(); + } else { + // A fresh world server, started from the binary on disk now: this is how a live update takes over + const auto& started = Game::im->StartNewInstance(source->GetMapID(), source->GetCloneID()); + if (!started) return refuse(eRefusal::MASTER_SHUTTING_DOWN); + target = started.get(); + migration.startedTarget = true; + } + + migration.target = target->GetInstanceID(); + // Seats for everyone there now; a few may still arrive while draining, the hard cap is the real limit + migration.reserved = source->GetCurrentClientCount(); + target->SetReserved(target->GetReserved() + migration.reserved); + source->SetIsDraining(true); + + LOG("Migration %u requested by %s: %s zone %u instance %u (%i player(s)) -> instance %u", migration.id, migration.by.c_str(), + KindName(migration.kind), migration.zone, migration.source, source->GetCurrentClientCount(), migration.target); + + auto& active = g_Active[migration.id] = migration; + if (target->GetIsReady()) { + SendOrder(active, *source, *target); + } else { + active.state = eState::STARTING_TARGET; + active.deadline = Clock::now() + TARGET_START_TIMEOUT; + Report(active, eState::STARTING_TARGET, "Starting instance " + std::to_string(active.target), + static_cast(source->GetCurrentClientCount())); + } + return eRefusal::NONE; +} + +void MigrationCoordinator::HandleStatus(const SystemAddress& from, const MigrationStatus& status) { + const auto it = g_Active.find(status.migrationId); + if (it == g_Active.end()) return; + auto& migration = it->second; + const auto& source = FindInstance(migration.zone, migration.source); + if (!source || source->GetSysAddr() != from) return; // only the source world runs it + + migration.moved = status.moved; + migration.failed = status.failed; + if (status.state == eState::DONE) return Finish(it, true, status.message); + if (status.state == eState::FAILED) return Finish(it, false, status.message); + migration.state = status.state; + Report(migration, status.state, status.message, status.remaining); +} + +void MigrationCoordinator::HandleCarriedState(const SystemAddress& from, const CarriedPlayerState& state, const unsigned char* data, uint32_t length) { + for (const auto& [id, migration] : g_Active) { + if (migration.zone != state.targetZone || migration.target != state.targetInstance) continue; + const auto& source = FindInstance(migration.zone, migration.source); + if (!source || source->GetSysAddr() != from) continue; + const auto& target = FindInstance(migration.zone, migration.target); + if (!target) return; + RakNet::BitStream forward(const_cast(data), length, false); + Game::server->Send(forward, target->GetSysAddr(), false); + return; + } +} + +void MigrationCoordinator::OnInstanceGone(const Instance& instance) { + for (auto it = g_Active.begin(); it != g_Active.end();) { + const auto& migration = it->second; + if (migration.zone != instance.GetMapID() || (migration.source != instance.GetInstanceID() && migration.target != instance.GetInstanceID())) { + ++it; + continue; + } + const bool wasSource = migration.source == instance.GetInstanceID(); + // Its players were saved when it stopped; with a shut down planned that is a finished drain + if (wasSource && migration.state == eState::MOVING && migration.shutdownSource) { + auto done = it++; + Finish(done, true, "The old instance stopped"); + continue; + } + if (!wasSource) SendCancel(migration); + auto failed = it++; + Finish(failed, false, wasSource ? "The instance being emptied stopped" : "The target instance stopped"); + } +} + +void MigrationCoordinator::Update() { + const auto now = Clock::now(); + for (auto it = g_Active.begin(); it != g_Active.end();) { + auto& migration = it->second; + auto current = it++; + if (migration.state == eState::STARTING_TARGET) { + const auto& source = FindInstance(migration.zone, migration.source); + const auto& target = FindInstance(migration.zone, migration.target); + if (!source || !target) { + Finish(current, false, !source ? "The instance being emptied stopped" : "The new instance stopped while starting"); + } else if (target->GetIsReady()) { + SendOrder(migration, *source, *target); + } else if (now > migration.deadline) { + Finish(current, false, "The new instance didn't start in time"); + } + } else if (now > migration.deadline) { + SendCancel(migration); + Finish(current, false, "Timed out moving players"); + } + } +} diff --git a/dMasterServer/MigrationCoordinator.h b/dMasterServer/MigrationCoordinator.h new file mode 100644 index 000000000..a2bb0f893 --- /dev/null +++ b/dMasterServer/MigrationCoordinator.h @@ -0,0 +1,37 @@ +#ifndef __MIGRATIONCOORDINATOR__H__ +#define __MIGRATIONCOORDINATOR__H__ + +#include + +#include "InstanceMigration.h" +#include "RakNetTypes.h" + +class Instance; + +/** + * Master's side of instance migrations (InstanceMigration.h): picks or starts the target instance, holds seats + * there, tells the source world to send its players over and passes progress on to whoever asked. The source + * instance is "draining" meanwhile, so no new player is sent to it. + */ +namespace MigrationCoordinator { + // Where statuses go (every world, so the GM who asked hears about it wherever they are); set once at startup + void SetReporter(std::function reporter); + + // A migration was asked for (a GM command, or later a dashboard); anything but NONE means it was refused (and + // reported as FAILED) + InstanceMigration::eRefusal Start(const InstanceMigrationRequest& request); + + // A world reported progress on a migration it is running + void HandleStatus(const SystemAddress& from, const MigrationStatus& status); + + // A source world sent state to carry over for one player: passed on to the target world as is + void HandleCarriedState(const SystemAddress& from, const CarriedPlayerState& state, const unsigned char* data, uint32_t length); + + // A world server went away; migrations it was part of end + void OnInstanceGone(const Instance& instance); + + // Every master frame: waits for targets to be ready and times out stuck migrations + void Update(); +} + +#endif //!__MIGRATIONCOORDINATOR__H__