diff --git a/src/common/Threading/PCQueue.h b/src/common/Threading/PCQueue.h index f66adb235e..e65a1d78fc 100644 --- a/src/common/Threading/PCQueue.h +++ b/src/common/Threading/PCQueue.h @@ -102,6 +102,15 @@ public: _condition.notify_all(); } + // Reopens the queue after Cancel()/Shutdown() so new consumers can be attached. + // Callers must make sure the previous consumers have already stopped. + void Reset() + { + std::lock_guard lock(_queueLock); + _cancel = false; + _shutdown = false; + } + private: template typename std::enable_if::value>::type DeleteQueuedObject(E& obj) diff --git a/src/server/database/Database/DatabaseWorkerPool.cpp b/src/server/database/Database/DatabaseWorkerPool.cpp index ee19a5a48c..c57a65b97e 100644 --- a/src/server/database/Database/DatabaseWorkerPool.cpp +++ b/src/server/database/Database/DatabaseWorkerPool.cpp @@ -91,6 +91,11 @@ uint32 DatabaseWorkerPool::Open() LOG_INFO("sql.driver", "Opening DatabasePool '{}'. Asynchronous connections: {}, synchronous connections: {}.", GetDatabaseName(), _async_threads, _synch_threads); + _queue->Cancel(); + _connections[IDX_ASYNC].clear(); + _connections[IDX_SYNCH].clear(); + _queue->Reset(); + uint32 error = OpenConnections(IDX_ASYNC, _async_threads); if (error)