Skip to content

Commit 85076ea

Browse files
authored
Fix[mqb]: FileStore must flush storage before control messages (#1202)
Signed-off-by: dorjesinpo <129227380+dorjesinpo@users.noreply.github.com>
1 parent f852d67 commit 85076ea

14 files changed

Lines changed: 123 additions & 146 deletions

src/groups/mqb/mqbblp/mqbblp_cluster.cpp

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -164,7 +164,9 @@ void Cluster::startDispatched(bsl::ostream* errorDescription, int* rc)
164164
bmqp_ctrlmsg::NodeStatusAdvisory& advisory =
165165
clusterMsg.choice().makeNodeStatusAdvisory();
166166
advisory.status() = bmqp_ctrlmsg::NodeStatus::E_STARTING;
167-
d_clusterData.messageTransmitter().broadcastMessage(controlMsg, true);
167+
d_clusterData.messageTransmitter().broadcastMessage(
168+
controlMsg,
169+
d_clusterData.transportManager());
168170

169171
// Start a StatMonitorSnapshotRecorder to track system stats during
170172
// recovery
@@ -327,6 +329,10 @@ void Cluster::stopDispatched()
327329
// guaranteed that this object will be kept alive.
328330

329331
// Teardown all ClusterNodeSession
332+
BALL_LOG_INFO << description() << " teardown "
333+
<< d_clusterData.membership().clusterNodeSessionMap().size()
334+
<< " ClusterNodeSession";
335+
330336
for (ClusterNodeSessionMapIter it =
331337
d_clusterData.membership().clusterNodeSessionMap().begin();
332338
it != d_clusterData.membership().clusterNodeSessionMap().end();
@@ -679,7 +685,9 @@ void Cluster::continueShutdownDispatched(
679685

680686
d_clusterData.membership().setSelfNodeStatus(
681687
bmqp_ctrlmsg::NodeStatus::E_UNAVAILABLE);
682-
d_clusterData.messageTransmitter().broadcastMessage(controlMsg, true);
688+
d_clusterData.messageTransmitter().broadcastMessage(
689+
controlMsg,
690+
d_clusterData.transportManager());
683691

684692
// Make sure all partitions done sending last sync points and advisories.
685693
// Synchronize with all Queue Dispatcher threads

src/groups/mqb/mqbblp/mqbblp_clusterorchestrator.cpp

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -797,8 +797,9 @@ void ClusterOrchestrator::transitionToAvailable()
797797
bmqp_ctrlmsg::NodeStatusAdvisory& advisory =
798798
clusterMsg.choice().makeNodeStatusAdvisory();
799799
advisory.status() = bmqp_ctrlmsg::NodeStatus::E_AVAILABLE;
800-
d_clusterData_p->messageTransmitter().broadcastMessage(controlMsg,
801-
true);
800+
d_clusterData_p->messageTransmitter().broadcastMessage(
801+
controlMsg,
802+
d_clusterData_p->transportManager());
802803

803804
d_wasAvailableAdvisorySent = true;
804805
}

src/groups/mqb/mqbblp/mqbblp_clusterproxy.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -111,7 +111,7 @@ void ClusterProxy::startDispatched()
111111
// and all session to build initial state and then schedule a refresh to
112112
// find active node if there is none after processing all pending events.
113113

114-
d_activeNodeManager.initialize(&d_clusterData.transportManager());
114+
d_activeNodeManager.initialize(d_clusterData.transportManager());
115115

116116
// Ready to read.
117117
d_clusterData.membership().netCluster()->enableRead();

src/groups/mqb/mqbc/mqbc_clusterdata.cpp

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -119,8 +119,7 @@ ClusterData::ClusterData(
119119
allocator))
120120
, d_cluster_p(cluster)
121121
, d_messageTransmitter(resources.blobSpPool(),
122-
cluster,
123-
transportManager,
122+
membership().netCluster(),
124123
allocator)
125124
, d_requestManager(bmqp::EventType::e_CONTROL,
126125
resources.blobSpPool(),

src/groups/mqb/mqbc/mqbc_clusterdata.h

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -26,14 +26,14 @@
2626

2727
// MQB
2828
#include <mqbc_clustermembership.h>
29-
#include <mqbc_controlmessagetransmitter.h>
3029
#include <mqbc_electorinfo.h>
3130
#include <mqbcfg_clusterquorummanager.h>
3231
#include <mqbcfg_messages.h>
3332
#include <mqbi_cluster.h>
3433
#include <mqbi_dispatcher.h>
3534
#include <mqbi_domain.h>
3635
#include <mqbnet_cluster.h>
36+
#include <mqbnet_controlmessagetransmitter.h>
3737
#include <mqbnet_elector.h>
3838
#include <mqbnet_multirequestmanager.h>
3939
#include <mqbnet_transportmanager.h>
@@ -185,7 +185,7 @@ class ClusterData {
185185
mqbi::Cluster* d_cluster_p;
186186

187187
/// Control message transmitter to use.
188-
ControlMessageTransmitter d_messageTransmitter;
188+
mqbnet::ControlMessageTransmitter d_messageTransmitter;
189189

190190
/// Request manager to use.
191191
RequestManagerType d_requestManager;
@@ -270,7 +270,7 @@ class ClusterData {
270270
mqbcfg::ClusterQuorumManager& quorumManager();
271271

272272
/// Get a modifiable reference to this object's messageTransmitter.
273-
ControlMessageTransmitter& messageTransmitter();
273+
mqbnet::ControlMessageTransmitter& messageTransmitter();
274274

275275
/// Get a modifiable reference to this object's requestManager.
276276
RequestManagerType& requestManager();
@@ -282,7 +282,7 @@ class ClusterData {
282282
mqbi::DomainFactory* domainFactory();
283283

284284
/// Get a modifiable reference to this object's transportManager.
285-
mqbnet::TransportManager& transportManager();
285+
mqbnet::TransportManager* transportManager();
286286

287287
/// Get a modifiable reference to this object's cluster stats.
288288
mqbstat::ClusterStats& stats();
@@ -403,7 +403,7 @@ inline mqbcfg::ClusterQuorumManager& ClusterData::quorumManager()
403403
return d_quorumManager;
404404
}
405405

406-
inline ControlMessageTransmitter& ClusterData::messageTransmitter()
406+
inline mqbnet::ControlMessageTransmitter& ClusterData::messageTransmitter()
407407
{
408408
return d_messageTransmitter;
409409
}
@@ -423,9 +423,9 @@ inline mqbi::DomainFactory* ClusterData::domainFactory()
423423
return d_domainFactory_p;
424424
}
425425

426-
inline mqbnet::TransportManager& ClusterData::transportManager()
426+
inline mqbnet::TransportManager* ClusterData::transportManager()
427427
{
428-
return *d_transportManager_p;
428+
return d_transportManager_p;
429429
}
430430

431431
inline mqbstat::ClusterStats& ClusterData::stats()

src/groups/mqb/mqbc/mqbc_storagemanager.cpp

Lines changed: 9 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -1613,8 +1613,7 @@ void StorageManager::do_replicaStateResponse(const EventWithData& event)
16131613
BSLS_ASSERT_SAFE(eventData.source()->nodeId() ==
16141614
d_partitionInfoVec[partitionId].primary()->nodeId());
16151615

1616-
d_clusterData_p->messageTransmitter().sendMessageSafe(controlMsg,
1617-
eventData.source());
1616+
fileStore(partitionId).sendMessage(controlMsg, eventData.source());
16181617

16191618
BALL_LOG_INFO << d_clusterData_p->identity().description()
16201619
<< ": sent response " << controlMsg
@@ -1673,8 +1672,7 @@ void StorageManager::do_failureReplicaStateResponse(const EventWithData& event)
16731672
<< eventData.source()->nodeDescription() << ".";
16741673
}
16751674

1676-
d_clusterData_p->messageTransmitter().sendMessageSafe(controlMsg,
1677-
eventData.source());
1675+
fileStore(partitionId).sendMessage(controlMsg, eventData.source());
16781676
}
16791677

16801678
void StorageManager::do_logFailureReplicaStateResponse(
@@ -1915,8 +1913,7 @@ void StorageManager::do_primaryStateResponse(const EventWithData& event)
19151913
response.partitionMaxFileSizes() = getSelfPartitionMaxFileSizes(
19161914
partitionId);
19171915

1918-
d_clusterData_p->messageTransmitter().sendMessageSafe(controlMsg,
1919-
eventData.source());
1916+
fileStore(partitionId).sendMessage(controlMsg, eventData.source());
19201917

19211918
BALL_LOG_INFO << d_clusterData_p->identity().description()
19221919
<< ": sent response " << controlMsg
@@ -1951,8 +1948,7 @@ void StorageManager::do_failurePrimaryStateResponse(const EventWithData& event)
19511948
response.code() = mqbi::ClusterErrorCode::e_NOT_PRIMARY;
19521949
response.message() = "Not a primary";
19531950

1954-
d_clusterData_p->messageTransmitter().sendMessageSafe(controlMsg,
1955-
eventData.source());
1951+
fileStore(partitionId).sendMessage(controlMsg, eventData.source());
19561952

19571953
BALL_LOG_INFO << d_clusterData_p->identity().description()
19581954
<< " Partition [" << partitionId
@@ -2181,8 +2177,7 @@ void StorageManager::do_replicaDataResponsePush(const EventWithData& event)
21812177
partitionId);
21822178
}
21832179

