Skip to content

Commit 8637d68

Browse files
fix: retain initialized DKG sessions during inventory reads
1 parent f35a24c commit 8637d68

4 files changed

Lines changed: 115 additions & 20 deletions

File tree

‎src/active/dkgsessionhandler.cpp‎

Lines changed: 27 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -87,16 +87,23 @@ bool ActiveDKGSessionHandler::InitNewQuorum(gsl::not_null<const CBlockIndex*> pQ
8787
return false;
8888
}
8989

90-
curSession = std::make_unique<ActiveDKGSession>(m_bls_worker, m_dmnman, m_dkgdbgman, m_qdkgsman, m_mn_metaman,
91-
m_qsnapman, m_mn_activeman, m_chainman, m_sporkman,
92-
pQuorumBaseBlockIndex, params);
90+
auto next_session = std::make_shared<ActiveDKGSession>(m_bls_worker, m_dmnman, m_dkgdbgman, m_qdkgsman,
91+
m_mn_metaman, m_qsnapman, m_mn_activeman, m_chainman,
92+
m_sporkman, pQuorumBaseBlockIndex, params);
9393

94-
if (!curSession->Init(m_mn_activeman.GetProTxHash(), quorumIndex)) {
94+
{
95+
std::shared_ptr<CDKGSession> previous_session;
96+
WITH_LOCK(cs_session, previous_session.swap(curSession));
97+
}
98+
99+
if (!next_session->Init(m_mn_activeman.GetProTxHash(), quorumIndex)) {
95100
LogPrintf("ActiveDKGSessionHandler::%s -- height[%d] quorum initialization failed for %s qi[%d]\n", __func__,
96101
pQuorumBaseBlockIndex->nHeight, params.name, quorumIndex);
97102
return false;
98103
}
99104

105+
WITH_LOCK(cs_session, curSession = std::move(next_session));
106+
100107
LogPrintf("ActiveDKGSessionHandler::%s -- height[%d] quorum initialization OK for %s qi[%d]\n", __func__, pQuorumBaseBlockIndex->nHeight, params.name, quorumIndex);
101108
return true;
102109
}
@@ -166,7 +173,8 @@ void ActiveDKGSessionHandler::WaitForNewQuorum(const uint256& oldQuorumHash) con
166173
void ActiveDKGSessionHandler::SleepBeforePhase(QuorumPhase curPhase, const uint256& expectedQuorumHash,
167174
double randomSleepFactor, const WhileWaitFunc& runWhileWaiting) const
168175
{
169-
if (!curSession->AreWeMember()) {
176+
const auto session = GetCurSession();
177+
if (!session || !session->AreWeMember()) {
170178
// Non-members do not participate and do not create any network load, no need to sleep.
171179
return;
172180
}
@@ -187,7 +195,7 @@ void ActiveDKGSessionHandler::SleepBeforePhase(QuorumPhase curPhase, const uint2
187195
// Don't expect perfect block times and thus reduce the phase time to be on the secure side (caller chooses factor)
188196
double adjustedPhaseSleepTimePerMember = phaseSleepTimePerMember * randomSleepFactor;
189197

190-
int64_t sleepTime = static_cast<int64_t>(adjustedPhaseSleepTimePerMember * curSession->GetMyMemberIndex().value_or(0));
198+
int64_t sleepTime = static_cast<int64_t>(adjustedPhaseSleepTimePerMember * session->GetMyMemberIndex().value_or(0));
191199
const auto endTime = SteadyClock::now() + std::chrono::milliseconds{sleepTime};
192200
int heightTmp{currentHeight.load()};
193201
int heightStart{heightTmp};
@@ -248,22 +256,31 @@ void ActiveDKGSessionHandler::HandlePhase(QuorumPhase curPhase, QuorumPhase next
248256

249257
bool ActiveDKGSessionHandler::GetContribution(const uint256& hash, CDKGContribution& ret) const
250258
{
251-
return curSession && curSession->GetContribution(hash, ret);
259+
const auto session = GetCurSession();
260+
return session && session->GetContribution(hash, ret);
252261
}
253262

254263
bool ActiveDKGSessionHandler::GetComplaint(const uint256& hash, CDKGComplaint& ret) const
255264
{
256-
return curSession && curSession->GetComplaint(hash, ret);
265+
const auto session = GetCurSession();
266+
return session && session->GetComplaint(hash, ret);
257267
}
258268

259269
bool ActiveDKGSessionHandler::GetJustification(const uint256& hash, CDKGJustification& ret) const
260270
{
261-
return curSession && curSession->GetJustification(hash, ret);
271+
const auto session = GetCurSession();
272+
return session && session->GetJustification(hash, ret);
262273
}
263274

264275
bool ActiveDKGSessionHandler::GetPrematureCommitment(const uint256& hash, CDKGPrematureCommitment& ret) const
265276
{
266-
return curSession && curSession->GetPrematureCommitment(hash, ret);
277+
const auto session = GetCurSession();
278+
return session && session->GetPrematureCommitment(hash, ret);
279+
}
280+
281+
std::shared_ptr<CDKGSession> ActiveDKGSessionHandler::GetCurSession() const
282+
{
283+
return WITH_LOCK(cs_session, return curSession);
267284
}
268285

269286
QuorumPhase ActiveDKGSessionHandler::GetPhase() const

‎src/active/dkgsessionhandler.h‎

Lines changed: 13 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -61,7 +61,8 @@ class ActiveDKGSessionHandler final : public llmq::CDKGSessionHandler
6161
private:
6262
std::atomic<bool> stopRequested{false};
6363
std::atomic<int> currentHeight{-1};
64-
std::unique_ptr<CDKGSession> curSession{nullptr};
64+
mutable Mutex cs_session;
65+
std::shared_ptr<CDKGSession> curSession GUARDED_BY(cs_session);
6566

6667
mutable Mutex cs_phase_qhash;
6768
QuorumPhase phase GUARDED_BY(cs_phase_qhash){QuorumPhase::Idle};
@@ -81,10 +82,12 @@ class ActiveDKGSessionHandler final : public llmq::CDKGSessionHandler
8182

8283
public:
8384
//! CDKGSessionHandler
84-
bool GetContribution(const uint256& hash, CDKGContribution& ret) const override;
85-
bool GetComplaint(const uint256& hash, CDKGComplaint& ret) const override;
86-
bool GetJustification(const uint256& hash, CDKGJustification& ret) const override;
87-
bool GetPrematureCommitment(const uint256& hash, CDKGPrematureCommitment& ret) const override;
85+
bool GetContribution(const uint256& hash, CDKGContribution& ret) const override EXCLUSIVE_LOCKS_REQUIRED(!cs_session);
86+
bool GetComplaint(const uint256& hash, CDKGComplaint& ret) const override EXCLUSIVE_LOCKS_REQUIRED(!cs_session);
87+
bool GetJustification(const uint256& hash, CDKGJustification& ret) const override
88+
EXCLUSIVE_LOCKS_REQUIRED(!cs_session);
89+
bool GetPrematureCommitment(const uint256& hash, CDKGPrematureCommitment& ret) const override
90+
EXCLUSIVE_LOCKS_REQUIRED(!cs_session);
8891
QuorumPhase GetPhase() const override EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash);
8992
void UpdatedBlockTip(const CBlockIndex* pindexNew) override EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash);
9093

@@ -96,12 +99,12 @@ class ActiveDKGSessionHandler final : public llmq::CDKGSessionHandler
9699
int QuorumIndex() const { return quorumIndex; }
97100
bool QuorumsWatch() const { return m_quorums_watch; }
98101
uint256 GetCurrentQuorumHash() const EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash);
99-
CDKGSession* GetCurSession() { return curSession.get(); }
102+
std::shared_ptr<CDKGSession> GetCurSession() const EXCLUSIVE_LOCKS_REQUIRED(!cs_session);
100103

101104
void RequestStop() { stopRequested = true; }
102105
bool IsStopRequested() const { return stopRequested; }
103106

104-
bool InitNewQuorum(gsl::not_null<const CBlockIndex*> pQuorumBaseBlockIndex);
107+
bool InitNewQuorum(gsl::not_null<const CBlockIndex*> pQuorumBaseBlockIndex) EXCLUSIVE_LOCKS_REQUIRED(!cs_session);
105108

106109
/**
107110
* @param curPhase current QuorumPhase
@@ -114,10 +117,11 @@ class ActiveDKGSessionHandler final : public llmq::CDKGSessionHandler
114117
const WhileWaitFunc& shouldNotWait = [] { return false; }) const EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash);
115118
void WaitForNewQuorum(const uint256& oldQuorumHash) const EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash);
116119
void SleepBeforePhase(QuorumPhase curPhase, const uint256& expectedQuorumHash, double randomSleepFactor,
117-
const WhileWaitFunc& runWhileWaiting) const EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash);
120+
const WhileWaitFunc& runWhileWaiting) const
121+
EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash, !cs_session);
118122
void HandlePhase(QuorumPhase curPhase, QuorumPhase nextPhase, const uint256& expectedQuorumHash,
119123
double randomSleepFactor, const StartPhaseFunc& startPhaseFunc,
120-
const WhileWaitFunc& runWhileWaiting) EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash);
124+
const WhileWaitFunc& runWhileWaiting) EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash, !cs_session);
121125

122126
private:
123127
std::pair<QuorumPhase, uint256> GetPhaseAndQuorumHash() const EXCLUSIVE_LOCKS_REQUIRED(!cs_phase_qhash);

‎src/llmq/net_dkg.cpp‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -851,7 +851,7 @@ void NetDKG::HandleDKGRound(ActiveDKGSessionHandler& handler)
851851

852852
active.dkgdbgman.MarkPhaseAdvanced(handler.params.type, handler.QuorumIndex(), QuorumPhase::Initialized);
853853

854-
auto* curSession = handler.GetCurSession();
854+
const auto curSession = handler.GetCurSession();
855855
if (handler.params.is_single_member()) {
856856
auto finalCommitment = curSession->FinalizeSingleCommitment();
857857
if (!finalCommitment.IsNull()) { // it can be null only if we are not member

‎src/test/llmq_dkg_tests.cpp‎

Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,12 +2,20 @@
22
// Distributed under the MIT software license, see the accompanying
33
// file COPYING or http://www.opensource.org/licenses/mit-license.php.
44

5+
#include <active/dkgsessionhandler.h>
6+
#include <active/masternode.h>
7+
#include <dbwrapper.h>
8+
#include <llmq/context.h>
9+
#include <llmq/debug.h>
510
#include <llmq/dkgsession.h>
611
#include <llmq/dkgsessionhandler.h>
12+
#include <llmq/dkgsessionmgr.h>
713
#include <protocol.h>
814
#include <streams.h>
15+
#include <test/util/setup_common.h>
916
#include <util/helpers.h>
1017
#include <util/std23.h>
18+
#include <validation.h>
1119

1220
#include <boost/test/unit_test.hpp>
1321

@@ -120,4 +128,70 @@ BOOST_AUTO_TEST_CASE(pending_messages_quota_is_per_protx)
120128
BOOST_CHECK_EQUAL(pending.PopPendingMessages(6).size(), 5U);
121129
}
122130

131+
namespace {
132+
struct DKGLifetimeSetup : RegTestingSetup {
133+
DKGLifetimeSetup() :
134+
RegTestingSetup({"-dip3params=0:0"})
135+
{
136+
}
137+
};
138+
} // namespace
139+
140+
BOOST_FIXTURE_TEST_CASE(inventory_session_survives_replacement, DKGLifetimeSetup)
141+
{
142+
auto params = Params().GetLLMQ(Consensus::LLMQType::LLMQ_TEST).value();
143+
// Isolate session lifetime from masternode registration and DKG participation.
144+
params.minSize = 0;
145+
CBLSSecretKey operator_key;
146+
operator_key.MakeNewKey();
147+
auto& ctx = *m_node.llmq_ctx;
148+
CActiveMasternodeManager active{*m_node.connman, *m_node.dmnman, operator_key};
149+
llmq::CDKGDebugManager debug{*m_node.dmnman, *ctx.qsnapman, *m_node.chainman};
150+
llmq::CDKGSessionManager manager{*m_node.dmnman,
151+
*ctx.qsnapman,
152+
*m_node.chainman,
153+
*m_node.sporkman,
154+
{.path = m_args.GetDataDirBase() / "dkg_sessions", .memory = true, .wipe = true}};
155+
llmq::ActiveDKGSessionHandler handler{*ctx.bls_worker,
156+
*m_node.dmnman,
157+
*m_node.mn_metaman,
158+
debug,
159+
manager,
160+
*ctx.quorum_block_processor,
161+
*ctx.qsnapman,
162+
active,
163+
*m_node.chainman,
164+
*m_node.sporkman,
165+
params,
166+
false,
167+
0};
168+
const auto* base = WITH_LOCK(cs_main, return m_node.chainman->ActiveTip());
169+
BOOST_CHECK(!handler.GetCurSession());
170+
BOOST_REQUIRE(handler.InitNewQuorum(base));
171+
auto snapshot = handler.GetCurSession();
172+
BOOST_REQUIRE(snapshot);
173+
BOOST_REQUIRE(handler.InitNewQuorum(base));
174+
BOOST_CHECK(snapshot != handler.GetCurSession());
175+
BOOST_CHECK(snapshot->BlockIndex() == base);
176+
llmq::CDKGContribution contribution;
177+
llmq::CDKGComplaint complaint;
178+
llmq::CDKGJustification justification;
179+
llmq::CDKGPrematureCommitment commitment;
180+
BOOST_CHECK(!snapshot->GetContribution(uint256::ONE, contribution));
181+
BOOST_CHECK(!snapshot->GetComplaint(uint256::ONE, complaint));
182+
BOOST_CHECK(!snapshot->GetJustification(uint256::ONE, justification));
183+
BOOST_CHECK(!snapshot->GetPrematureCommitment(uint256::ONE, commitment));
184+
BOOST_CHECK(!handler.GetContribution(uint256::ONE, contribution));
185+
BOOST_CHECK(!handler.GetComplaint(uint256::ONE, complaint));
186+
BOOST_CHECK(!handler.GetJustification(uint256::ONE, justification));
187+
BOOST_CHECK(!handler.GetPrematureCommitment(uint256::ONE, commitment));
188+
189+
// Failed initialization retires the current session instead of serving prior-round data.
190+
params.minSize = 1;
191+
BOOST_CHECK(!handler.InitNewQuorum(base));
192+
BOOST_CHECK(!handler.GetCurSession());
193+
BOOST_CHECK(!handler.GetContribution(uint256::ONE, contribution));
194+
BOOST_CHECK(snapshot->BlockIndex() == base);
195+
}
196+
123197
BOOST_AUTO_TEST_SUITE_END()

0 commit comments

Comments
 (0)