From e1823bb2db751a7cc0a90a8543e778449ebf7d84 Mon Sep 17 00:00:00 2001 From: Keleborn <22352763+Celandriel@users.noreply.github.com> Date: Mon, 14 Sep 2026 00:20:44 -0700 Subject: [PATCH] feat(Core/Database): enable modules to own their database (#27544) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: UndeadRogue <254514851+UndeadRogue@users.noreply.github.com> Co-authored-by: UltraNix Co-authored-by: 郑佩茹 Co-authored-by: Yunfan Li Co-authored-by: kadeshar Co-authored-by: sudlud Co-authored-by: Claude Opus 5 (1M context) --- src/server/apps/worldserver/Main.cpp | 9 + .../database/Database/DatabaseUpdatePool.h | 37 +++ .../Database/DatabaseWorkerPoolAdapter.h | 49 ++++ .../database/Database/ModuleDatabasePool.cpp | 264 ++++++++++++++++++ .../database/Database/ModuleDatabasePool.h | 122 ++++++++ .../database/Database/MySQLConnection.h | 1 + src/server/database/Database/Transaction.h | 1 + src/server/database/Updater/DBUpdater.cpp | 171 +++++++++--- src/server/database/Updater/DBUpdater.h | 33 ++- .../ScriptDefines/DatabaseScript.cpp | 25 ++ .../Scripting/ScriptDefines/DatabaseScript.h | 39 +++ src/server/game/Scripting/ScriptMgr.h | 5 + src/server/game/World/World.cpp | 1 + src/server/scripts/Commands/cs_server.cpp | 7 + 14 files changed, 710 insertions(+), 54 deletions(-) create mode 100644 src/server/database/Database/DatabaseUpdatePool.h create mode 100644 src/server/database/Database/DatabaseWorkerPoolAdapter.h create mode 100644 src/server/database/Database/ModuleDatabasePool.cpp create mode 100644 src/server/database/Database/ModuleDatabasePool.h diff --git a/src/server/apps/worldserver/Main.cpp b/src/server/apps/worldserver/Main.cpp index 4dac150fd0..c445ad24ce 100644 --- a/src/server/apps/worldserver/Main.cpp +++ b/src/server/apps/worldserver/Main.cpp @@ -446,6 +446,9 @@ bool StartDB() if (!loader.Load()) return false; + if (!sScriptMgr->OnModuleDatabasesLoading()) + return false; + ///- Get the realm Id from the configuration file realm.Id.Realm = sConfigMgr->GetOption("RealmID", 1); if (!realm.Id.Realm) @@ -494,6 +497,8 @@ void StopDB() WorldDatabase.Close(); LoginDatabase.Close(); + sScriptMgr->OnModuleDatabasesClosing(); + MySQL::Library_End(); } @@ -584,6 +589,8 @@ void WorldUpdateLoop() CharacterDatabase.WarnAboutSyncQueries(true); WorldDatabase.WarnAboutSyncQueries(true); + sScriptMgr->OnDatabaseWarnAboutSyncQueries(true); + ///- While we have not World::m_stopEvent, update the world while (!World::IsStopped()) { @@ -613,6 +620,8 @@ void WorldUpdateLoop() #endif } + sScriptMgr->OnDatabaseWarnAboutSyncQueries(false); + LoginDatabase.WarnAboutSyncQueries(false); CharacterDatabase.WarnAboutSyncQueries(false); WorldDatabase.WarnAboutSyncQueries(false); diff --git a/src/server/database/Database/DatabaseUpdatePool.h b/src/server/database/Database/DatabaseUpdatePool.h new file mode 100644 index 0000000000..ed084dfb91 --- /dev/null +++ b/src/server/database/Database/DatabaseUpdatePool.h @@ -0,0 +1,37 @@ +/* + * This file is part of the AzerothCore Project. See AUTHORS file for Copyright information + * + * This program is free software; you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation; either version 2 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, but WITHOUT + * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or + * FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License for + * more details. + * + * You should have received a copy of the GNU General Public License along + * with this program. If not, see . + */ + +#ifndef DATABASE_UPDATE_POOL_H +#define DATABASE_UPDATE_POOL_H + +#include "Define.h" +#include "MySQLConnection.h" +#include "QueryResult.h" +#include + +// Minimal pool interface the DB updater operates on. Core pools reach it through +// DatabaseWorkerPoolAdapter; modules can implement it directly (see ModuleDatabasePool). +struct AC_DATABASE_API DatabaseUpdatePool +{ + virtual ~DatabaseUpdatePool() = default; + + virtual void DirectExecute(std::string_view query) = 0; + virtual QueryResult Query(std::string_view query) = 0; + virtual MySQLConnectionInfo const* GetConnectionInfo() const = 0; +}; + +#endif diff --git a/src/server/database/Database/DatabaseWorkerPoolAdapter.h b/src/server/database/Database/DatabaseWorkerPoolAdapter.h new file mode 100644 index 0000000000..d277a488e2 --- /dev/null +++ b/src/server/database/Database/DatabaseWorkerPoolAdapter.h @@ -0,0 +1,49 @@ +/* + * This file is part of the AzerothCore Project. See AUTHORS file for Copyright information + * + * This program is free software; you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation; either version 2 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, but WITHOUT + * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or + * FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License for + * more details. + * + * You should have received a copy of the GNU General Public License along + * with this program. If not, see . + */ + +#ifndef DATABASE_WORKER_POOL_ADAPTER_H +#define DATABASE_WORKER_POOL_ADAPTER_H + +#include "DatabaseUpdatePool.h" +#include "DatabaseWorkerPool.h" + +template +class DatabaseWorkerPoolAdapter : public DatabaseUpdatePool +{ +public: + DatabaseWorkerPoolAdapter(DatabaseWorkerPool& pool) : _pool(pool) {} + + void DirectExecute(std::string_view query) override + { + _pool.DirectExecute(query); + } + + QueryResult Query(std::string_view query) override + { + return _pool.Query(query); + } + + MySQLConnectionInfo const* GetConnectionInfo() const override + { + return _pool.GetConnectionInfo(); + } + +private: + DatabaseWorkerPool& _pool; +}; + +#endif diff --git a/src/server/database/Database/ModuleDatabasePool.cpp b/src/server/database/Database/ModuleDatabasePool.cpp new file mode 100644 index 0000000000..9d71e78d86 --- /dev/null +++ b/src/server/database/Database/ModuleDatabasePool.cpp @@ -0,0 +1,264 @@ +/* + * This file is part of the AzerothCore Project. See AUTHORS file for Copyright information + * + * This program is free software; you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation; either version 2 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, but WITHOUT + * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or + * FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License for + * more details. + * + * You should have received a copy of the GNU General Public License along + * with this program. If not, see . + */ + +#include "ModuleDatabasePool.h" +#include "Errors.h" +#include "Log.h" +#include "MySQLConnection.h" +#include "MySQLPreparedStatement.h" +#include "QueryResult.h" +#include "Transaction.h" +#include +#include +#include +#include + +ModuleDatabasePool::ModuleDatabasePool() + : _connectionInfo(""), _synchThreads(0) +{ +} + +ModuleDatabasePool::~ModuleDatabasePool() +{ + Close(); +} + +void ModuleDatabasePool::SetConnectionInfo(std::string_view infoString, uint8 synchThreads) +{ + _connectionInfo = MySQLConnectionInfo(infoString); + _synchThreads = synchThreads; +} + +uint32 ModuleDatabasePool::Open() +{ + if (!_synchThreads) + { + LOG_ERROR("sql.driver", "ModuleDatabasePool: database `{}` was configured with 0 synchronous connections, " + "at least one is required.", _connectionInfo.database); + return CR_UNKNOWN_ERROR; + } + + Close(); + + for (uint8 i = 0; i < _synchThreads; ++i) + { + auto conn = std::unique_ptr(CreateConnection(_connectionInfo)); + uint32 result = conn->Open(); + if (result != 0) + { + LOG_ERROR("sql.driver", "ModuleDatabasePool: could not open connection {}/{} to database `{}`, error {}", + i + 1, _synchThreads, _connectionInfo.database, result); + Close(); + return result; + } + + _connections.push_back(std::move(conn)); + } + + return 0; +} + +bool ModuleDatabasePool::PrepareStatements() +{ + for (auto const& conn : _connections) + { + conn->LockIfReady(); + if (!conn->PrepareStatements()) + { + conn->Unlock(); + Close(); + return false; + } + + conn->Unlock(); + } + + if (!_connections.empty()) + { + MySQLConnection const* conn = _connections.front().get(); + _preparedStatementSize.assign(conn->m_stmts.size(), 0); + for (std::size_t i = 0; i < conn->m_stmts.size(); ++i) + { + if (MySQLPreparedStatement* stmt = conn->m_stmts[i].get()) + { + uint32 const paramCount = stmt->GetParameterCount(); + ASSERT(paramCount < std::numeric_limits::max()); + _preparedStatementSize[i] = static_cast(paramCount); + } + } + } + + return true; +} + +void ModuleDatabasePool::Close() +{ + _connections.clear(); + + _preparedStatementSize.clear(); +} + +void ModuleDatabasePool::Execute(std::string_view sql) +{ + // Synchronous for now - kept separate from DirectExecute so async execution + // can be added later without touching callers. + DirectExecute(sql); +} + +void ModuleDatabasePool::DirectExecute(std::string_view sql) +{ + if (sql.empty()) + return; + + if (_connections.empty()) + return; + + MySQLConnection* conn = GetFreeConnection(); + conn->Execute(sql); + conn->Unlock(); +} + +QueryResult ModuleDatabasePool::Query(std::string_view sql) +{ + if (_connections.empty()) + return QueryResult(nullptr); + + MySQLConnection* conn = GetFreeConnection(); + ResultSet* result = conn->Query(sql); + conn->Unlock(); + + // Mirror DatabaseWorkerPool::Query semantics: nullptr for empty results, + // and the first row loaded before the result is handed out. + if (!result || !result->GetRowCount() || !result->NextRow()) + { + delete result; + return QueryResult(nullptr); + } + + return QueryResult(result); +} + +void ModuleDatabasePool::Execute(PreparedStatementBase* stmt) +{ + if (_connections.empty()) + { + delete stmt; + return; + } + + MySQLConnection* conn = GetFreeConnection(); + conn->Execute(stmt); + conn->Unlock(); + + delete stmt; +} + +PreparedQueryResult ModuleDatabasePool::Query(PreparedStatementBase* stmt) +{ + if (_connections.empty()) + { + delete stmt; + return PreparedQueryResult(nullptr); + } + + MySQLConnection* conn = GetFreeConnection(); + PreparedResultSet* result = conn->Query(stmt); + conn->Unlock(); + + //! Delete proxy-class. Not needed anymore + delete stmt; + + if (!result || !result->GetRowCount()) + { + delete result; + return PreparedQueryResult(nullptr); + } + + return PreparedQueryResult(result); +} + +uint8 ModuleDatabasePool::GetPreparedStatementParamCount(uint32 index) const +{ + return index < _preparedStatementSize.size() ? _preparedStatementSize[index] : 0; +} + +void ModuleDatabasePool::DirectCommitTransaction(std::shared_ptr transaction) +{ + if (_connections.empty()) + return; + + MySQLConnection* conn = GetFreeConnection(); + int errorCode = conn->ExecuteTransaction(transaction); + if (!errorCode) + { + conn->Unlock(); + return; + } + + //! Handle MySQL Errno 1213 without extending deadlock to the core itself + if (errorCode == ER_LOCK_DEADLOCK) + { + uint8 constexpr loopBreaker = 5; + for (uint8 i = 0; i < loopBreaker; ++i) + { + if (!conn->ExecuteTransaction(transaction)) + break; + } + } + + transaction->Cleanup(); + conn->Unlock(); +} + +void ModuleDatabasePool::KeepAlive() +{ + //! Ping connections that are not busy; a locked connection is in use and alive. + for (auto const& conn : _connections) + { + if (conn->LockIfReady()) + { + conn->Ping(); + conn->Unlock(); + } + } +} + +MySQLConnection* ModuleDatabasePool::GetFreeConnection() +{ + uint8 i = 0; + auto const num_cons = _connections.size(); + MySQLConnection* connection = nullptr; + + //! Block forever until a connection is free + for (;;) + { + connection = _connections[++i % num_cons].get(); + //! Must be matched with connection->Unlock() or you will get deadlocks + if (connection->LockIfReady()) + break; + + if (i % num_cons == 0) + std::this_thread::yield(); + } + + return connection; +} + +MySQLConnectionInfo const* ModuleDatabasePool::GetConnectionInfo() const +{ + return &_connectionInfo; +} diff --git a/src/server/database/Database/ModuleDatabasePool.h b/src/server/database/Database/ModuleDatabasePool.h new file mode 100644 index 0000000000..ba0bb168e2 --- /dev/null +++ b/src/server/database/Database/ModuleDatabasePool.h @@ -0,0 +1,122 @@ +/* + * This file is part of the AzerothCore Project. See AUTHORS file for Copyright information + * + * This program is free software; you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation; either version 2 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, but WITHOUT + * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or + * FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License for + * more details. + * + * You should have received a copy of the GNU General Public License along + * with this program. If not, see . + */ + +#ifndef MODULE_DATABASE_POOL_H +#define MODULE_DATABASE_POOL_H + +#include "DatabaseEnvFwd.h" +#include "DatabaseUpdatePool.h" +#include "Define.h" +#include "MySQLConnection.h" +#include "PreparedStatement.h" +#include "StringFormat.h" +#include +#include +#include + +class TransactionBase; + +// Base class for module-owned database pools. A module derives from this, +// implements CreateConnection with its own MySQLConnection subclass (carrying +// the module's prepared statements), and gets open/execute/query plus DBUpdater +// compatibility without any core-side registration. +// +// Connections are synchronous; asynchronous execution can be added in a +// follow-up without changing this interface. DoPrepareStatements must mark every +// statement CONNECTION_SYNCH: a CONNECTION_ASYNC one is skipped on these +// connections and asserts on first use. +class AC_DATABASE_API ModuleDatabasePool : public DatabaseUpdatePool +{ +public: + ModuleDatabasePool(); + virtual ~ModuleDatabasePool(); + + void SetConnectionInfo(std::string_view infoString, uint8 synchThreads); + + //! Opens the configured number of synchronous connections. + //! Returns 0 on success, or the MySQL error code of the first failed connection. + uint32 Open(); + + //! Prepares the connection statements. Call after the schema exists + //! (post create/populate/update), mirroring DatabaseLoader's ordering. + bool PrepareStatements(); + + void Close(); + + void Execute(std::string_view sql); + void DirectExecute(std::string_view sql) override; + QueryResult Query(std::string_view sql) override; + MySQLConnectionInfo const* GetConnectionInfo() const override; + + //! Format variants, mirroring DatabaseWorkerPool. + template + void Execute(std::string_view sql, Args&&... args) + { + if (sql.empty()) + return; + + Execute(std::string_view(Acore::StringFormat(sql, std::forward(args)...))); + } + + template + void DirectExecute(std::string_view sql, Args&&... args) + { + if (sql.empty()) + return; + + DirectExecute(std::string_view(Acore::StringFormat(sql, std::forward(args)...))); + } + + template + QueryResult Query(std::string_view sql, Args&&... args) + { + if (sql.empty()) + return QueryResult(nullptr); + + return Query(std::string_view(Acore::StringFormat(sql, std::forward(args)...))); + } + + //! Prepared statements. The index space is defined by the module's connection + //! class (DoPrepareStatements); parameter counts are recorded by PrepareStatements(), + //! so building one before that call yields a zero-parameter statement. + //! Both calls consume (delete) the statement, mirroring DatabaseWorkerPool. + void Execute(PreparedStatementBase* stmt); + PreparedQueryResult Query(PreparedStatementBase* stmt); + + //! Parameter count for a prepared statement index, for constructing typed + //! PreparedStatement objects module-side. + [[nodiscard]] uint8 GetPreparedStatementParamCount(uint32 index) const; + + //! Synchronously commits the transaction on a free connection. + void DirectCommitTransaction(std::shared_ptr transaction); + + //! Pings every idle connection to keep them alive. + void KeepAlive(); + +protected: + virtual MySQLConnection* CreateConnection(MySQLConnectionInfo& connInfo) = 0; + +private: + MySQLConnection* GetFreeConnection(); + + MySQLConnectionInfo _connectionInfo; + std::vector> _connections; + std::vector _preparedStatementSize; + uint8 _synchThreads; +}; + +#endif diff --git a/src/server/database/Database/MySQLConnection.h b/src/server/database/Database/MySQLConnection.h index e1b51836af..4fe8eec279 100644 --- a/src/server/database/Database/MySQLConnection.h +++ b/src/server/database/Database/MySQLConnection.h @@ -56,6 +56,7 @@ class AC_DATABASE_API MySQLConnection template friend class DatabaseWorkerPool; +friend class ModuleDatabasePool; friend class PingOperation; public: diff --git a/src/server/database/Database/Transaction.h b/src/server/database/Database/Transaction.h index e3417cff5b..5e6f350feb 100644 --- a/src/server/database/Database/Transaction.h +++ b/src/server/database/Database/Transaction.h @@ -31,6 +31,7 @@ class AC_DATABASE_API TransactionBase { friend class TransactionTask; friend class MySQLConnection; + friend class ModuleDatabasePool; template friend class DatabaseWorkerPool; diff --git a/src/server/database/Updater/DBUpdater.cpp b/src/server/database/Updater/DBUpdater.cpp index 554e568d84..a101da023c 100644 --- a/src/server/database/Updater/DBUpdater.cpp +++ b/src/server/database/Updater/DBUpdater.cpp @@ -20,6 +20,7 @@ #include "Config.h" #include "DatabaseEnv.h" #include "DatabaseLoader.h" +#include "DatabaseWorkerPoolAdapter.h" #include "Log.h" #include "StartProcess.h" #include "UpdateFetcher.h" @@ -93,10 +94,16 @@ std::string DBUpdater::GetTableName() return "Auth"; } +template<> +std::string DBUpdater::GetSourceDirectory() +{ + return BuiltInConfig::GetSourceDirectory(); +} + template<> std::string DBUpdater::GetBaseFilesDirectory() { - return BuiltInConfig::GetSourceDirectory() + "/data/sql/base/db_auth/"; + return DBUpdater::GetSourceDirectory() + "/data/sql/base/db_auth/"; } template<> @@ -126,10 +133,16 @@ std::string DBUpdater::GetTableName() return "World"; } +template<> +std::string DBUpdater::GetSourceDirectory() +{ + return BuiltInConfig::GetSourceDirectory(); +} + template<> std::string DBUpdater::GetBaseFilesDirectory() { - return BuiltInConfig::GetSourceDirectory() + "/data/sql/base/db_world/"; + return DBUpdater::GetSourceDirectory() + "/data/sql/base/db_world/"; } template<> @@ -159,10 +172,16 @@ std::string DBUpdater::GetTableName() return "Character"; } +template<> +std::string DBUpdater::GetSourceDirectory() +{ + return BuiltInConfig::GetSourceDirectory(); +} + template<> std::string DBUpdater::GetBaseFilesDirectory() { - return BuiltInConfig::GetSourceDirectory() + "/data/sql/base/db_characters/"; + return DBUpdater::GetSourceDirectory() + "/data/sql/base/db_characters/"; } template<> @@ -186,8 +205,19 @@ BaseLocation DBUpdater::GetBaseLocationType() return LOCATION_REPOSITORY; } -template -bool DBUpdater::Create(DatabaseWorkerPool& pool) +namespace +{ + +using Path = std::filesystem::path; + +QueryResult Retrieve(DatabaseUpdatePool& pool, std::string const& query); +void Apply(DatabaseUpdatePool& pool, std::string const& query); +void ApplyFile(DatabaseUpdatePool& pool, Path const& path); +void ApplyFile(DatabaseUpdatePool& pool, std::string const& host, std::string const& user, + std::string const& password, std::string const& port_or_socket, std::string const& database, + std::string const& ssl, Path const& path); + +bool CreateDatabase(DatabaseUpdatePool& pool) { LOG_WARN("sql.updates", "Database \"{}\" does not exist", pool.GetConnectionInfo()->database); @@ -220,8 +250,9 @@ bool DBUpdater::Create(DatabaseWorkerPool& pool) try { - DBUpdater::ApplyFile(pool, pool.GetConnectionInfo()->host, pool.GetConnectionInfo()->user, pool.GetConnectionInfo()->password, - pool.GetConnectionInfo()->port_or_socket, "", pool.GetConnectionInfo()->ssl, temp); + ApplyFile(pool, pool.GetConnectionInfo()->host, pool.GetConnectionInfo()->user, + pool.GetConnectionInfo()->password, pool.GetConnectionInfo()->port_or_socket, "", + pool.GetConnectionInfo()->ssl, temp); } catch (UpdateException&) { @@ -236,15 +267,14 @@ bool DBUpdater::Create(DatabaseWorkerPool& pool) return true; } -template -bool DBUpdater::Update(DatabaseWorkerPool& pool, std::string_view modulesList /*= {}*/) +bool UpdateDatabase(DatabaseUpdatePool& pool, DBUpdaterInfo const& info, std::string_view modulesList) { if (!DBUpdaterUtil::CheckExecutable()) return false; - LOG_INFO("sql.updates", "Updating {} database...", DBUpdater::GetTableName()); + LOG_INFO("sql.updates", "Updating {} database...", info.displayName); - Path const sourceDirectory(BuiltInConfig::GetSourceDirectory()); + Path const sourceDirectory(info.sourceDirectory); if (!is_directory(sourceDirectory)) { @@ -255,16 +285,16 @@ bool DBUpdater::Update(DatabaseWorkerPool& pool, std::string_view modulesL auto CheckUpdateTable = [&](std::string const& tableName) { - auto checkTable = DBUpdater::Retrieve(pool, Acore::StringFormat("SHOW TABLES LIKE '{}'", tableName)); + auto checkTable = Retrieve(pool, Acore::StringFormat("SHOW TABLES LIKE '{}'", tableName)); if (!checkTable) { LOG_WARN("sql.updates", "> Table '{}' not exist! Try add based table", tableName); - Path const temp(GetBaseFilesDirectory() + tableName + ".sql"); + Path const temp = Path(info.baseFilesDirectory) / (tableName + ".sql"); try { - DBUpdater::ApplyFile(pool, temp); + ApplyFile(pool, temp); } catch (UpdateException&) { @@ -281,9 +311,9 @@ bool DBUpdater::Update(DatabaseWorkerPool& pool, std::string_view modulesL if (!CheckUpdateTable("updates") || !CheckUpdateTable("updates_include")) return false; - UpdateFetcher updateFetcher(sourceDirectory, [&](std::string const & query) { DBUpdater::Apply(pool, query); }, - [&](Path const & file) { DBUpdater::ApplyFile(pool, file); }, - [&](std::string const & query) -> QueryResult { return DBUpdater::Retrieve(pool, query); }, DBUpdater::GetDBModuleName(), modulesList); + UpdateFetcher updateFetcher(sourceDirectory, [&](std::string const & query) { Apply(pool, query); }, + [&](Path const & file) { ApplyFile(pool, file); }, + [&](std::string const & query) -> QueryResult { return Retrieve(pool, query); }, info.dbModuleName, modulesList); UpdateResult result; try @@ -299,27 +329,27 @@ bool DBUpdater::Update(DatabaseWorkerPool& pool, std::string_view modulesL return false; } - std::string const info = Acore::StringFormat("Containing {} new and {} archived updates.", result.recent, result.archived); + std::string const summary = Acore::StringFormat("Containing {} new and {} archived updates.", result.recent, result.archived); if (!result.updated) - LOG_INFO("sql.updates", ">> {} database is up-to-date! {}", DBUpdater::GetTableName(), info); + LOG_INFO("sql.updates", ">> {} database is up-to-date! {}", info.displayName, summary); else - LOG_INFO("sql.updates", ">> Applied {} {}. {}", result.updated, result.updated == 1 ? "query" : "queries", info); + LOG_INFO("sql.updates", ">> Applied {} {}. {}", result.updated, result.updated == 1 ? "query" : "queries", summary); LOG_INFO("sql.updates", " "); return true; } -template -bool DBUpdater::Update(DatabaseWorkerPool& pool, std::vector const* setDirectories) +bool UpdateDatabase(DatabaseUpdatePool& pool, DBUpdaterInfo const& info, + std::vector const* setDirectories) { if (!DBUpdaterUtil::CheckExecutable()) { return false; } - Path const sourceDirectory(BuiltInConfig::GetSourceDirectory()); + Path const sourceDirectory(info.sourceDirectory); if (!is_directory(sourceDirectory)) { return false; @@ -327,13 +357,13 @@ bool DBUpdater::Update(DatabaseWorkerPool& pool, std::vector auto CheckUpdateTable = [&](std::string const& tableName) { - auto checkTable = DBUpdater::Retrieve(pool, Acore::StringFormat("SHOW TABLES LIKE '{}'", tableName)); + auto checkTable = Retrieve(pool, Acore::StringFormat("SHOW TABLES LIKE '{}'", tableName)); if (!checkTable) { - Path const temp(GetBaseFilesDirectory() + tableName + ".sql"); + Path const temp = Path(info.baseFilesDirectory) / (tableName + ".sql"); try { - DBUpdater::ApplyFile(pool, temp); + ApplyFile(pool, temp); } catch (UpdateException&) { @@ -351,9 +381,9 @@ bool DBUpdater::Update(DatabaseWorkerPool& pool, std::vector return false; } - UpdateFetcher updateFetcher(sourceDirectory, [&](std::string const & query) { DBUpdater::Apply(pool, query); }, - [&](Path const & file) { DBUpdater::ApplyFile(pool, file); }, - [&](std::string const & query) -> QueryResult { return DBUpdater::Retrieve(pool, query); }, DBUpdater::GetDBModuleName(), setDirectories); + UpdateFetcher updateFetcher(sourceDirectory, [&](std::string const & query) { Apply(pool, query); }, + [&](Path const & file) { ApplyFile(pool, file); }, + [&](std::string const & query) -> QueryResult { return Retrieve(pool, query); }, info.dbModuleName, setDirectories); UpdateResult result; try @@ -372,8 +402,7 @@ bool DBUpdater::Update(DatabaseWorkerPool& pool, std::vector return true; } -template -bool DBUpdater::Populate(DatabaseWorkerPool& pool) +bool PopulateDatabase(DatabaseUpdatePool& pool, DBUpdaterInfo const& info) { { QueryResult const result = Retrieve(pool, "SHOW TABLES"); @@ -384,9 +413,9 @@ bool DBUpdater::Populate(DatabaseWorkerPool& pool) if (!DBUpdaterUtil::CheckExecutable()) return false; - LOG_INFO("sql.updates", "Database {} is empty, auto populating it...", DBUpdater::GetTableName()); + LOG_INFO("sql.updates", "Database {} is empty, auto populating it...", info.displayName); - std::string const DirPathStr = DBUpdater::GetBaseFilesDirectory(); + std::string const DirPathStr = info.baseFilesDirectory; Path const DirPath(DirPathStr); if (!std::filesystem::is_directory(DirPath)) @@ -445,28 +474,26 @@ bool DBUpdater::Populate(DatabaseWorkerPool& pool) return true; } -template -QueryResult DBUpdater::Retrieve(DatabaseWorkerPool& pool, std::string const& query) +QueryResult Retrieve(DatabaseUpdatePool& pool, std::string const& query) { - return pool.Query(query.c_str()); + return pool.Query(query); } -template -void DBUpdater::Apply(DatabaseWorkerPool& pool, std::string const& query) +void Apply(DatabaseUpdatePool& pool, std::string const& query) { - pool.DirectExecute(query.c_str()); + pool.DirectExecute(query); } -template -void DBUpdater::ApplyFile(DatabaseWorkerPool& pool, Path const& path) +void ApplyFile(DatabaseUpdatePool& pool, Path const& path) { - DBUpdater::ApplyFile(pool, pool.GetConnectionInfo()->host, pool.GetConnectionInfo()->user, pool.GetConnectionInfo()->password, - pool.GetConnectionInfo()->port_or_socket, pool.GetConnectionInfo()->database, pool.GetConnectionInfo()->ssl, path); + ApplyFile(pool, pool.GetConnectionInfo()->host, pool.GetConnectionInfo()->user, + pool.GetConnectionInfo()->password, pool.GetConnectionInfo()->port_or_socket, + pool.GetConnectionInfo()->database, pool.GetConnectionInfo()->ssl, path); } -template -void DBUpdater::ApplyFile(DatabaseWorkerPool& pool, std::string const& host, std::string const& user, - std::string const& password, std::string const& port_or_socket, std::string const& database, std::string const& ssl, Path const& path) +void ApplyFile(DatabaseUpdatePool& pool, std::string const& host, std::string const& user, + std::string const& password, std::string const& port_or_socket, std::string const& database, + std::string const& ssl, Path const& path) { std::string configTempDir = sConfigMgr->GetOption("TempDir", ""); @@ -554,6 +581,58 @@ void DBUpdater::ApplyFile(DatabaseWorkerPool& pool, std::string const& hos } } +} // anonymous namespace + +template +DBUpdaterInfo DBUpdater::GetUpdaterInfo() +{ + return { GetTableName(), GetSourceDirectory(), GetBaseFilesDirectory(), GetDBModuleName() }; +} + +template +bool DBUpdater::Create(DatabaseWorkerPool& pool) +{ + DatabaseWorkerPoolAdapter adapter(pool); + return CreateDatabase(adapter); +} + +template +bool DBUpdater::Update(DatabaseWorkerPool& pool, std::string_view modulesList /*= {}*/) +{ + DatabaseWorkerPoolAdapter adapter(pool); + return UpdateDatabase(adapter, GetUpdaterInfo(), modulesList); +} + +template +bool DBUpdater::Update(DatabaseWorkerPool& pool, std::vector const* setDirectories) +{ + DatabaseWorkerPoolAdapter adapter(pool); + return UpdateDatabase(adapter, GetUpdaterInfo(), setDirectories); +} + +template +bool DBUpdater::Populate(DatabaseWorkerPool& pool) +{ + DatabaseWorkerPoolAdapter adapter(pool); + return PopulateDatabase(adapter, GetUpdaterInfo()); +} + +bool ModuleDBUpdater::Create(DatabaseUpdatePool& pool) +{ + return CreateDatabase(pool); +} + +bool ModuleDBUpdater::Update(DatabaseUpdatePool& pool, DBUpdaterInfo const& info, + std::string_view modulesList /*= {}*/) +{ + return UpdateDatabase(pool, info, modulesList); +} + +bool ModuleDBUpdater::Populate(DatabaseUpdatePool& pool, DBUpdaterInfo const& info) +{ + return PopulateDatabase(pool, info); +} + template class AC_DATABASE_API DBUpdater; template class AC_DATABASE_API DBUpdater; template class AC_DATABASE_API DBUpdater; diff --git a/src/server/database/Updater/DBUpdater.h b/src/server/database/Updater/DBUpdater.h index 506c549c14..77c1306e50 100644 --- a/src/server/database/Updater/DBUpdater.h +++ b/src/server/database/Updater/DBUpdater.h @@ -19,13 +19,13 @@ #define DBUpdater_h__ #include "DatabaseEnv.h" +#include "DatabaseUpdatePool.h" #include "Define.h" #include "QueryResult.h" #include #include - -template -class DatabaseWorkerPool; +#include +#include namespace boost { @@ -71,6 +71,15 @@ private: static uint32& failed_updates(); }; +// Runtime metadata the updater needs about one database, core or module owned. +struct DBUpdaterInfo +{ + std::string displayName; // name used in log output + std::string sourceDirectory; // root directory holding the sql tree + std::string baseFilesDirectory; // base *.sql files, trailing separator optional + std::string dbModuleName; // update-fetcher module name, must be lowercase +}; + template class AC_DATABASE_API DBUpdater { @@ -79,6 +88,7 @@ public: static inline std::string GetConfigEntry(); static inline std::string GetTableName(); + static std::string GetSourceDirectory(); static std::string GetBaseFilesDirectory(); static bool IsEnabled(uint32 const updateMask); static BaseLocation GetBaseLocationType(); @@ -91,11 +101,18 @@ public: static std::string GetDBModuleName(); private: - static QueryResult Retrieve(DatabaseWorkerPool& pool, std::string const& query); - static void Apply(DatabaseWorkerPool& pool, std::string const& query); - static void ApplyFile(DatabaseWorkerPool& pool, Path const& path); - static void ApplyFile(DatabaseWorkerPool& pool, std::string const& host, std::string const& user, - std::string const& password, std::string const& port_or_socket, std::string const& database, std::string const& ssl, Path const& path); + static DBUpdaterInfo GetUpdaterInfo(); +}; + +// Non-template updater entry points for module-owned pools (see ModuleDatabasePool). +// Mirrors the DBUpdater flow: Create the schema when missing, Populate an empty +// database from the base files, then apply pending updates through the UpdateFetcher. +class AC_DATABASE_API ModuleDBUpdater +{ +public: + static bool Create(DatabaseUpdatePool& pool); + static bool Update(DatabaseUpdatePool& pool, DBUpdaterInfo const& info, std::string_view modulesList = {}); + static bool Populate(DatabaseUpdatePool& pool, DBUpdaterInfo const& info); }; #endif // DBUpdater_h__ diff --git a/src/server/game/Scripting/ScriptDefines/DatabaseScript.cpp b/src/server/game/Scripting/ScriptDefines/DatabaseScript.cpp index ebe4432a64..af01ca420b 100644 --- a/src/server/game/Scripting/ScriptDefines/DatabaseScript.cpp +++ b/src/server/game/Scripting/ScriptDefines/DatabaseScript.cpp @@ -19,6 +19,11 @@ #include "ScriptMgr.h" #include "ScriptMgrMacros.h" +bool ScriptMgr::OnModuleDatabasesLoading() +{ + CALL_ENABLED_BOOLEAN_HOOKS(DatabaseScript, DATABASEHOOK_ON_MODULE_DATABASES_LOADING, !script->OnModuleDatabasesLoading()); +} + void ScriptMgr::OnAfterDatabasesLoaded(uint32 updateFlags) { CALL_ENABLED_HOOKS(DatabaseScript, DATABASEHOOK_ON_AFTER_DATABASES_LOADED, script->OnAfterDatabasesLoaded(updateFlags)); @@ -29,6 +34,26 @@ void ScriptMgr::OnAfterDatabaseLoadCreatureTemplates(std::vectorOnAfterDatabaseLoadCreatureTemplates(creatureTemplates)); } +void ScriptMgr::OnModuleDatabasesKeepAlive() +{ + CALL_ENABLED_HOOKS(DatabaseScript, DATABASEHOOK_ON_MODULE_DATABASES_KEEPALIVE, script->OnModuleDatabasesKeepAlive()); +} + +void ScriptMgr::OnModuleDatabasesClosing() +{ + CALL_ENABLED_HOOKS(DatabaseScript, DATABASEHOOK_ON_MODULE_DATABASES_CLOSING, script->OnModuleDatabasesClosing()); +} + +void ScriptMgr::OnDatabaseWarnAboutSyncQueries(bool apply) +{ + CALL_ENABLED_HOOKS(DatabaseScript, DATABASEHOOK_ON_DATABASE_WARN_ABOUT_SYNC_QUERIES, script->OnDatabaseWarnAboutSyncQueries(apply)); +} + +void ScriptMgr::OnDatabaseGetDBRevision(std::map& revisions) +{ + CALL_ENABLED_HOOKS(DatabaseScript, DATABASEHOOK_ON_DATABASE_GET_DB_REVISION, script->OnDatabaseGetDBRevision(revisions)); +} + DatabaseScript::DatabaseScript(char const* name, std::vector enabledHooks) : ScriptObject(name, DATABASEHOOK_END) { diff --git a/src/server/game/Scripting/ScriptDefines/DatabaseScript.h b/src/server/game/Scripting/ScriptDefines/DatabaseScript.h index c4d3408a2f..ef5c7dc64d 100644 --- a/src/server/game/Scripting/ScriptDefines/DatabaseScript.h +++ b/src/server/game/Scripting/ScriptDefines/DatabaseScript.h @@ -19,12 +19,19 @@ #define SCRIPT_OBJECT_DATABASE_SCRIPT_H_ #include "ScriptObject.h" +#include +#include #include enum DatabaseHook { DATABASEHOOK_ON_AFTER_DATABASES_LOADED, DATABASEHOOK_ON_AFTER_DATABASE_LOAD_CREATURETEMPLATES, + DATABASEHOOK_ON_MODULE_DATABASES_LOADING, + DATABASEHOOK_ON_MODULE_DATABASES_KEEPALIVE, + DATABASEHOOK_ON_MODULE_DATABASES_CLOSING, + DATABASEHOOK_ON_DATABASE_WARN_ABOUT_SYNC_QUERIES, + DATABASEHOOK_ON_DATABASE_GET_DB_REVISION, DATABASEHOOK_END }; @@ -52,6 +59,38 @@ public: */ virtual void OnAfterDatabaseLoadCreatureTemplates(std::vector /*creatureTemplates*/) { } + /** + * @brief Called once the core databases are up, so a module can open a database of its own. + * Runs before the rest of the world loads, unlike OnAfterDatabasesLoaded which reports the + * finished core load. + * + * @return false to abort startup, e.g. when the module's own database failed to open + */ + [[nodiscard]] virtual bool OnModuleDatabasesLoading() { return true; } + + /** + * @brief Called on the world's keep-alive tick, alongside the core pools being pinged. + */ + virtual void OnModuleDatabasesKeepAlive() { } + + /** + * @brief Called after the core databases are closed, so a module can close its own. + */ + virtual void OnModuleDatabasesClosing() { } + + /** + * @brief Called when the core turns its synchronous-query warning on or off. + * + * @param apply True when the warning is being enabled + */ + virtual void OnDatabaseWarnAboutSyncQueries(bool /*apply*/) { } + + /** + * @brief Called by .server info to collect the revision of a module-owned database. + * + * @param revisions Revision string to report, keyed by module name + */ + virtual void OnDatabaseGetDBRevision(std::map& /*revisions*/) { } }; #endif diff --git a/src/server/game/Scripting/ScriptMgr.h b/src/server/game/Scripting/ScriptMgr.h index 4bedab5f8f..f4187a3a85 100644 --- a/src/server/game/Scripting/ScriptMgr.h +++ b/src/server/game/Scripting/ScriptMgr.h @@ -712,8 +712,13 @@ public: /* CommandSC */ public: /* DatabaseScript */ + bool OnModuleDatabasesLoading(); void OnAfterDatabasesLoaded(uint32 updateFlags); void OnAfterDatabaseLoadCreatureTemplates(std::vector creatureTemplateStore); + void OnModuleDatabasesKeepAlive(); + void OnModuleDatabasesClosing(); + void OnDatabaseWarnAboutSyncQueries(bool apply); + void OnDatabaseGetDBRevision(std::map& revisions); public: /* WorldObjectScript */ diff --git a/src/server/game/World/World.cpp b/src/server/game/World/World.cpp index a7a44942e5..fc5a74803e 100644 --- a/src/server/game/World/World.cpp +++ b/src/server/game/World/World.cpp @@ -1322,6 +1322,7 @@ void World::Update(uint32 diff) CharacterDatabase.KeepAlive(); LoginDatabase.KeepAlive(); WorldDatabase.KeepAlive(); + sScriptMgr->OnModuleDatabasesKeepAlive(); } { diff --git a/src/server/scripts/Commands/cs_server.cpp b/src/server/scripts/Commands/cs_server.cpp index 05370ba028..fe8ccf8d6c 100644 --- a/src/server/scripts/Commands/cs_server.cpp +++ b/src/server/scripts/Commands/cs_server.cpp @@ -27,6 +27,7 @@ #include "MySQLThreading.h" #include "RBAC.h" #include "Realm.h" +#include "ScriptMgr.h" #include "StringConvert.h" #include "UpdateTime.h" #include "VMapFactory.h" @@ -34,6 +35,7 @@ #include "WorldSessionMgr.h" #include #include +#include #include #include #include @@ -216,6 +218,11 @@ public: handler->PSendSysMessage("Using World DB: {}", sWorld->GetDBVersion()); + std::map moduleDBRevisions; + sScriptMgr->OnDatabaseGetDBRevision(moduleDBRevisions); + for (auto const& [moduleName, revision] : moduleDBRevisions) + handler->PSendSysMessage("Using {} DB Revision: {}", moduleName, revision); + std::string lldb = "No updates found!"; if (QueryResult resL = LoginDatabase.Query("SELECT name FROM updates ORDER BY name DESC LIMIT 1")) {