From 6a23255c58e816b91174afa37ae2f6e404f9c0e2 Mon Sep 17 00:00:00 2001 From: xdustinface Date: Fri, 10 Jul 2020 01:42:43 +0200 Subject: [PATCH 1/5] llqm: Fix thread handling in CDKGSessionManager and CDKGSessionHandler --- src/llmq/quorums_dkgsessionhandler.cpp | 37 ++++++++++++++++++++------ src/llmq/quorums_dkgsessionhandler.h | 3 +++ src/llmq/quorums_dkgsessionmgr.cpp | 19 ++++++++----- src/llmq/quorums_dkgsessionmgr.h | 4 +-- src/llmq/quorums_init.cpp | 4 +-- 5 files changed, 49 insertions(+), 18 deletions(-) diff --git a/src/llmq/quorums_dkgsessionhandler.cpp b/src/llmq/quorums_dkgsessionhandler.cpp index 2907f8f722cc..b4178d1b7de4 100644 --- a/src/llmq/quorums_dkgsessionhandler.cpp +++ b/src/llmq/quorums_dkgsessionhandler.cpp @@ -95,18 +95,10 @@ CDKGSessionHandler::CDKGSessionHandler(const Consensus::LLMQParams& _params, ctp pendingJustifications((size_t)_params.size * 2, MSG_QUORUM_JUSTIFICATION), pendingPrematureCommitments((size_t)_params.size * 2, MSG_QUORUM_PREMATURE_COMMITMENT) { - phaseHandlerThread = std::thread([this] { - RenameThread(strprintf("dash-q-phase-%d", (uint8_t)params.type).c_str()); - PhaseHandlerThread(); - }); } CDKGSessionHandler::~CDKGSessionHandler() { - stopRequested = true; - if (phaseHandlerThread.joinable()) { - phaseHandlerThread.join(); - } } void CDKGSessionHandler::UpdatedBlockTip(const CBlockIndex* pindexNew) @@ -145,6 +137,35 @@ void CDKGSessionHandler::ProcessMessage(CNode* pfrom, const std::string& strComm } } +void CDKGSessionHandler::StartThread() +{ + auto threadName = [&]() -> auto { + switch (params.type) { + case Consensus::LLMQ_50_60: + return "q-phase-1"; + case Consensus::LLMQ_400_60: + return "q-phase-2"; + case Consensus::LLMQ_400_85: + return "q-phase-3"; + case Consensus::LLMQ_TEST: + return "q-phase-100"; + case Consensus::LLMQ_DEVNET: + return "q-phase-101"; + default: + throw std::runtime_error("Tried to start a CDKGSessionHandler thread for LLMQ_NONE."); + } + }; + phaseHandlerThread = std::thread(&TraceThread >, threadName(), std::function(std::bind(&CDKGSessionHandler::PhaseHandlerThread, this))); +} + +void CDKGSessionHandler::StopThread() +{ + stopRequested = true; + if (phaseHandlerThread.joinable()) { + phaseHandlerThread.join(); + } +} + bool CDKGSessionHandler::InitNewQuorum(const CBlockIndex* pindexQuorum) { //AssertLockHeld(cs_main); diff --git a/src/llmq/quorums_dkgsessionhandler.h b/src/llmq/quorums_dkgsessionhandler.h index 7ee399973bd8..e8c18ffe70a7 100644 --- a/src/llmq/quorums_dkgsessionhandler.h +++ b/src/llmq/quorums_dkgsessionhandler.h @@ -126,6 +126,9 @@ class CDKGSessionHandler void UpdatedBlockTip(const CBlockIndex *pindexNew); void ProcessMessage(CNode* pfrom, const std::string& strCommand, CDataStream& vRecv, CConnman& connman); + void StartThread(); + void StopThread(); + private: bool InitNewQuorum(const CBlockIndex* pindexQuorum); diff --git a/src/llmq/quorums_dkgsessionmgr.cpp b/src/llmq/quorums_dkgsessionmgr.cpp index 820b1794266c..2621ac02b5f5 100644 --- a/src/llmq/quorums_dkgsessionmgr.cpp +++ b/src/llmq/quorums_dkgsessionmgr.cpp @@ -24,26 +24,33 @@ CDKGSessionManager::CDKGSessionManager(CDBWrapper& _llmqDb, CBLSWorker& _blsWork llmqDb(_llmqDb), blsWorker(_blsWorker) { + for (const auto& qt : Params().GetConsensus().llmqs) { + dkgSessionHandlers.emplace(std::piecewise_construct, + std::forward_as_tuple(qt.first), + std::forward_as_tuple(qt.second, messageHandlerPool, blsWorker, *this)); + } } CDKGSessionManager::~CDKGSessionManager() { } -void CDKGSessionManager::StartMessageHandlerPool() +void CDKGSessionManager::StartThreads() { - for (const auto& qt : Params().GetConsensus().llmqs) { - dkgSessionHandlers.emplace(std::piecewise_construct, - std::forward_as_tuple(qt.first), - std::forward_as_tuple(qt.second, messageHandlerPool, blsWorker, *this)); + for (auto& it : dkgSessionHandlers) { + it.second.StartThread(); } messageHandlerPool.resize(2); RenameThreadPool(messageHandlerPool, "dash-q-msg"); } -void CDKGSessionManager::StopMessageHandlerPool() +void CDKGSessionManager::StopThreads() { + for (auto& it : dkgSessionHandlers) { + it.second.StopThread(); + } + messageHandlerPool.stop(true); } diff --git a/src/llmq/quorums_dkgsessionmgr.h b/src/llmq/quorums_dkgsessionmgr.h index ca13475daf17..0661db5005cf 100644 --- a/src/llmq/quorums_dkgsessionmgr.h +++ b/src/llmq/quorums_dkgsessionmgr.h @@ -50,8 +50,8 @@ class CDKGSessionManager CDKGSessionManager(CDBWrapper& _llmqDb, CBLSWorker& _blsWorker); ~CDKGSessionManager(); - void StartMessageHandlerPool(); - void StopMessageHandlerPool(); + void StartThreads(); + void StopThreads(); void UpdatedBlockTip(const CBlockIndex *pindexNew, bool fInitialDownload); diff --git a/src/llmq/quorums_init.cpp b/src/llmq/quorums_init.cpp index 23b8b5ebdc9d..85b412b8573f 100644 --- a/src/llmq/quorums_init.cpp +++ b/src/llmq/quorums_init.cpp @@ -70,7 +70,7 @@ void StartLLMQSystem() blsWorker->Start(); } if (quorumDKGSessionManager) { - quorumDKGSessionManager->StartMessageHandlerPool(); + quorumDKGSessionManager->StartThreads(); } if (quorumSigSharesManager) { quorumSigSharesManager->RegisterAsRecoveredSigsListener(); @@ -97,7 +97,7 @@ void StopLLMQSystem() quorumSigSharesManager->UnregisterAsRecoveredSigsListener(); } if (quorumDKGSessionManager) { - quorumDKGSessionManager->StopMessageHandlerPool(); + quorumDKGSessionManager->StopThreads(); } if (blsWorker) { blsWorker->Stop(); From 6fedfb71760e27d5ca3fd67f5daef043fcb8f103 Mon Sep 17 00:00:00 2001 From: xdustinface Date: Fri, 10 Jul 2020 01:45:04 +0200 Subject: [PATCH 2/5] llmq: Removed unused thread_pool from CDKGSessionManager --- src/llmq/quorums_dkgsessionhandler.cpp | 3 +-- src/llmq/quorums_dkgsessionhandler.h | 3 +-- src/llmq/quorums_dkgsessionmgr.cpp | 7 +------ src/llmq/quorums_dkgsessionmgr.h | 1 - 4 files changed, 3 insertions(+), 11 deletions(-) diff --git a/src/llmq/quorums_dkgsessionhandler.cpp b/src/llmq/quorums_dkgsessionhandler.cpp index b4178d1b7de4..f6b214728e3f 100644 --- a/src/llmq/quorums_dkgsessionhandler.cpp +++ b/src/llmq/quorums_dkgsessionhandler.cpp @@ -84,9 +84,8 @@ void CDKGPendingMessages::Clear() ////// -CDKGSessionHandler::CDKGSessionHandler(const Consensus::LLMQParams& _params, ctpl::thread_pool& _messageHandlerPool, CBLSWorker& _blsWorker, CDKGSessionManager& _dkgManager) : +CDKGSessionHandler::CDKGSessionHandler(const Consensus::LLMQParams& _params, CBLSWorker& _blsWorker, CDKGSessionManager& _dkgManager) : params(_params), - messageHandlerPool(_messageHandlerPool), blsWorker(_blsWorker), dkgManager(_dkgManager), curSession(std::make_shared(_params, _blsWorker, _dkgManager)), diff --git a/src/llmq/quorums_dkgsessionhandler.h b/src/llmq/quorums_dkgsessionhandler.h index e8c18ffe70a7..1b7430ae2a42 100644 --- a/src/llmq/quorums_dkgsessionhandler.h +++ b/src/llmq/quorums_dkgsessionhandler.h @@ -103,7 +103,6 @@ class CDKGSessionHandler std::atomic stopRequested{false}; const Consensus::LLMQParams& params; - ctpl::thread_pool& messageHandlerPool; CBLSWorker& blsWorker; CDKGSessionManager& dkgManager; @@ -120,7 +119,7 @@ class CDKGSessionHandler CDKGPendingMessages pendingPrematureCommitments; public: - CDKGSessionHandler(const Consensus::LLMQParams& _params, ctpl::thread_pool& _messageHandlerPool, CBLSWorker& blsWorker, CDKGSessionManager& _dkgManager); + CDKGSessionHandler(const Consensus::LLMQParams& _params, CBLSWorker& blsWorker, CDKGSessionManager& _dkgManager); ~CDKGSessionHandler(); void UpdatedBlockTip(const CBlockIndex *pindexNew); diff --git a/src/llmq/quorums_dkgsessionmgr.cpp b/src/llmq/quorums_dkgsessionmgr.cpp index 2621ac02b5f5..b4d55c143862 100644 --- a/src/llmq/quorums_dkgsessionmgr.cpp +++ b/src/llmq/quorums_dkgsessionmgr.cpp @@ -27,7 +27,7 @@ CDKGSessionManager::CDKGSessionManager(CDBWrapper& _llmqDb, CBLSWorker& _blsWork for (const auto& qt : Params().GetConsensus().llmqs) { dkgSessionHandlers.emplace(std::piecewise_construct, std::forward_as_tuple(qt.first), - std::forward_as_tuple(qt.second, messageHandlerPool, blsWorker, *this)); + std::forward_as_tuple(qt.second, blsWorker, *this)); } } @@ -40,9 +40,6 @@ void CDKGSessionManager::StartThreads() for (auto& it : dkgSessionHandlers) { it.second.StartThread(); } - - messageHandlerPool.resize(2); - RenameThreadPool(messageHandlerPool, "dash-q-msg"); } void CDKGSessionManager::StopThreads() @@ -50,8 +47,6 @@ void CDKGSessionManager::StopThreads() for (auto& it : dkgSessionHandlers) { it.second.StopThread(); } - - messageHandlerPool.stop(true); } void CDKGSessionManager::UpdatedBlockTip(const CBlockIndex* pindexNew, bool fInitialDownload) diff --git a/src/llmq/quorums_dkgsessionmgr.h b/src/llmq/quorums_dkgsessionmgr.h index 0661db5005cf..4e0d20d2f1ea 100644 --- a/src/llmq/quorums_dkgsessionmgr.h +++ b/src/llmq/quorums_dkgsessionmgr.h @@ -23,7 +23,6 @@ class CDKGSessionManager private: CDBWrapper& llmqDb; CBLSWorker& blsWorker; - ctpl::thread_pool messageHandlerPool; std::map dkgSessionHandlers; From 106177577c8cd52a6fa81815758f46553161dbad Mon Sep 17 00:00:00 2001 From: UdjinM6 Date: Mon, 13 Jul 2020 23:47:26 +0300 Subject: [PATCH 3/5] Tweak `CDKGSessionHandler::StartThread()` --- src/llmq/quorums_dkgsessionhandler.cpp | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/src/llmq/quorums_dkgsessionhandler.cpp b/src/llmq/quorums_dkgsessionhandler.cpp index f6b214728e3f..b41c15e09896 100644 --- a/src/llmq/quorums_dkgsessionhandler.cpp +++ b/src/llmq/quorums_dkgsessionhandler.cpp @@ -138,6 +138,10 @@ void CDKGSessionHandler::ProcessMessage(CNode* pfrom, const std::string& strComm void CDKGSessionHandler::StartThread() { + if (phaseHandlerThread.joinable()) { + throw std::runtime_error("Tried to start an already started CDKGSessionHandler thread."); + } + auto threadName = [&]() -> auto { switch (params.type) { case Consensus::LLMQ_50_60: @@ -150,8 +154,10 @@ void CDKGSessionHandler::StartThread() return "q-phase-100"; case Consensus::LLMQ_DEVNET: return "q-phase-101"; - default: + case Consensus::LLMQ_NONE: throw std::runtime_error("Tried to start a CDKGSessionHandler thread for LLMQ_NONE."); + default: + throw std::runtime_error("Tried to start a CDKGSessionHandler thread for an unknown LLMQ type."); } }; phaseHandlerThread = std::thread(&TraceThread >, threadName(), std::function(std::bind(&CDKGSessionHandler::PhaseHandlerThread, this))); From 01033067dd75509418c61a4e3b846373c38dfabf Mon Sep 17 00:00:00 2001 From: xdustinface Date: Tue, 14 Jul 2020 13:07:05 +0200 Subject: [PATCH 4/5] llmq: Simplify CDKGSessionHandler's thread naming --- src/llmq/quorums_dkgsessionhandler.cpp | 21 ++------------------- 1 file changed, 2 insertions(+), 19 deletions(-) diff --git a/src/llmq/quorums_dkgsessionhandler.cpp b/src/llmq/quorums_dkgsessionhandler.cpp index b41c15e09896..ede13f105438 100644 --- a/src/llmq/quorums_dkgsessionhandler.cpp +++ b/src/llmq/quorums_dkgsessionhandler.cpp @@ -142,25 +142,8 @@ void CDKGSessionHandler::StartThread() throw std::runtime_error("Tried to start an already started CDKGSessionHandler thread."); } - auto threadName = [&]() -> auto { - switch (params.type) { - case Consensus::LLMQ_50_60: - return "q-phase-1"; - case Consensus::LLMQ_400_60: - return "q-phase-2"; - case Consensus::LLMQ_400_85: - return "q-phase-3"; - case Consensus::LLMQ_TEST: - return "q-phase-100"; - case Consensus::LLMQ_DEVNET: - return "q-phase-101"; - case Consensus::LLMQ_NONE: - throw std::runtime_error("Tried to start a CDKGSessionHandler thread for LLMQ_NONE."); - default: - throw std::runtime_error("Tried to start a CDKGSessionHandler thread for an unknown LLMQ type."); - } - }; - phaseHandlerThread = std::thread(&TraceThread >, threadName(), std::function(std::bind(&CDKGSessionHandler::PhaseHandlerThread, this))); + std::string threadName = strprintf("q-phase-%d", params.type); + phaseHandlerThread = std::thread(&TraceThread >, threadName, std::function(std::bind(&CDKGSessionHandler::PhaseHandlerThread, this))); } void CDKGSessionHandler::StopThread() From ef66d06620d5dd003dc0ab9384f59ed0f26c2e01 Mon Sep 17 00:00:00 2001 From: xdustinface Date: Tue, 14 Jul 2020 16:31:57 +0200 Subject: [PATCH 5/5] llmq: Make sure CDKGSessionHandler uses a valid LLMQ type Co-Authored-By: UdjinM6 --- src/llmq/quorums_dkgsessionhandler.cpp | 3 +++ 1 file changed, 3 insertions(+) diff --git a/src/llmq/quorums_dkgsessionhandler.cpp b/src/llmq/quorums_dkgsessionhandler.cpp index ede13f105438..42ad891f0e46 100644 --- a/src/llmq/quorums_dkgsessionhandler.cpp +++ b/src/llmq/quorums_dkgsessionhandler.cpp @@ -94,6 +94,9 @@ CDKGSessionHandler::CDKGSessionHandler(const Consensus::LLMQParams& _params, CBL pendingJustifications((size_t)_params.size * 2, MSG_QUORUM_JUSTIFICATION), pendingPrematureCommitments((size_t)_params.size * 2, MSG_QUORUM_PREMATURE_COMMITMENT) { + if (params.type == Consensus::LLMQ_NONE) { + throw std::runtime_error("Can't initialize CDKGSessionHandler with LLMQ_NONE type."); + } } CDKGSessionHandler::~CDKGSessionHandler()