From 22aaf101fcea2f87cc1d4b2e78ae28e06ac045de Mon Sep 17 00:00:00 2001 From: Robert Nagy Date: Mon, 20 Jul 2026 11:35:55 +0200 Subject: [PATCH 1/3] Expose exact transaction log positions --- db/db_impl/db_impl_write.cc | 1 + db/db_log_iter_test.cc | 182 ++++++++++++++++++++++++++++++ db/transaction_log_impl.cc | 111 ++++++++++++++++-- db/transaction_log_impl.h | 26 ++++- include/rocksdb/transaction_log.h | 42 +++++-- 5 files changed, 343 insertions(+), 19 deletions(-) diff --git a/db/db_impl/db_impl_write.cc b/db/db_impl/db_impl_write.cc index 8a4c5ec9be6c..039f469cafe1 100644 --- a/db/db_impl/db_impl_write.cc +++ b/db/db_impl/db_impl_write.cc @@ -768,6 +768,7 @@ Status DBImpl::WriteImpl(const WriteOptions& write_options, wal_context.need_wal_sync, wal_context.need_wal_dir_sync, last_sequence + 1, *wal_context.wal_file_number_size); + TEST_SYNC_POINT("DBImpl::WriteImpl:AfterWriteWAL"); } } else { if (status.ok() && !write_options.disableWAL) { diff --git a/db/db_log_iter_test.cc b/db/db_log_iter_test.cc index 62b1f893d5c2..b1a11d5984cd 100644 --- a/db/db_log_iter_test.cc +++ b/db/db_log_iter_test.cc @@ -11,6 +11,9 @@ // in Release build. // which is a pity, it is a good test +#include +#include + #include "db/db_test_util.h" #include "env/mock_env.h" #include "port/stack_trace.h" @@ -332,6 +335,185 @@ TEST_F(DBTestXactLogIterator, TransactionLogIteratorBlobs) { "Delete(0, key2)", handler.seen); } + +TEST_F(DBTestXactLogIterator, + TransactionLogIteratorLogDataOnlyTailAndWalPosition) { + Options options = OptionsForLogIterTest(); + DestroyAndReopen(options); + ASSERT_OK(Put("key1", "value1")); + + auto iter = OpenTransactionLogIter(0); + BatchResult initial = iter->GetBatch(); + ASSERT_EQ(initial.writeBatchPtr->Count(), 1U); + iter->Next(); + ASSERT_FALSE(iter->Valid()); + ASSERT_OK(iter->status()); + + WriteBatch first_log_data; + ASSERT_OK(first_log_data.PutLogData(Slice("log-data-1"))); + ASSERT_OK(dbfull()->Write(WriteOptions(), &first_log_data)); + + iter->Next(); + ASSERT_TRUE(iter->Valid()); + BatchResult first = iter->GetBatch(); + ASSERT_EQ(first.sequence, 2U); + ASSERT_EQ(first.writeBatchPtr->Count(), 0U); + ASSERT_GT(first.wal_file_number, 0U); + ASSERT_LT(first.wal_record_offset, first.wal_record_end); + + iter->Next(); + ASSERT_FALSE(iter->Valid()); + ASSERT_OK(iter->status()); + + WriteBatch second_log_data; + ASSERT_OK(second_log_data.PutLogData(Slice("log-data-2"))); + ASSERT_OK(dbfull()->Write(WriteOptions(), &second_log_data)); + + iter->Next(); + ASSERT_TRUE(iter->Valid()); + BatchResult second = iter->GetBatch(); + ASSERT_EQ(second.sequence, 2U); + ASSERT_EQ(second.writeBatchPtr->Count(), 0U); + ASSERT_EQ(first.wal_file_number, second.wal_file_number); + ASSERT_LT(first.wal_record_offset, second.wal_record_offset); + ASSERT_LE(first.wal_record_end, second.wal_record_offset); + ASSERT_NE(first.wal_record_checksum, second.wal_record_checksum); + + const uint64_t first_file_number = first.wal_file_number; + const uint64_t first_offset = first.wal_record_offset; + const uint64_t first_end = first.wal_record_end; + const uint64_t first_checksum = first.wal_record_checksum; + const uint64_t second_file_number = second.wal_file_number; + const uint64_t second_offset = second.wal_record_offset; + const uint64_t second_end = second.wal_record_end; + const uint64_t second_checksum = second.wal_record_checksum; + + BatchResult move_constructed(std::move(first)); + ASSERT_EQ(move_constructed.wal_file_number, first_file_number); + ASSERT_EQ(move_constructed.wal_record_offset, first_offset); + ASSERT_EQ(move_constructed.wal_record_end, first_end); + ASSERT_EQ(move_constructed.wal_record_checksum, first_checksum); + + BatchResult moved; + moved = std::move(second); + ASSERT_EQ(moved.wal_file_number, second_file_number); + ASSERT_EQ(moved.wal_record_offset, second_offset); + ASSERT_EQ(moved.wal_record_end, second_end); + ASSERT_EQ(moved.wal_record_checksum, second_checksum); + + iter->Next(); + ASSERT_FALSE(iter->Valid()); + ASSERT_OK(iter->status()); + ASSERT_EQ(dbfull()->GetLatestSequenceNumber(), 1U); + + auto replay = OpenTransactionLogIter(0); + ASSERT_EQ(replay->GetBatch().writeBatchPtr->Count(), 1U); + replay->Next(); + ASSERT_TRUE(replay->Valid()); + BatchResult replay_first = replay->GetBatch(); + ASSERT_EQ(replay_first.wal_file_number, first_file_number); + ASSERT_EQ(replay_first.wal_record_offset, first_offset); + ASSERT_EQ(replay_first.wal_record_end, first_end); + ASSERT_EQ(replay_first.wal_record_checksum, first_checksum); + + replay->Next(); + ASSERT_TRUE(replay->Valid()); + BatchResult replay_second = replay->GetBatch(); + ASSERT_EQ(replay_second.wal_file_number, second_file_number); + ASSERT_EQ(replay_second.wal_record_offset, second_offset); + ASSERT_EQ(replay_second.wal_record_end, second_end); + ASSERT_EQ(replay_second.wal_record_checksum, second_checksum); + replay->Next(); + ASSERT_FALSE(replay->Valid()); + ASSERT_OK(replay->status()); +} + +TEST_F(DBTestXactLogIterator, + TransactionLogIteratorLogDataOnlyAtInitialSequence) { + Options options = OptionsForLogIterTest(); + DestroyAndReopen(options); + + WriteBatch first_log_data; + ASSERT_OK(first_log_data.PutLogData(Slice("log-data-1"))); + ASSERT_OK(dbfull()->Write(WriteOptions(), &first_log_data)); + WriteBatch second_log_data; + ASSERT_OK(second_log_data.PutLogData(Slice("log-data-2"))); + ASSERT_OK(dbfull()->Write(WriteOptions(), &second_log_data)); + ASSERT_EQ(dbfull()->GetLatestSequenceNumber(), 0U); + + auto iter = OpenTransactionLogIter(0); + BatchResult first = iter->GetBatch(); + ASSERT_EQ(first.sequence, 1U); + ASSERT_EQ(first.writeBatchPtr->Count(), 0U); + + iter->Next(); + ASSERT_TRUE(iter->Valid()); + BatchResult second = iter->GetBatch(); + ASSERT_EQ(second.sequence, 1U); + ASSERT_EQ(second.writeBatchPtr->Count(), 0U); + ASSERT_EQ(first.wal_file_number, second.wal_file_number); + ASSERT_LT(first.wal_record_offset, second.wal_record_offset); + ASSERT_NE(first.wal_record_checksum, second.wal_record_checksum); + + iter->Next(); + ASSERT_FALSE(iter->Valid()); + ASSERT_OK(iter->status()); +} + +#ifndef NDEBUG +TEST_F(DBTestXactLogIterator, + TransactionLogIteratorDoesNotExposeUnpublishedConsumingBatch) { + Options options = OptionsForLogIterTest(); + DestroyAndReopen(options); + ASSERT_OK(Put("key1", "value1")); + + auto iter = OpenTransactionLogIter(0); + ASSERT_EQ(iter->GetBatch().writeBatchPtr->Count(), 1U); + iter->Next(); + ASSERT_FALSE(iter->Valid()); + ASSERT_OK(iter->status()); + + std::atomic wal_written{false}; + std::atomic allow_publish{false}; + SyncPoint::GetInstance()->SetCallBack( + "DBImpl::WriteImpl:AfterWriteWAL", [&](void*) { + wal_written.store(true, std::memory_order_release); + while (!allow_publish.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + }); + SyncPoint::GetInstance()->EnableProcessing(); + + Status write_status; + std::thread writer([&]() { write_status = Put("key2", "value2"); }); + while (!wal_written.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + + iter->Next(); + const bool exposed_before_publish = iter->Valid(); + const Status status_before_publish = iter->status(); + + allow_publish.store(true, std::memory_order_release); + writer.join(); + SyncPoint::GetInstance()->DisableProcessing(); + SyncPoint::GetInstance()->ClearAllCallBacks(); + + ASSERT_OK(write_status); + ASSERT_FALSE(exposed_before_publish); + ASSERT_OK(status_before_publish); + ASSERT_EQ(dbfull()->GetLatestSequenceNumber(), 2U); + + iter->Next(); + ASSERT_TRUE(iter->Valid()); + BatchResult published = iter->GetBatch(); + ASSERT_EQ(published.sequence, 2U); + ASSERT_EQ(published.writeBatchPtr->Count(), 1U); + iter->Next(); + ASSERT_FALSE(iter->Valid()); + ASSERT_OK(iter->status()); +} +#endif } // namespace ROCKSDB_NAMESPACE int main(int argc, char** argv) { diff --git a/db/transaction_log_impl.cc b/db/transaction_log_impl.cc index 7c8e20435e2d..af5b6f3ac366 100644 --- a/db/transaction_log_impl.cc +++ b/db/transaction_log_impl.cc @@ -9,6 +9,7 @@ #include "db/write_batch_internal.h" #include "file/sequence_file_reader.h" +#include "util/coding.h" #include "util/defer.h" namespace ROCKSDB_NAMESPACE { @@ -32,7 +33,20 @@ TransactionLogIteratorImpl::TransactionLogIteratorImpl( is_valid_(false), current_file_index_(0), current_batch_seq_(0), - current_last_seq_(0) { + current_last_seq_(0), + current_batch_wal_file_number_(0), + current_batch_wal_record_offset_(0), + current_batch_wal_record_end_(0), + current_batch_wal_record_checksum_(0), + current_record_wal_file_number_(0), + current_record_wal_record_offset_(0), + current_record_wal_record_end_(0), + current_record_wal_record_checksum_(0), + has_unpublished_record_(false), + unpublished_record_wal_file_number_(0), + unpublished_record_wal_record_offset_(0), + unpublished_record_wal_record_end_(0), + unpublished_record_wal_record_checksum_(0) { assert(files_ != nullptr); assert(versions_ != nullptr); assert(!seq_per_batch_); @@ -75,6 +89,10 @@ BatchResult TransactionLogIteratorImpl::GetBatch() { assert(is_valid_); // cannot call in a non valid state. BatchResult result; result.sequence = current_batch_seq_; + result.wal_file_number = current_batch_wal_file_number_; + result.wal_record_offset = current_batch_wal_record_offset_; + result.wal_record_end = current_batch_wal_record_end_; + result.wal_record_checksum = current_batch_wal_record_checksum_; result.writeBatchPtr = std::move(current_batch_); return result; } @@ -83,12 +101,65 @@ Status TransactionLogIteratorImpl::status() { return current_status_; } bool TransactionLogIteratorImpl::Valid() { return started_ && is_valid_; } +bool TransactionLogIteratorImpl::IsRecordPublished(const Slice& record) const { + if (record.size() < WriteBatchInternal::kHeader) { + // Let the caller preserve the existing corruption-reporting behavior. + return true; + } + + const uint64_t count = DecodeFixed32(record.data() + 8); + if (count == 0) { + return true; + } + + const SequenceNumber sequence = DecodeFixed64(record.data()); + const SequenceNumber published_sequence = versions_->LastSequence(); + return sequence <= published_sequence && + count - 1 <= published_sequence - sequence; +} + bool TransactionLogIteratorImpl::RestrictedRead(Slice* record) { - // Don't read if no more complete entries to read from logs - if (current_last_seq_ >= versions_->LastSequence()) { + if (has_unpublished_record_) { + Slice unpublished_record(unpublished_record_); + if (!IsRecordPublished(unpublished_record)) { + return false; + } + + scratch_ = std::move(unpublished_record_); + has_unpublished_record_ = false; + *record = Slice(scratch_); + current_record_wal_file_number_ = unpublished_record_wal_file_number_; + current_record_wal_record_offset_ = unpublished_record_wal_record_offset_; + current_record_wal_record_end_ = unpublished_record_wal_record_end_; + current_record_wal_record_checksum_ = + unpublished_record_wal_record_checksum_; + return true; + } + + uint64_t record_checksum = 0; + if (!current_log_reader_->ReadRecord( + record, &scratch_, WALRecoveryMode::kTolerateCorruptedTailRecords, + &record_checksum)) { return false; } - return current_log_reader_->ReadRecord(record, &scratch_); + + current_record_wal_file_number_ = current_log_reader_->GetLogNumber(); + current_record_wal_record_offset_ = current_log_reader_->LastRecordOffset(); + current_record_wal_record_end_ = current_log_reader_->LastRecordEnd(); + current_record_wal_record_checksum_ = record_checksum; + + if (!IsRecordPublished(*record)) { + unpublished_record_.assign(record->data(), record->size()); + has_unpublished_record_ = true; + unpublished_record_wal_file_number_ = current_record_wal_file_number_; + unpublished_record_wal_record_offset_ = current_record_wal_record_offset_; + unpublished_record_wal_record_end_ = current_record_wal_record_end_; + unpublished_record_wal_record_checksum_ = + current_record_wal_record_checksum_; + return false; + } + + return true; } void TransactionLogIteratorImpl::SeekToStartSequence(uint64_t start_file_index, @@ -112,12 +183,14 @@ void TransactionLogIteratorImpl::SeekToStartSequence(uint64_t start_file_index, } else if (!current_status_.ok()) { return; } - Status s = - OpenLogReader(files_->at(static_cast(start_file_index)).get()); - if (!s.ok()) { - current_status_ = s; - reporter_.Info(current_status_.ToString().c_str()); - return; + if (!has_unpublished_record_) { + Status s = + OpenLogReader(files_->at(static_cast(start_file_index)).get()); + if (!s.ok()) { + current_status_ = s; + reporter_.Info(current_status_.ToString().c_str()); + return; + } } while (RestrictedRead(&record)) { if (record.size() < WriteBatchInternal::kHeader) { @@ -146,6 +219,11 @@ void TransactionLogIteratorImpl::SeekToStartSequence(uint64_t start_file_index, } } + if (has_unpublished_record_) { + current_status_ = Status::OK(); + return; + } + // Could not find start sequence in first file. Normally this must be the // only file. Otherwise log the error and let the iterator return next entry // If strict is set, we want to seek exactly till the start sequence and it @@ -203,6 +281,15 @@ void TransactionLogIteratorImpl::NextImpl(bool internal) { } } + if (has_unpublished_record_) { + // A complete record has reached the WAL, but its sequence range is not + // published yet. Keep the iterator resumable at this exact reader + // position instead of opening another WAL or reporting a false gap. + is_valid_ = false; + current_status_ = Status::OK(); + return; + } + // Open the next file if (current_file_index_ < files_->size() - 1) { ++current_file_index_; @@ -276,6 +363,10 @@ void TransactionLogIteratorImpl::UpdateCurrentWriteBatch(const Slice& record) { assert(current_last_seq_ <= versions_->LastSequence()); current_batch_ = std::move(batch); + current_batch_wal_file_number_ = current_record_wal_file_number_; + current_batch_wal_record_offset_ = current_record_wal_record_offset_; + current_batch_wal_record_end_ = current_record_wal_record_end_; + current_batch_wal_record_checksum_ = current_record_wal_record_checksum_; is_valid_ = true; current_status_ = Status::OK(); } diff --git a/db/transaction_log_impl.h b/db/transaction_log_impl.h index 437d3a881d79..ebe7c014ecd0 100644 --- a/db/transaction_log_impl.h +++ b/db/transaction_log_impl.h @@ -109,8 +109,32 @@ class TransactionLogIteratorImpl : public TransactionLogIterator { SequenceNumber current_batch_seq_; // sequence number at start of current batch SequenceNumber current_last_seq_; // last sequence in the current batch - // Reads from transaction log only if the writebatch record has been written + uint64_t current_batch_wal_file_number_; + uint64_t current_batch_wal_record_offset_; + uint64_t current_batch_wal_record_end_; + uint64_t current_batch_wal_record_checksum_; + + // Metadata for the record most recently returned by RestrictedRead(). + uint64_t current_record_wal_file_number_; + uint64_t current_record_wal_record_offset_; + uint64_t current_record_wal_record_end_; + uint64_t current_record_wal_record_checksum_; + + // A complete consuming record can reach the WAL before its ending sequence + // is published. Keep it here so a later Next() can return it without either + // exposing it early or losing the reader position. + bool has_unpublished_record_; + std::string unpublished_record_; + uint64_t unpublished_record_wal_file_number_; + uint64_t unpublished_record_wal_record_offset_; + uint64_t unpublished_record_wal_record_end_; + uint64_t unpublished_record_wal_record_checksum_; + + // Reads the next complete transaction-log record. Zero-count records are + // immediately visible. Consuming records are returned only after their + // ending sequence has been published. bool RestrictedRead(Slice* record); + bool IsRecordPublished(const Slice& record) const; // Seeks to starting_sequence_number_ reading from start_file_index in files_. // If strict is set, then must get a batch starting with // starting_sequence_number_. diff --git a/include/rocksdb/transaction_log.h b/include/rocksdb/transaction_log.h index d3e335748634..8741af340455 100644 --- a/include/rocksdb/transaction_log.h +++ b/include/rocksdb/transaction_log.h @@ -62,6 +62,20 @@ using LogFile = WalFile; struct BatchResult { SequenceNumber sequence = 0; + + // The WAL file number containing this batch. Together with + // wal_record_offset, this identifies the batch while the WAL is retained. + uint64_t wal_file_number = 0; + + // The physical byte offset of the first fragment of this logical WAL record. + uint64_t wal_record_offset = 0; + + // The first physical byte offset after this logical WAL record. + uint64_t wal_record_end = 0; + + // XXH3 checksum of the logical WAL record contents. + uint64_t wal_record_checksum = 0; + std::unique_ptr writeBatchPtr; // Add empty __ctor and __dtor for the rule of five @@ -75,12 +89,23 @@ struct BatchResult { BatchResult& operator=(const BatchResult&) = delete; - BatchResult(BatchResult&& bResult) - : sequence(std::move(bResult.sequence)), + BatchResult(BatchResult&& bResult) noexcept + : sequence(bResult.sequence), + wal_file_number(bResult.wal_file_number), + wal_record_offset(bResult.wal_record_offset), + wal_record_end(bResult.wal_record_end), + wal_record_checksum(bResult.wal_record_checksum), writeBatchPtr(std::move(bResult.writeBatchPtr)) {} - BatchResult& operator=(BatchResult&& bResult) { - sequence = std::move(bResult.sequence); + BatchResult& operator=(BatchResult&& bResult) noexcept { + if (this == &bResult) { + return *this; + } + sequence = bResult.sequence; + wal_file_number = bResult.wal_file_number; + wal_record_offset = bResult.wal_record_offset; + wal_record_end = bResult.wal_record_end; + wal_record_checksum = bResult.wal_record_checksum; writeBatchPtr = std::move(bResult.writeBatchPtr); return *this; } @@ -99,12 +124,13 @@ class TransactionLogIterator { // Can read data from a valid iterator. virtual bool Valid() = 0; - // Moves the iterator to the next WriteBatch. - // REQUIRES: Valid() to be true. + // Moves the iterator to the next WriteBatch. At a clean WAL tail, callers + // may append more writes and call Next() again to continue tailing. + // REQUIRES: Valid() is true, or Valid() is false and status() is OK. virtual void Next() = 0; - // Returns ok if the iterator is valid. - // Returns the Error when something has gone wrong. + // Returns OK if the iterator is valid or stopped at a clean WAL tail. + // Returns the error when something has gone wrong. virtual Status status() = 0; // If valid return's the current write_batch and the sequence number of the From d4a79a5de0b4fa22f5911426353014e80d242bda Mon Sep 17 00:00:00 2001 From: Robert Nagy Date: Mon, 20 Jul 2026 12:13:13 +0200 Subject: [PATCH 2/3] Make WAL cursor checksums opt in --- db/db_log_iter_test.cc | 46 ++++++++++++++++++++++++------- db/transaction_log_impl.cc | 6 ++-- include/rocksdb/transaction_log.h | 13 +++++++-- 3 files changed, 50 insertions(+), 15 deletions(-) diff --git a/db/db_log_iter_test.cc b/db/db_log_iter_test.cc index b1a11d5984cd..dd13b4a52d9e 100644 --- a/db/db_log_iter_test.cc +++ b/db/db_log_iter_test.cc @@ -12,12 +12,14 @@ // which is a pity, it is a good test #include +#include #include #include "db/db_test_util.h" #include "env/mock_env.h" #include "port/stack_trace.h" #include "util/atomic.h" +#include "util/defer.h" namespace ROCKSDB_NAMESPACE { @@ -27,9 +29,11 @@ class DBTestXactLogIterator : public DBTestBase { : DBTestBase("db_log_iter_test", /*env_do_fsync=*/true) {} std::unique_ptr OpenTransactionLogIter( - const SequenceNumber seq) { + const SequenceNumber seq, const bool include_wal_record_checksum = false) { std::unique_ptr iter; - Status status = dbfull()->GetUpdatesSince(seq, &iter); + TransactionLogIterator::ReadOptions read_options; + read_options.include_wal_record_checksum_ = include_wal_record_checksum; + Status status = dbfull()->GetUpdatesSince(seq, &iter, read_options); EXPECT_OK(status); EXPECT_TRUE(iter->Valid()); return iter; @@ -342,7 +346,10 @@ TEST_F(DBTestXactLogIterator, DestroyAndReopen(options); ASSERT_OK(Put("key1", "value1")); - auto iter = OpenTransactionLogIter(0); + auto default_iter = OpenTransactionLogIter(0); + ASSERT_EQ(default_iter->GetBatch().wal_record_checksum, 0U); + + auto iter = OpenTransactionLogIter(0, true); BatchResult initial = iter->GetBatch(); ASSERT_EQ(initial.writeBatchPtr->Count(), 1U); iter->Next(); @@ -406,7 +413,7 @@ TEST_F(DBTestXactLogIterator, ASSERT_OK(iter->status()); ASSERT_EQ(dbfull()->GetLatestSequenceNumber(), 1U); - auto replay = OpenTransactionLogIter(0); + auto replay = OpenTransactionLogIter(0, true); ASSERT_EQ(replay->GetBatch().writeBatchPtr->Count(), 1U); replay->Next(); ASSERT_TRUE(replay->Valid()); @@ -441,7 +448,7 @@ TEST_F(DBTestXactLogIterator, ASSERT_OK(dbfull()->Write(WriteOptions(), &second_log_data)); ASSERT_EQ(dbfull()->GetLatestSequenceNumber(), 0U); - auto iter = OpenTransactionLogIter(0); + auto iter = OpenTransactionLogIter(0, true); BatchResult first = iter->GetBatch(); ASSERT_EQ(first.sequence, 1U); ASSERT_EQ(first.writeBatchPtr->Count(), 0U); @@ -475,20 +482,40 @@ TEST_F(DBTestXactLogIterator, std::atomic wal_written{false}; std::atomic allow_publish{false}; + std::atomic publish_wait_timed_out{false}; + std::thread writer; + Defer cleanup([&]() { + allow_publish.store(true, std::memory_order_release); + if (writer.joinable()) { + writer.join(); + } + SyncPoint::GetInstance()->DisableProcessing(); + SyncPoint::GetInstance()->ClearAllCallBacks(); + }); SyncPoint::GetInstance()->SetCallBack( "DBImpl::WriteImpl:AfterWriteWAL", [&](void*) { wal_written.store(true, std::memory_order_release); - while (!allow_publish.load(std::memory_order_acquire)) { + const auto deadline = + std::chrono::steady_clock::now() + std::chrono::seconds(10); + while (!allow_publish.load(std::memory_order_acquire) && + std::chrono::steady_clock::now() < deadline) { std::this_thread::yield(); } + publish_wait_timed_out.store( + !allow_publish.load(std::memory_order_acquire), + std::memory_order_release); }); SyncPoint::GetInstance()->EnableProcessing(); Status write_status; - std::thread writer([&]() { write_status = Put("key2", "value2"); }); - while (!wal_written.load(std::memory_order_acquire)) { + writer = std::thread([&]() { write_status = Put("key2", "value2"); }); + const auto wal_write_deadline = + std::chrono::steady_clock::now() + std::chrono::seconds(10); + while (!wal_written.load(std::memory_order_acquire) && + std::chrono::steady_clock::now() < wal_write_deadline) { std::this_thread::yield(); } + ASSERT_TRUE(wal_written.load(std::memory_order_acquire)); iter->Next(); const bool exposed_before_publish = iter->Valid(); @@ -496,10 +523,9 @@ TEST_F(DBTestXactLogIterator, allow_publish.store(true, std::memory_order_release); writer.join(); - SyncPoint::GetInstance()->DisableProcessing(); - SyncPoint::GetInstance()->ClearAllCallBacks(); ASSERT_OK(write_status); + ASSERT_FALSE(publish_wait_timed_out.load(std::memory_order_acquire)); ASSERT_FALSE(exposed_before_publish); ASSERT_OK(status_before_publish); ASSERT_EQ(dbfull()->GetLatestSequenceNumber(), 2U); diff --git a/db/transaction_log_impl.cc b/db/transaction_log_impl.cc index af5b6f3ac366..e03c2a59ae0c 100644 --- a/db/transaction_log_impl.cc +++ b/db/transaction_log_impl.cc @@ -107,7 +107,8 @@ bool TransactionLogIteratorImpl::IsRecordPublished(const Slice& record) const { return true; } - const uint64_t count = DecodeFixed32(record.data() + 8); + const uint64_t count = + DecodeFixed32(record.data() + sizeof(SequenceNumber)); if (count == 0) { return true; } @@ -139,7 +140,8 @@ bool TransactionLogIteratorImpl::RestrictedRead(Slice* record) { uint64_t record_checksum = 0; if (!current_log_reader_->ReadRecord( record, &scratch_, WALRecoveryMode::kTolerateCorruptedTailRecords, - &record_checksum)) { + read_options_.include_wal_record_checksum_ ? &record_checksum + : nullptr)) { return false; } diff --git a/include/rocksdb/transaction_log.h b/include/rocksdb/transaction_log.h index 8741af340455..dda582347ccc 100644 --- a/include/rocksdb/transaction_log.h +++ b/include/rocksdb/transaction_log.h @@ -73,7 +73,8 @@ struct BatchResult { // The first physical byte offset after this logical WAL record. uint64_t wal_record_end = 0; - // XXH3 checksum of the logical WAL record contents. + // XXH3 checksum of the logical WAL record contents, or zero when + // ReadOptions::include_wal_record_checksum_ is false. uint64_t wal_record_checksum = 0; std::unique_ptr writeBatchPtr; @@ -145,10 +146,16 @@ class TransactionLogIterator { // Default: true bool verify_checksums_; - ReadOptions() : verify_checksums_(true) {} + // If true, calculate and populate BatchResult::wal_record_checksum. + // Default: false + bool include_wal_record_checksum_; + + ReadOptions() + : verify_checksums_(true), include_wal_record_checksum_(false) {} explicit ReadOptions(bool verify_checksums) - : verify_checksums_(verify_checksums) {} + : verify_checksums_(verify_checksums), + include_wal_record_checksum_(false) {} }; }; } // namespace ROCKSDB_NAMESPACE From 4e743e4b7ccd332438b1d32ee2cb6c81155a3f83 Mon Sep 17 00:00:00 2001 From: Robert Nagy Date: Mon, 20 Jul 2026 12:50:20 +0200 Subject: [PATCH 3/3] Preserve transaction log iterator ABI --- db/db_log_iter_test.cc | 154 ++++++++++++++++++++---------- db/transaction_log_impl.cc | 54 +++++++---- db/transaction_log_impl.h | 6 +- include/rocksdb/transaction_log.h | 69 ++++++------- 4 files changed, 174 insertions(+), 109 deletions(-) diff --git a/db/db_log_iter_test.cc b/db/db_log_iter_test.cc index dd13b4a52d9e..ac4907516c25 100644 --- a/db/db_log_iter_test.cc +++ b/db/db_log_iter_test.cc @@ -16,10 +16,12 @@ #include #include "db/db_test_util.h" +#include "db/write_batch_internal.h" #include "env/mock_env.h" #include "port/stack_trace.h" #include "util/atomic.h" #include "util/defer.h" +#include "util/xxhash.h" namespace ROCKSDB_NAMESPACE { @@ -29,15 +31,20 @@ class DBTestXactLogIterator : public DBTestBase { : DBTestBase("db_log_iter_test", /*env_do_fsync=*/true) {} std::unique_ptr OpenTransactionLogIter( - const SequenceNumber seq, const bool include_wal_record_checksum = false) { + const SequenceNumber seq) { std::unique_ptr iter; - TransactionLogIterator::ReadOptions read_options; - read_options.include_wal_record_checksum_ = include_wal_record_checksum; - Status status = dbfull()->GetUpdatesSince(seq, &iter, read_options); + Status status = dbfull()->GetUpdatesSince(seq, &iter); EXPECT_OK(status); EXPECT_TRUE(iter->Valid()); return iter; } + + TransactionLogPositionV1 CurrentPosition(TransactionLogIterator* iter) { + TransactionLogPositionV1 position; + EXPECT_OK(GetTransactionLogIteratorBatchPositionV1( + iter, &position, /*include_checksum=*/true)); + return position; + } }; namespace { @@ -343,30 +350,65 @@ TEST_F(DBTestXactLogIterator, TransactionLogIteratorBlobs) { TEST_F(DBTestXactLogIterator, TransactionLogIteratorLogDataOnlyTailAndWalPosition) { Options options = OptionsForLogIterTest(); +#ifdef ZSTD + options.wal_compression = kZSTD; +#endif DestroyAndReopen(options); ASSERT_OK(Put("key1", "value1")); - auto default_iter = OpenTransactionLogIter(0); - ASSERT_EQ(default_iter->GetBatch().wal_record_checksum, 0U); - - auto iter = OpenTransactionLogIter(0, true); + auto iter = OpenTransactionLogIter(0); + const TransactionLogPositionV1 initial_position = + CurrentPosition(iter.get()); + TransactionLogPositionV1 repeated_initial_position; + ASSERT_OK(GetTransactionLogIteratorBatchPositionV1( + iter.get(), &repeated_initial_position, /*include_checksum=*/false)); + ASSERT_EQ(repeated_initial_position.wal_file_number, + initial_position.wal_file_number); + ASSERT_EQ(repeated_initial_position.wal_record_offset, + initial_position.wal_record_offset); + ASSERT_EQ(repeated_initial_position.wal_record_end, + initial_position.wal_record_end); + ASSERT_EQ(repeated_initial_position.wal_record_checksum, 0U); + ASSERT_TRUE(GetTransactionLogIteratorBatchPositionV1( + nullptr, &repeated_initial_position, + /*include_checksum=*/true) + .IsInvalidArgument()); + ASSERT_TRUE(GetTransactionLogIteratorBatchPositionV1( + iter.get(), nullptr, /*include_checksum=*/true) + .IsInvalidArgument()); BatchResult initial = iter->GetBatch(); ASSERT_EQ(initial.writeBatchPtr->Count(), 1U); + const Slice initial_contents = + WriteBatchInternal::Contents(initial.writeBatchPtr.get()); + ASSERT_EQ(initial_position.wal_record_checksum, + XXH3_64bits(initial_contents.data(), initial_contents.size())); + TransactionLogPositionV1 moved_out_position; + ASSERT_TRUE(GetTransactionLogIteratorBatchPositionV1( + iter.get(), &moved_out_position, + /*include_checksum=*/true) + .IsInvalidArgument()); iter->Next(); ASSERT_FALSE(iter->Valid()); ASSERT_OK(iter->status()); + Random random(301); + const std::string first_log_payload = random.RandomString(128 * 1024); WriteBatch first_log_data; - ASSERT_OK(first_log_data.PutLogData(Slice("log-data-1"))); + ASSERT_OK(first_log_data.PutLogData(Slice(first_log_payload))); ASSERT_OK(dbfull()->Write(WriteOptions(), &first_log_data)); iter->Next(); ASSERT_TRUE(iter->Valid()); + const TransactionLogPositionV1 first_position = CurrentPosition(iter.get()); BatchResult first = iter->GetBatch(); ASSERT_EQ(first.sequence, 2U); ASSERT_EQ(first.writeBatchPtr->Count(), 0U); - ASSERT_GT(first.wal_file_number, 0U); - ASSERT_LT(first.wal_record_offset, first.wal_record_end); + const Slice first_contents = + WriteBatchInternal::Contents(first.writeBatchPtr.get()); + ASSERT_EQ(first_position.wal_record_checksum, + XXH3_64bits(first_contents.data(), first_contents.size())); + ASSERT_GT(first_position.wal_file_number, 0U); + ASSERT_LT(first_position.wal_record_offset, first_position.wal_record_end); iter->Next(); ASSERT_FALSE(iter->Valid()); @@ -378,58 +420,65 @@ TEST_F(DBTestXactLogIterator, iter->Next(); ASSERT_TRUE(iter->Valid()); + const TransactionLogPositionV1 second_position = + CurrentPosition(iter.get()); BatchResult second = iter->GetBatch(); ASSERT_EQ(second.sequence, 2U); ASSERT_EQ(second.writeBatchPtr->Count(), 0U); - ASSERT_EQ(first.wal_file_number, second.wal_file_number); - ASSERT_LT(first.wal_record_offset, second.wal_record_offset); - ASSERT_LE(first.wal_record_end, second.wal_record_offset); - ASSERT_NE(first.wal_record_checksum, second.wal_record_checksum); - - const uint64_t first_file_number = first.wal_file_number; - const uint64_t first_offset = first.wal_record_offset; - const uint64_t first_end = first.wal_record_end; - const uint64_t first_checksum = first.wal_record_checksum; - const uint64_t second_file_number = second.wal_file_number; - const uint64_t second_offset = second.wal_record_offset; - const uint64_t second_end = second.wal_record_end; - const uint64_t second_checksum = second.wal_record_checksum; - - BatchResult move_constructed(std::move(first)); - ASSERT_EQ(move_constructed.wal_file_number, first_file_number); - ASSERT_EQ(move_constructed.wal_record_offset, first_offset); - ASSERT_EQ(move_constructed.wal_record_end, first_end); - ASSERT_EQ(move_constructed.wal_record_checksum, first_checksum); - - BatchResult moved; - moved = std::move(second); - ASSERT_EQ(moved.wal_file_number, second_file_number); - ASSERT_EQ(moved.wal_record_offset, second_offset); - ASSERT_EQ(moved.wal_record_end, second_end); - ASSERT_EQ(moved.wal_record_checksum, second_checksum); + ASSERT_EQ(first_position.wal_file_number, second_position.wal_file_number); + ASSERT_LT(first_position.wal_record_offset, + second_position.wal_record_offset); + ASSERT_LE(first_position.wal_record_end, + second_position.wal_record_offset); + ASSERT_NE(first_position.wal_record_checksum, + second_position.wal_record_checksum); iter->Next(); ASSERT_FALSE(iter->Valid()); ASSERT_OK(iter->status()); ASSERT_EQ(dbfull()->GetLatestSequenceNumber(), 1U); - auto replay = OpenTransactionLogIter(0, true); + auto replay = OpenTransactionLogIter(0); + const TransactionLogPositionV1 replay_initial_position = + CurrentPosition(replay.get()); + ASSERT_EQ(replay_initial_position.wal_file_number, + initial_position.wal_file_number); + ASSERT_EQ(replay_initial_position.wal_record_offset, + initial_position.wal_record_offset); + ASSERT_EQ(replay_initial_position.wal_record_end, + initial_position.wal_record_end); + ASSERT_EQ(replay_initial_position.wal_record_checksum, + initial_position.wal_record_checksum); ASSERT_EQ(replay->GetBatch().writeBatchPtr->Count(), 1U); replay->Next(); ASSERT_TRUE(replay->Valid()); + const TransactionLogPositionV1 replay_first_position = + CurrentPosition(replay.get()); BatchResult replay_first = replay->GetBatch(); - ASSERT_EQ(replay_first.wal_file_number, first_file_number); - ASSERT_EQ(replay_first.wal_record_offset, first_offset); - ASSERT_EQ(replay_first.wal_record_end, first_end); - ASSERT_EQ(replay_first.wal_record_checksum, first_checksum); + ASSERT_EQ(replay_first_position.wal_file_number, + first_position.wal_file_number); + ASSERT_EQ(replay_first_position.wal_record_offset, + first_position.wal_record_offset); + ASSERT_EQ(replay_first_position.wal_record_end, + first_position.wal_record_end); + ASSERT_EQ(replay_first_position.wal_record_checksum, + first_position.wal_record_checksum); + ASSERT_EQ(replay_first.sequence, first.sequence); replay->Next(); ASSERT_TRUE(replay->Valid()); + const TransactionLogPositionV1 replay_second_position = + CurrentPosition(replay.get()); BatchResult replay_second = replay->GetBatch(); - ASSERT_EQ(replay_second.wal_file_number, second_file_number); - ASSERT_EQ(replay_second.wal_record_offset, second_offset); - ASSERT_EQ(replay_second.wal_record_end, second_end); - ASSERT_EQ(replay_second.wal_record_checksum, second_checksum); + ASSERT_EQ(replay_second_position.wal_file_number, + second_position.wal_file_number); + ASSERT_EQ(replay_second_position.wal_record_offset, + second_position.wal_record_offset); + ASSERT_EQ(replay_second_position.wal_record_end, + second_position.wal_record_end); + ASSERT_EQ(replay_second_position.wal_record_checksum, + second_position.wal_record_checksum); + ASSERT_EQ(replay_second.sequence, second.sequence); replay->Next(); ASSERT_FALSE(replay->Valid()); ASSERT_OK(replay->status()); @@ -448,19 +497,24 @@ TEST_F(DBTestXactLogIterator, ASSERT_OK(dbfull()->Write(WriteOptions(), &second_log_data)); ASSERT_EQ(dbfull()->GetLatestSequenceNumber(), 0U); - auto iter = OpenTransactionLogIter(0, true); + auto iter = OpenTransactionLogIter(0); + const TransactionLogPositionV1 first_position = CurrentPosition(iter.get()); BatchResult first = iter->GetBatch(); ASSERT_EQ(first.sequence, 1U); ASSERT_EQ(first.writeBatchPtr->Count(), 0U); iter->Next(); ASSERT_TRUE(iter->Valid()); + const TransactionLogPositionV1 second_position = + CurrentPosition(iter.get()); BatchResult second = iter->GetBatch(); ASSERT_EQ(second.sequence, 1U); ASSERT_EQ(second.writeBatchPtr->Count(), 0U); - ASSERT_EQ(first.wal_file_number, second.wal_file_number); - ASSERT_LT(first.wal_record_offset, second.wal_record_offset); - ASSERT_NE(first.wal_record_checksum, second.wal_record_checksum); + ASSERT_EQ(first_position.wal_file_number, second_position.wal_file_number); + ASSERT_LT(first_position.wal_record_offset, + second_position.wal_record_offset); + ASSERT_NE(first_position.wal_record_checksum, + second_position.wal_record_checksum); iter->Next(); ASSERT_FALSE(iter->Valid()); diff --git a/db/transaction_log_impl.cc b/db/transaction_log_impl.cc index e03c2a59ae0c..ca752fb47554 100644 --- a/db/transaction_log_impl.cc +++ b/db/transaction_log_impl.cc @@ -11,6 +11,7 @@ #include "file/sequence_file_reader.h" #include "util/coding.h" #include "util/defer.h" +#include "util/xxhash.h" namespace ROCKSDB_NAMESPACE { @@ -37,16 +38,13 @@ TransactionLogIteratorImpl::TransactionLogIteratorImpl( current_batch_wal_file_number_(0), current_batch_wal_record_offset_(0), current_batch_wal_record_end_(0), - current_batch_wal_record_checksum_(0), current_record_wal_file_number_(0), current_record_wal_record_offset_(0), current_record_wal_record_end_(0), - current_record_wal_record_checksum_(0), has_unpublished_record_(false), unpublished_record_wal_file_number_(0), unpublished_record_wal_record_offset_(0), - unpublished_record_wal_record_end_(0), - unpublished_record_wal_record_checksum_(0) { + unpublished_record_wal_record_end_(0) { assert(files_ != nullptr); assert(versions_ != nullptr); assert(!seq_per_batch_); @@ -89,14 +87,42 @@ BatchResult TransactionLogIteratorImpl::GetBatch() { assert(is_valid_); // cannot call in a non valid state. BatchResult result; result.sequence = current_batch_seq_; - result.wal_file_number = current_batch_wal_file_number_; - result.wal_record_offset = current_batch_wal_record_offset_; - result.wal_record_end = current_batch_wal_record_end_; - result.wal_record_checksum = current_batch_wal_record_checksum_; result.writeBatchPtr = std::move(current_batch_); return result; } +Status TransactionLogIteratorImpl::GetBatchPosition( + TransactionLogPositionV1* position, bool include_checksum) const { + if (position == nullptr) { + return Status::InvalidArgument("Transaction log position is null"); + } + if (!is_valid_ || !current_status_.ok() || current_batch_ == nullptr) { + return Status::InvalidArgument( + "Transaction log iterator is not positioned at a batch"); + } + + position->wal_file_number = current_batch_wal_file_number_; + position->wal_record_offset = current_batch_wal_record_offset_; + position->wal_record_end = current_batch_wal_record_end_; + position->wal_record_checksum = 0; + if (include_checksum) { + const Slice contents = WriteBatchInternal::Contents(current_batch_.get()); + position->wal_record_checksum = + XXH3_64bits(contents.data(), contents.size()); + } + return Status::OK(); +} + +Status GetTransactionLogIteratorBatchPositionV1( + const TransactionLogIterator* iterator, + TransactionLogPositionV1* position, bool include_checksum) { + if (iterator == nullptr) { + return Status::InvalidArgument("Transaction log iterator is null"); + } + return static_cast(iterator) + ->GetBatchPosition(position, include_checksum); +} + Status TransactionLogIteratorImpl::status() { return current_status_; } bool TransactionLogIteratorImpl::Valid() { return started_ && is_valid_; } @@ -132,23 +158,18 @@ bool TransactionLogIteratorImpl::RestrictedRead(Slice* record) { current_record_wal_file_number_ = unpublished_record_wal_file_number_; current_record_wal_record_offset_ = unpublished_record_wal_record_offset_; current_record_wal_record_end_ = unpublished_record_wal_record_end_; - current_record_wal_record_checksum_ = - unpublished_record_wal_record_checksum_; return true; } - uint64_t record_checksum = 0; if (!current_log_reader_->ReadRecord( - record, &scratch_, WALRecoveryMode::kTolerateCorruptedTailRecords, - read_options_.include_wal_record_checksum_ ? &record_checksum - : nullptr)) { + record, &scratch_, + WALRecoveryMode::kTolerateCorruptedTailRecords)) { return false; } current_record_wal_file_number_ = current_log_reader_->GetLogNumber(); current_record_wal_record_offset_ = current_log_reader_->LastRecordOffset(); current_record_wal_record_end_ = current_log_reader_->LastRecordEnd(); - current_record_wal_record_checksum_ = record_checksum; if (!IsRecordPublished(*record)) { unpublished_record_.assign(record->data(), record->size()); @@ -156,8 +177,6 @@ bool TransactionLogIteratorImpl::RestrictedRead(Slice* record) { unpublished_record_wal_file_number_ = current_record_wal_file_number_; unpublished_record_wal_record_offset_ = current_record_wal_record_offset_; unpublished_record_wal_record_end_ = current_record_wal_record_end_; - unpublished_record_wal_record_checksum_ = - current_record_wal_record_checksum_; return false; } @@ -368,7 +387,6 @@ void TransactionLogIteratorImpl::UpdateCurrentWriteBatch(const Slice& record) { current_batch_wal_file_number_ = current_record_wal_file_number_; current_batch_wal_record_offset_ = current_record_wal_record_offset_; current_batch_wal_record_end_ = current_record_wal_record_end_; - current_batch_wal_record_checksum_ = current_record_wal_record_checksum_; is_valid_ = true; current_status_ = Status::OK(); } diff --git a/db/transaction_log_impl.h b/db/transaction_log_impl.h index ebe7c014ecd0..d5406ba2abbb 100644 --- a/db/transaction_log_impl.h +++ b/db/transaction_log_impl.h @@ -71,6 +71,9 @@ class TransactionLogIteratorImpl : public TransactionLogIterator { BatchResult GetBatch() override; + Status GetBatchPosition(TransactionLogPositionV1* position, + bool include_checksum) const; + private: const std::string& dir_; const ImmutableDBOptions* options_; @@ -112,13 +115,11 @@ class TransactionLogIteratorImpl : public TransactionLogIterator { uint64_t current_batch_wal_file_number_; uint64_t current_batch_wal_record_offset_; uint64_t current_batch_wal_record_end_; - uint64_t current_batch_wal_record_checksum_; // Metadata for the record most recently returned by RestrictedRead(). uint64_t current_record_wal_file_number_; uint64_t current_record_wal_record_offset_; uint64_t current_record_wal_record_end_; - uint64_t current_record_wal_record_checksum_; // A complete consuming record can reach the WAL before its ending sequence // is published. Keep it here so a later Next() can return it without either @@ -128,7 +129,6 @@ class TransactionLogIteratorImpl : public TransactionLogIterator { uint64_t unpublished_record_wal_file_number_; uint64_t unpublished_record_wal_record_offset_; uint64_t unpublished_record_wal_record_end_; - uint64_t unpublished_record_wal_record_checksum_; // Reads the next complete transaction-log record. Zero-count records are // immediately visible. Consuming records are returned only after their diff --git a/include/rocksdb/transaction_log.h b/include/rocksdb/transaction_log.h index dda582347ccc..6b5085c00c28 100644 --- a/include/rocksdb/transaction_log.h +++ b/include/rocksdb/transaction_log.h @@ -62,21 +62,6 @@ using LogFile = WalFile; struct BatchResult { SequenceNumber sequence = 0; - - // The WAL file number containing this batch. Together with - // wal_record_offset, this identifies the batch while the WAL is retained. - uint64_t wal_file_number = 0; - - // The physical byte offset of the first fragment of this logical WAL record. - uint64_t wal_record_offset = 0; - - // The first physical byte offset after this logical WAL record. - uint64_t wal_record_end = 0; - - // XXH3 checksum of the logical WAL record contents, or zero when - // ReadOptions::include_wal_record_checksum_ is false. - uint64_t wal_record_checksum = 0; - std::unique_ptr writeBatchPtr; // Add empty __ctor and __dtor for the rule of five @@ -90,28 +75,30 @@ struct BatchResult { BatchResult& operator=(const BatchResult&) = delete; - BatchResult(BatchResult&& bResult) noexcept - : sequence(bResult.sequence), - wal_file_number(bResult.wal_file_number), - wal_record_offset(bResult.wal_record_offset), - wal_record_end(bResult.wal_record_end), - wal_record_checksum(bResult.wal_record_checksum), + BatchResult(BatchResult&& bResult) + : sequence(std::move(bResult.sequence)), writeBatchPtr(std::move(bResult.writeBatchPtr)) {} - BatchResult& operator=(BatchResult&& bResult) noexcept { - if (this == &bResult) { - return *this; - } - sequence = bResult.sequence; - wal_file_number = bResult.wal_file_number; - wal_record_offset = bResult.wal_record_offset; - wal_record_end = bResult.wal_record_end; - wal_record_checksum = bResult.wal_record_checksum; + BatchResult& operator=(BatchResult&& bResult) { + sequence = std::move(bResult.sequence); writeBatchPtr = std::move(bResult.writeBatchPtr); return *this; } }; +// Exact physical position of a WriteBatch returned by a transaction log +// iterator. The WAL file number and record offset identify the batch while its +// WAL is retained. wal_record_end is the first physical byte after the logical +// WAL record. wal_record_checksum is the XXH3 checksum of the logical record +// contents when requested by GetTransactionLogIteratorBatchPositionV1(). +struct TransactionLogPositionV1 { + uint64_t wal_file_number = 0; + uint64_t wal_record_offset = 0; + uint64_t wal_record_end = 0; + uint64_t wal_record_checksum = 0; +}; +static_assert(sizeof(TransactionLogPositionV1) == 4 * sizeof(uint64_t)); + // A TransactionLogIterator is used to iterate over the transactions in a db. // One run of the iterator is continuous, i.e. the iterator will stop at the // beginning of any gap in sequences @@ -146,16 +133,22 @@ class TransactionLogIterator { // Default: true bool verify_checksums_; - // If true, calculate and populate BatchResult::wal_record_checksum. - // Default: false - bool include_wal_record_checksum_; - - ReadOptions() - : verify_checksums_(true), include_wal_record_checksum_(false) {} + ReadOptions() : verify_checksums_(true) {} explicit ReadOptions(bool verify_checksums) - : verify_checksums_(verify_checksums), - include_wal_record_checksum_(false) {} + : verify_checksums_(verify_checksums) {} }; }; + +// Returns the exact WAL position for the iterator's current batch without +// consuming it. wal_record_checksum is zero unless include_checksum is true. +// +// This extension supports iterators created by RocksDB's built-in DB +// implementation. Passing a custom TransactionLogIterator is unsupported. +// +// REQUIRES: Valid() is true, status() is OK, position is not null, and this is +// called before GetBatch(), which moves the current WriteBatch out. +Status GetTransactionLogIteratorBatchPositionV1( + const TransactionLogIterator* iterator, + TransactionLogPositionV1* position, bool include_checksum); } // namespace ROCKSDB_NAMESPACE