2184-
d_clusterData_p->messageTransmitter().sendMessageSafe(controlMsg,
2185-
destNode);
2180+
fileStore(partitionId).sendMessage(controlMsg, destNode);
21862181

21872182
BALL_LOG_INFO << d_clusterData_p->identity().description()
21882183
<< " Partition [" << partitionId << "]: " << "Sent response "
@@ -2394,8 +2389,7 @@ void StorageManager::do_replicaDataResponseDrop(const EventWithData& event)
23942389
eventData.partitionSeqNumDataRange().first;
23952390
response.endSequenceNumber() = eventData.partitionSeqNumDataRange().second;
23962391

2397-
d_clusterData_p->messageTransmitter().sendMessageSafe(controlMsg,
2398-
eventData.source());
2392+
fileStore(partitionId).sendMessage(controlMsg, eventData.source());
23992393

24002394
BALL_LOG_INFO << d_clusterData_p->identity().description()
24012395
<< " Partition [" << partitionId << "]: " << "Sent response "
@@ -2506,8 +2500,7 @@ void StorageManager::do_replicaDataResponsePull(const EventWithData& event)
25062500
eventData.partitionSeqNumDataRange().first;
25072501
response.endSequenceNumber() = eventData.partitionSeqNumDataRange().second;
25082502

