Skip to content

Commit 2eccee3

Browse files
committed
Refactor[mqbs::FileStore]: sequenceNumber -> writeHeadSeqNum
Signed-off-by: Yuan Jing Vincent Yan <yyan82@bloomberg.net>
1 parent fc86c5f commit 2eccee3

9 files changed

Lines changed: 69 additions & 67 deletions

src/groups/mqb/mqbblp/mqbblp_recoverymanager.cpp

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2011,7 +2011,7 @@ int RecoveryManager::syncPeerPartition(PrimarySyncContext* primarySyncCtx,
20112011

20122012
bmqp_ctrlmsg::PartitionSequenceNumber selfSequenceNum;
20132013
selfSequenceNum.primaryLeaseId() = fs->writeHeadLeaseId();
2014-
selfSequenceNum.sequenceNumber() = fs->sequenceNumber();
2014+
selfSequenceNum.sequenceNumber() = fs->writeHeadSeqNum();
20152015

20162016
const FileTransferInfo& fti = primarySyncCtx->fileTransferInfo();
20172017

@@ -4373,7 +4373,7 @@ void RecoveryManager::startPartitionPrimarySync(
43734373

43744374
bmqp_ctrlmsg::PartitionSequenceNumber tmp;
43754375
tmp.primaryLeaseId() = fs->writeHeadLeaseId();
4376-
tmp.sequenceNumber() = fs->sequenceNumber();
4376+
tmp.sequenceNumber() = fs->writeHeadSeqNum();
43774377
primarySyncCtx.setSelfPartitionSequenceNum(tmp);
43784378

43794379
if (!fs->syncPoints().empty()) {
@@ -4469,7 +4469,7 @@ void RecoveryManager::processPartitionSyncStateRequest(
44694469

44704470
response.partitionId() = req.partitionId();
44714471
response.primaryLeaseId() = fs->writeHeadLeaseId();
4472-
response.sequenceNum() = fs->sequenceNumber();
4472+
response.sequenceNum() = fs->writeHeadSeqNum();
44734473
if (!fs->syncPoints().empty()) {
44744474
response.lastSyncPointOffsetPair() = fs->syncPoints().back();
44754475
}
@@ -4578,7 +4578,7 @@ void RecoveryManager::processPartitionSyncDataRequest(
45784578

45794579
bmqp_ctrlmsg::PartitionSequenceNumber selfPSN;
45804580
selfPSN.primaryLeaseId() = fs->writeHeadLeaseId();
4581-
selfPSN.sequenceNumber() = fs->sequenceNumber();
4581+
selfPSN.sequenceNumber() = fs->writeHeadSeqNum();
45824582

45834583
if (requesterUptoPSN <= requesterPSN) {
45844584
BALL_LOG_WARN << d_clusterData_p->identity().description()

src/groups/mqb/mqbblp/mqbblp_storagemanager.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -364,7 +364,7 @@ void StorageManager::onPartitionRecovery(
364364
: "**none**")
365365
<< ", "
366366
<< mqbs::printPSN(fs->writeHeadLeaseId(),
367-
fs->sequenceNumber())
367+
fs->writeHeadSeqNum())
368368
<< ")";
369369
}
370370
}

src/groups/mqb/mqbc/mqbc_recoverymanager.cpp

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -232,7 +232,7 @@ void RecoveryManager::setExpectedDataChunkRange(
232232
BSLS_ASSERT_SAFE(receiveDataCtx.d_currPSN.primaryLeaseId() ==
233233
fs.writeHeadLeaseId());
234234
BSLS_ASSERT_SAFE(receiveDataCtx.d_currPSN.sequenceNumber() ==
235-
fs.sequenceNumber());
235+
fs.writeHeadSeqNum());
236236
}
237237
else {
238238
BSLS_ASSERT_SAFE(recoveryCtx.d_mappedJournalFd.isValid() &&
@@ -646,12 +646,12 @@ int RecoveryManager::processReceiveDataChunks(
646646
BSLS_ASSERT_SAFE(receiveDataCtx.d_currPSN.primaryLeaseId() ==
647647
fs->writeHeadLeaseId());
648648
BSLS_ASSERT_SAFE(receiveDataCtx.d_currPSN.sequenceNumber() ==
649-
fs->sequenceNumber());
649+
fs->writeHeadSeqNum());
650650

651651
fs->processStorageEvent(blob, true /* isPartitionSyncEvent */, source);
652652

653653
receiveDataCtx.d_currPSN.primaryLeaseId() = fs->writeHeadLeaseId();
654-
receiveDataCtx.d_currPSN.sequenceNumber() = fs->sequenceNumber();
654+
receiveDataCtx.d_currPSN.sequenceNumber() = fs->writeHeadSeqNum();
655655

656656
if (receiveDataCtx.d_currPSN == receiveDataCtx.d_endPSN) {
657657
receiveDataCtx.d_expectChunks = false;

src/groups/mqb/mqbc/mqbc_storagemanager.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1925,7 +1925,7 @@ void StorageManager::do_storeSelfSeq(
19251925
BSLS_ASSERT_SAFE(fs);
19261926
if (fs->isOpen()) {
19271927
nodePSNCtx.d_PSN.primaryLeaseId() = fs->writeHeadLeaseId();
1928-
nodePSNCtx.d_PSN.sequenceNumber() = fs->sequenceNumber();
1928+
nodePSNCtx.d_PSN.sequenceNumber() = fs->writeHeadSeqNum();
19291929
}
19301930
else {
19311931
const int rc = d_recoveryManager_mp->recoverPSN(&nodePSNCtx.d_PSN,

src/groups/mqb/mqbc/mqbc_storagemanager.t.cpp

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2924,7 +2924,7 @@ static void test18_primaryHealingWatchdogRetry()
29242924
const int k_PRIMARY_LEASE_ID =
29252925
storageManager.fileStore(k_PARTITION_ID).writeHeadLeaseId();
29262926
const int k_PRIMARY_SEQ_NUM =
2927-
storageManager.fileStore(k_PARTITION_ID).sequenceNumber();
2927+
storageManager.fileStore(k_PARTITION_ID).writeHeadSeqNum();
29282928

29292929
static const int k_REQUEST_ID = 1;
29302930
bmqp_ctrlmsg::ControlMessage message;
@@ -3555,7 +3555,7 @@ static void test23_replicaHealingReceivesReplicaDataRqstDropInvalidPid()
35553555

35563556
// 5. Send a storage event (PUT) and verify it is buffered, not processed
35573557
const bsls::Types::Uint64 seqNumBefore =
3558-
storageManager.fileStore(k_PARTITION_ID).sequenceNumber();
3558+
storageManager.fileStore(k_PARTITION_ID).writeHeadSeqNum();
35593559

35603560
bmqp::StorageEventBuilder seb(mqbs::FileStoreProtocol::k_VERSION,
35613561
bmqp::EventType::e_STORAGE,
@@ -3593,8 +3593,9 @@ static void test23_replicaHealingReceivesReplicaDataRqstDropInvalidPid()
35933593

35943594
// Sequence number has not advanced; this proves that we did not process
35953595
// the PUT.
3596-
BMQTST_ASSERT_EQ(storageManager.fileStore(k_PARTITION_ID).sequenceNumber(),
3597-
seqNumBefore);
3596+
BMQTST_ASSERT_EQ(
3597+
storageManager.fileStore(k_PARTITION_ID).writeHeadSeqNum(),
3598+
seqNumBefore);
35983599

35993600
BMQTST_ASSERT_EQ(storageManager.partitionHealthState(k_PARTITION_ID),
36003601
mqbc::PartitionFSM::State::e_REPLICA_HEALING);

src/groups/mqb/mqbc/mqbc_storageutil.cpp

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3735,7 +3735,7 @@ void StorageUtil::forceIssueAdvisoryAndSyncPt(mqbc::ClusterData* clusterData,
37353735
<< fs->config().partitionId()
37363736
<< "]: successfully issued a forced SyncPt: "
37373737
<< mqbs::printPSN(fs->writeHeadLeaseId(),
3738-
fs->sequenceNumber())
3738+
fs->writeHeadSeqNum())
37393739
<< ".";
37403740
}
37413741
else {
@@ -3744,7 +3744,7 @@ void StorageUtil::forceIssueAdvisoryAndSyncPt(mqbc::ClusterData* clusterData,
37443744
<< "]: failed to force-issue SyncPt, rc: " << rc
37453745
<< ", current PSN: "
37463746
<< mqbs::printPSN(fs->writeHeadLeaseId(),
3747-
fs->sequenceNumber());
3747+
fs->writeHeadSeqNum());
37483748
}
37493749
}
37503750

0 commit comments

Comments
 (0)