Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions db/db_impl/db_impl_write.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
262 changes: 262 additions & 0 deletions db/db_log_iter_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,17 @@
// in Release build.
// which is a pity, it is a good test

#include <atomic>
#include <chrono>
#include <thread>

#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 {

Expand All @@ -31,6 +38,13 @@ class DBTestXactLogIterator : public DBTestBase {
EXPECT_TRUE(iter->Valid());
return iter;
}

TransactionLogPositionV1 CurrentPosition(TransactionLogIterator* iter) {
TransactionLogPositionV1 position;
EXPECT_OK(GetTransactionLogIteratorBatchPositionV1(
iter, &position, /*include_checksum=*/true));
return position;
}
};

namespace {
Expand Down Expand Up @@ -332,6 +346,254 @@ TEST_F(DBTestXactLogIterator, TransactionLogIteratorBlobs) {
"Delete(0, key2)",
handler.seen);
}

TEST_F(DBTestXactLogIterator,
TransactionLogIteratorLogDataOnlyTailAndWalPosition) {
Options options = OptionsForLogIterTest();
#ifdef ZSTD
options.wal_compression = kZSTD;
#endif
DestroyAndReopen(options);
ASSERT_OK(Put("key1", "value1"));

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(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);
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());
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());
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_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);
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_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_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());
}

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);
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_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());
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<bool> wal_written{false};
std::atomic<bool> allow_publish{false};
std::atomic<bool> 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);
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;
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();
const Status status_before_publish = iter->status();

allow_publish.store(true, std::memory_order_release);
writer.join();

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);

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) {
Expand Down
Loading