2509-
d_clusterData_p->messageTransmitter().sendMessageSafe(controlMsg,
2510-
destNode);
2503+
fileStore(partitionId).sendMessage(controlMsg, destNode);
25112504

25122505
BALL_LOG_INFO << d_clusterData_p->identity().description()
25132506
<< " Partition [" << partitionId << "]: Sent response "
@@ -2546,8 +2539,7 @@ void StorageManager::do_failureReplicaDataResponsePull(
25462539
status.code() = mqbi::ClusterErrorCode::e_STORAGE_FAILURE;
25472540
status.message() = "Failed to send data chunks";
25482541

2549-
d_clusterData_p->messageTransmitter().sendMessageSafe(controlMsg,
2550-
destNode);
2542+
fileStore(partitionId).sendMessage(controlMsg, destNode);
25512543

25522544
BALL_LOG_INFO << d_clusterData_p->identity().description()
25532545
<< " Partition [" << partitionId
@@ -3919,8 +3911,7 @@ void StorageManager::do_replicaDataResponseResize(const EventWithData& event)
39193911
response.partitionMaxFileSizes() = getSelfPartitionMaxFileSizes(
39203912
partitionId);
39213913

3922-
d_clusterData_p->messageTransmitter().sendMessageSafe(controlMsg,
3923-
eventData.source());
3914+
fileStore(partitionId).sendMessage(controlMsg, eventData.source());
39243915

39253916
BALL_LOG_INFO << d_clusterData_p->identity().description()
39263917
<< " Partition [" << partitionId << "]: " << "Sent response "

src/groups/mqb/mqbc/mqbc_storageutil.cpp

Lines changed: 9 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1412,15 +1412,15 @@ StorageUtil::findMinReqDiskSpace(const mqbcfg::PartitionConfig& config)
14121412
return minimumRequiredDiskSpace;
14131413
}
14141414

