Skip to content

Commit 66fb6a2

Browse files
committed
Refactor[mqbs::FileStore]: Encapsulate write-head seqNum increment
Signed-off-by: Yuan Jing Vincent Yan <yyan82@bloomberg.net>
1 parent 2eccee3 commit 66fb6a2

2 files changed

Lines changed: 15 additions & 17 deletions

File tree

src/groups/mqb/mqbs/mqbs_filestore.cpp

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1007,7 +1007,7 @@ int FileStore::openInRecoveryMode(bsl::ostream& errorDescription,
10071007

10081008
bmqp_ctrlmsg::SyncPoint syncPoint;
10091009
syncPoint.primaryLeaseId() = d_writeHeadLeaseId;
1010-
syncPoint.sequenceNum() = currentSeqNumRef() + 1;
1010+
syncPoint.sequenceNum() = writeHeadSeqNum() + 1;
10111011
syncPoint.dataFileOffsetDwords() = fileSetSp->d_data.d_filePosition /
10121012
bmqp::Protocol::k_DWORD_SIZE;
10131013
syncPoint.qlistFileOffsetWords() =
@@ -3300,7 +3300,7 @@ int FileStore::rollover()
33003300

33013301
bmqp_ctrlmsg::SyncPoint syncPt;
33023302
syncPt.primaryLeaseId() = d_writeHeadLeaseId;
3303-
syncPt.sequenceNum() = currentSeqNumRef() + 1;
3303+
syncPt.sequenceNum() = writeHeadSeqNum() + 1;
33043304
syncPt.dataFileOffsetDwords() = activeFileSet->d_data.d_filePosition /
33053305
bmqp::Protocol::k_DWORD_SIZE;
33063306
syncPt.qlistFileOffsetWords() =
@@ -3747,7 +3747,7 @@ void FileStore::writeQueueOpRecordImpl(DataStoreRecordHandle* handle,
37473747
new (qRec.get()) QueueOpRecord();
37483748
qRec->header()
37493749
.setPrimaryLeaseId(d_writeHeadLeaseId)
3750-
.setSequenceNumber(++currentSeqNumRef())
3750+
.setSequenceNumber(incrementWriteHeadSeqNum())
37513751
.setTimestamp(timestamp);
37523752
qRec->setQueueKey(queueKey).setType(queueOpFlag);
37533753
qRec->setStartSequenceNumber(startSequenceNum);
@@ -4104,7 +4104,7 @@ void FileStore::issueSyncPointIfNeeded()
41044104

41054105
bmqp_ctrlmsg::SyncPoint sp;
41064106
sp.primaryLeaseId() = d_writeHeadLeaseId;
4107-
sp.sequenceNum() = currentSeqNumRef() + 1;
4107+
sp.sequenceNum() = writeHeadSeqNum() + 1;
41084108
sp.dataFileOffsetDwords() = fs->d_data.d_filePosition /
41094109
bmqp::Protocol::k_DWORD_SIZE;
41104110
if (d_qListAware) {
@@ -5752,7 +5752,7 @@ int FileStore::writeMessageRecord(mqbi::StorageMessageAttributes* attributes,
57525752
new (msgRec.get()) MessageRecord();
57535753
msgRec->header()
57545754
.setPrimaryLeaseId(d_writeHeadLeaseId)
5755-
.setSequenceNumber(++currentSeqNumRef())
5755+
.setSequenceNumber(incrementWriteHeadSeqNum())
57565756
.setTimestamp(attributes->arrivalTimestamp());
57575757
msgRec->setRefCount(attributes->refCount())
57585758
.setQueueKey(queueKey)
@@ -6030,7 +6030,7 @@ int FileStore::writeQueueCreationRecord(DataStoreRecordHandle* handle,
60306030
new (queueOpRec.get()) QueueOpRecord();
60316031
queueOpRec->header()
60326032
.setPrimaryLeaseId(d_writeHeadLeaseId)
6033-
.setSequenceNumber(++currentSeqNumRef())
6033+
.setSequenceNumber(incrementWriteHeadSeqNum())
60346034
.setTimestamp(timestamp);
60356035
queueOpRec->setQueueKey(queueKey).setType(
60366036
isNewQueue ? QueueOpType::e_CREATION : QueueOpType::e_ADDITION);
@@ -6167,7 +6167,7 @@ int FileStore::writeConfirmRecord(DataStoreRecordHandle* handle,
61676167
new (confRec.get()) ConfirmRecord();
61686168
confRec->header()
61696169
.setPrimaryLeaseId(d_writeHeadLeaseId)
6170-
.setSequenceNumber(++currentSeqNumRef())
6170+
.setSequenceNumber(incrementWriteHeadSeqNum())
61716171
.setTimestamp(timestamp);
61726172
confRec->setQueueKey(queueKey).setMessageGUID(guid);
61736173

@@ -6235,7 +6235,7 @@ int FileStore::writeDeletionRecord(const bmqt::MessageGUID& guid,
62356235
new (delRec.get()) DeletionRecord();
62366236
delRec->header()
62376237
.setPrimaryLeaseId(d_writeHeadLeaseId)
6238-
.setSequenceNumber(++currentSeqNumRef())
6238+
.setSequenceNumber(incrementWriteHeadSeqNum())
62396239
.setTimestamp(timestamp);
62406240
delRec->setDeletionRecordFlag(deletionFlag)
62416241
.setQueueKey(queueKey)
@@ -6264,7 +6264,7 @@ int FileStore::writeSyncPointRecord(const bmqp_ctrlmsg::SyncPoint& syncPoint,
62646264
}
62656265

62666266
// Update the PSN only when record writing is guaranteed.
6267-
++currentSeqNumRef();
6267+
incrementWriteHeadSeqNum();
62686268

62696269
// Local refs for convenience.
62706270

@@ -6655,7 +6655,7 @@ int FileStore::processRecoveryEvent(const bsl::shared_ptr<bdlbb::Blob>& blob)
66556655
return rc_INVALID_SEQ_NUM; // RETURN
66566656
}
66576657

6658-
++currentSeqNumRef();
6658+
incrementWriteHeadSeqNum();
66596659
}
66606660
}
66616661

@@ -6855,7 +6855,7 @@ int FileStore::issueSyncPoint()
68556855

68566856
bmqp_ctrlmsg::SyncPoint syncPoint;
68576857
syncPoint.primaryLeaseId() = d_writeHeadLeaseId;
6858-
syncPoint.sequenceNum() = currentSeqNumRef() + 1;
6858+
syncPoint.sequenceNum() = writeHeadSeqNum() + 1;
68596859
syncPoint.dataFileOffsetDwords() = fs->d_data.d_filePosition /
68606860
bmqp::Protocol::k_DWORD_SIZE;
68616861
syncPoint.qlistFileOffsetWords() = d_qListAware

src/groups/mqb/mqbs/mqbs_filestore.h

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -391,10 +391,8 @@ class FileStore BSLS_KEYWORD_FINAL : public DataStore {
391391
private:
392392
// PRIVATE MANIPULATORS
393393

394-
/// Return a mutable reference to the current sequence number entry in
395-
/// `d_highestSeqNums`, i.e., `d_highestSeqNums[d_writeHeadLeaseId]`.
396-
/// Note that this will insert a zero entry if one does not exist.
397-
bsls::Types::Uint64& currentSeqNumRef();
394+
/// Increment the write head's sequence number and return its new value.
395+
bsls::Types::Uint64 incrementWriteHeadSeqNum();
398396

399397
/// Move the internal write cursor to the specified `leaseId` and `seqNum`.
400398
void setWriteHead(unsigned int leaseId, bsls::Types::Uint64 seqNum);
@@ -1163,9 +1161,9 @@ inline FileStore::NodeContext::NodeContext(BlobSpPool* blobSpPool_p,
11631161
// ---------------
11641162

11651163
// PRIVATE MANIPULATORS
1166-
inline bsls::Types::Uint64& FileStore::currentSeqNumRef()
1164+
inline bsls::Types::Uint64 FileStore::incrementWriteHeadSeqNum()
11671165
{
1168-
return d_highestSeqNums[d_writeHeadLeaseId];
1166+
return ++d_highestSeqNums.at(d_writeHeadLeaseId);
11691167
}
11701168

11711169
inline void FileStore::setWriteHead(unsigned int leaseId,

0 commit comments

Comments
 (0)