1415-
void StorageUtil::transitionToActivePrimary(PartitionInfo* partitionInfo,
1416-
mqbc::ClusterData* clusterData,
1417-
int partitionId)
1415+
void StorageUtil::transitionToActivePrimary(PartitionInfo* partitionInfo,
1416+
mqbs::FileStore* fs,
1417+
int partitionId)
14181418
{
14191419
// executed by *QUEUE_DISPATCHER* thread associated with 'partitionId'
14201420

14211421
// PRECONDITIONS
14221422
BSLS_ASSERT_SAFE(partitionInfo);
1423-
BSLS_ASSERT_SAFE(clusterData);
1423+
BSLS_ASSERT_SAFE(fs);
14241424

14251425
partitionInfo->setPrimaryStatus(bmqp_ctrlmsg::PrimaryStatus::E_ACTIVE);
14261426

@@ -1434,7 +1434,7 @@ void StorageUtil::transitionToActivePrimary(PartitionInfo* partitionInfo,
14341434
primaryAdv.primaryLeaseId() = partitionInfo->primaryLeaseId();
14351435
primaryAdv.status() = bmqp_ctrlmsg::PrimaryStatus::E_ACTIVE;
14361436

1437-
clusterData->messageTransmitter().broadcastMessageSafe(controlMsg);
1437+
fs->broadcastMessage(controlMsg);
14381438
}
14391439

14401440
void StorageUtil::onPartitionPrimarySync(
@@ -1505,7 +1505,7 @@ void StorageUtil::onPartitionPrimarySync(
15051505

15061506
// Broadcast self as active primary of this partition. This must be done
15071507
// before invoking 'FileStore::setActivePrimary'.
1508-
transitionToActivePrimary(pinfo, clusterData, partitionId);
1508+
transitionToActivePrimary(pinfo, fs, partitionId);
15091509

15101510
partitionPrimaryStatusCb(partitionId, status, pinfo->primaryLeaseId());
15111511

@@ -3412,7 +3412,7 @@ void StorageUtil::processShutdownEventDispatched(ClusterData* clusterData,
34123412
primaryAdv.primaryLeaseId() = pinfo->primaryLeaseId();
34133413
primaryAdv.status() = bmqp_ctrlmsg::PrimaryStatus::E_PASSIVE;
34143414

3415-
clusterData->messageTransmitter().broadcastMessageSafe(controlMsg);
3415+
fs->broadcastMessage(controlMsg);
34163416
}
34173417

34183418
// Notify partition that self node is shutting down (this occurs
@@ -3862,11 +3862,10 @@ void StorageUtil::forceIssueAdvisoryAndSyncPt(mqbc::ClusterData* clusterData,
38623862
primaryAdv.status() = bmqp_ctrlmsg::PrimaryStatus::E_ACTIVE;
38633863

38643864
if (destination) {
3865-
clusterData->messageTransmitter().sendMessageSafe(controlMsg,
3866-
destination);
3865+
fs->sendMessage(controlMsg, destination);
38673866
}
38683867
else {
3869-
clusterData->messageTransmitter().broadcastMessageSafe(controlMsg);
3868+
fs->broadcastMessage(controlMsg);
38703869
}
38713870
const int rc = fs->issueSyncPoint();
38723871
bmqp_ctrlmsg::PartitionSequenceNumber psn;

src/groups/mqb/mqbc/mqbc_storageutil.h

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -540,14 +540,13 @@ struct StorageUtil {
540540

541541
/// Transition self to active primary of the specified `partitionId` and
542542
/// load this info into the specified `partitionInfo`. Then, broadcast
543-
/// a primary status advisory to peers using the specified
544-
/// `clusterData`.
543+
/// a primary status advisory to peers via the specified `fs`.
545544
///
546545
/// THREAD: Executed by the queue dispatcher thread associated with
547546
/// 'partitionId'.
548-
static void transitionToActivePrimary(PartitionInfo* partitionInfo,
549-
mqbc::ClusterData* clusterData,
550-
int partitionId);
547+
static void transitionToActivePrimary(PartitionInfo* partitionInfo,
548+
mqbs::FileStore* fs,
549+
int partitionId);
551550

552551
/// Callback executed after primary sync for the specified 'partitionId'
553552
/// is complete with the specified 'status'. Use the specified 'fs',

src/groups/mqb/mqbc/package/mqbc.mem

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@ mqbc_clusterstateledgerutil
1010
mqbc_clusterstatemanager
1111
mqbc_clusterstatetable
1212
mqbc_clusterutil
13-
mqbc_controlmessagetransmitter
1413
mqbc_electorinfo
1514
mqbc_incoreclusterstateledger
1615
mqbc_incoreclusterstateledgeriterator

0 commit comments

Comments
 (0)