| // Licensed to the Apache Software Foundation (ASF) under one |
| // or more contributor license agreements. See the NOTICE file |
| // distributed with this work for additional information |
| // regarding copyright ownership. The ASF licenses this file |
| // to you under the Apache License, Version 2.0 (the |
| // "License"); you may not use this file except in compliance |
| // with the License. You may obtain a copy of the License at |
| // |
| // http://www.apache.org/licenses/LICENSE-2.0 |
| // |
| // Unless required by applicable law or agreed to in writing, |
| // software distributed under the License is distributed on an |
| // "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
| // KIND, either express or implied. See the License for the |
| // specific language governing permissions and limitations |
| // under the License. |
| |
| #include "cloud/cloud_tablet.h" |
| |
| #include <gen_cpp/olap_file.pb.h> |
| #include <gtest/gtest-message.h> |
| #include <gtest/gtest-test-part.h> |
| #include <gtest/gtest.h> |
| |
| #include <chrono> |
| #include <cstdint> |
| #include <ranges> |
| |
| #include "cloud/cloud_meta_mgr.h" |
| #include "cloud/cloud_storage_engine.h" |
| #include "cloud/cloud_warm_up_manager.h" |
| #include "common/config.h" |
| #include "cpp/sync_point.h" |
| #include "storage/rowset/rowset.h" |
| #include "storage/rowset/rowset_factory.h" |
| #include "storage/rowset/rowset_meta.h" |
| #include "storage/tablet/tablet_meta.h" |
| #include "util/uid_util.h" |
| |
| namespace doris { |
| |
| using namespace std::chrono; |
| |
| class CloudTabletWarmUpStateTest : public testing::Test { |
| public: |
| CloudTabletWarmUpStateTest() : _engine(CloudStorageEngine(EngineOptions {})) {} |
| |
| void SetUp() override { |
| _tablet_meta.reset(new TabletMeta(1, 2, 15673, 15674, 4, 5, TTabletSchema(), 6, {{7, 8}}, |
| UniqueId(9, 10), TTabletType::TABLET_TYPE_DISK, |
| TCompressionType::LZ4F)); |
| _tablet = |
| std::make_shared<CloudTablet>(_engine, std::make_shared<TabletMeta>(*_tablet_meta)); |
| } |
| void TearDown() override {} |
| |
| RowsetSharedPtr create_rowset(Version version, int num_segments = 1) { |
| auto rs_meta = std::make_shared<RowsetMeta>(); |
| rs_meta->set_rowset_type(BETA_ROWSET); |
| rs_meta->set_version(version); |
| rs_meta->set_rowset_id(_engine.next_rowset_id()); |
| rs_meta->set_num_segments(num_segments); |
| RowsetSharedPtr rowset; |
| Status st = RowsetFactory::create_rowset(nullptr, "", rs_meta, &rowset); |
| if (!st.ok()) { |
| return nullptr; |
| } |
| return rowset; |
| } |
| |
| protected: |
| std::string _json_rowset_meta; |
| TabletMetaSharedPtr _tablet_meta; |
| std::shared_ptr<CloudTablet> _tablet; |
| CloudStorageEngine _engine; |
| }; |
| |
| // Test get_rowset_warmup_state for non-existent rowset |
| TEST_F(CloudTabletWarmUpStateTest, TestGetRowsetWarmupStateNonExistent) { |
| auto rowset = create_rowset(Version(1, 1)); |
| ASSERT_NE(rowset, nullptr); |
| |
| auto non_existent_id = _engine.next_rowset_id(); |
| |
| WarmUpState state = _tablet->get_rowset_warmup_state(non_existent_id); |
| WarmUpState expected_state = WarmUpState {WarmUpTriggerSource::NONE, WarmUpProgress::NONE}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test add_rowset_warmup_state with TRIGGERED_BY_EVENT_DRIVEN state |
| TEST_F(CloudTabletWarmUpStateTest, TestAddRowsetWarmupStateTriggeredByJob) { |
| auto rowset = create_rowset(Version(1, 1), 5); |
| ASSERT_NE(rowset, nullptr); |
| |
| bool result = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(result); |
| |
| // Verify the state is correctly set |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test add_rowset_warmup_state with TRIGGERED_BY_SYNC_ROWSET state |
| TEST_F(CloudTabletWarmUpStateTest, TestAddRowsetWarmupStateTriggeredBySyncRowset) { |
| auto rowset = create_rowset(Version(2, 2), 3); |
| ASSERT_NE(rowset, nullptr); |
| |
| bool result = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::SYNC_ROWSET); |
| EXPECT_TRUE(result); |
| |
| // Verify the state is correctly set |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test adding duplicate rowset warmup state should fail |
| TEST_F(CloudTabletWarmUpStateTest, TestAddDuplicateRowsetWarmupState) { |
| auto rowset = create_rowset(Version(3, 3), 2); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First addition should succeed |
| bool result1 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(result1); |
| |
| // Second addition should fail |
| bool result2 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::SYNC_ROWSET); |
| EXPECT_FALSE(result2); |
| |
| // State should remain the original one |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test complete_rowset_segment_warmup for non-existent rowset |
| TEST_F(CloudTabletWarmUpStateTest, TestCompleteRowsetSegmentWarmupNonExistent) { |
| auto non_existent_id = _engine.next_rowset_id(); |
| |
| WarmUpState result = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::SYNC_ROWSET, non_existent_id, Status::OK(), 1, 0); |
| WarmUpState expected_state = WarmUpState {WarmUpTriggerSource::NONE, WarmUpProgress::NONE}; |
| EXPECT_EQ(result, expected_state); |
| } |
| |
| // Test complete_rowset_segment_warmup with partial completion |
| TEST_F(CloudTabletWarmUpStateTest, TestCompleteRowsetSegmentWarmupPartial) { |
| auto rowset = create_rowset(Version(4, 4), 3); |
| ASSERT_NE(rowset, nullptr); |
| |
| // Add rowset warmup state |
| bool add_result = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(add_result); |
| |
| // Complete one segment, should still be in TRIGGERED_BY_EVENT_DRIVEN state |
| WarmUpState result1 = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), Status::OK(), 1, 0); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(result1, expected_state); |
| |
| // Complete second segment, should still be in TRIGGERED_BY_EVENT_DRIVEN state |
| WarmUpState result2 = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), Status::OK(), 1, 0); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(result2, expected_state); |
| |
| // Verify current state is still TRIGGERED_BY_EVENT_DRIVEN |
| WarmUpState current_state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(current_state, expected_state); |
| } |
| |
| // Test complete_rowset_segment_warmup with full completion |
| TEST_F(CloudTabletWarmUpStateTest, TestCompleteRowsetSegmentWarmupFull) { |
| auto rowset = create_rowset(Version(5, 5), 2); |
| ASSERT_NE(rowset, nullptr); |
| |
| // Add rowset warmup state |
| bool add_result = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::SYNC_ROWSET); |
| EXPECT_TRUE(add_result); |
| |
| // Complete first segment |
| WarmUpState result1 = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::SYNC_ROWSET, rowset->rowset_id(), Status::OK(), 1, 0); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(result1, expected_state); |
| |
| // Complete second segment, should transition to DONE state |
| WarmUpState result2 = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::SYNC_ROWSET, rowset->rowset_id(), Status::OK(), 1, 0); |
| expected_state = WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DONE}; |
| EXPECT_EQ(result2, expected_state); |
| |
| // Verify final state is DONE |
| WarmUpState final_state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DONE}; |
| EXPECT_EQ(final_state, expected_state); |
| } |
| |
| // Test complete_rowset_segment_warmup with inverted index file, partial completion |
| TEST_F(CloudTabletWarmUpStateTest, TestCompleteRowsetSegmentWarmupWithInvertedIndexPartial) { |
| auto rowset = create_rowset(Version(6, 6), 1); |
| ASSERT_NE(rowset, nullptr); |
| |
| // Add rowset warmup state |
| bool add_result = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(add_result); |
| |
| EXPECT_TRUE(_tablet->update_rowset_warmup_state_inverted_idx_num( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), 1)); |
| EXPECT_TRUE(_tablet->update_rowset_warmup_state_inverted_idx_num( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), 1)); |
| |
| // Complete one segment file |
| WarmUpState result1 = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), Status::OK(), 1, 0); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(result1, expected_state); |
| |
| // Complete inverted index file, should still be in TRIGGERED_BY_EVENT_DRIVEN state |
| WarmUpState result2 = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), Status::OK(), 0, 1); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(result2, expected_state); |
| |
| // Verify current state is still TRIGGERED_BY_EVENT_DRIVEN |
| WarmUpState current_state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(current_state, expected_state); |
| } |
| |
| // Test complete_rowset_segment_warmup with inverted index file, full completion |
| TEST_F(CloudTabletWarmUpStateTest, TestCompleteRowsetSegmentWarmupWithInvertedIndexFull) { |
| auto rowset = create_rowset(Version(6, 6), 1); |
| ASSERT_NE(rowset, nullptr); |
| |
| // Add rowset warmup state |
| bool add_result = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(add_result); |
| |
| EXPECT_TRUE(_tablet->update_rowset_warmup_state_inverted_idx_num( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), 1)); |
| |
| // Complete segment file |
| WarmUpState result1 = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), Status::OK(), 1, 0); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(result1, expected_state); |
| |
| // Complete inverted index file |
| WarmUpState result2 = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), Status::OK(), 0, 1); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(result2, expected_state); |
| |
| // Verify final state is DONE |
| WarmUpState final_state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(final_state, expected_state); |
| } |
| |
| // Test complete_rowset_segment_warmup with error status |
| TEST_F(CloudTabletWarmUpStateTest, TestCompleteRowsetSegmentWarmupWithError) { |
| auto rowset = create_rowset(Version(6, 6), 1); |
| ASSERT_NE(rowset, nullptr); |
| |
| // Add rowset warmup state |
| bool add_result = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(add_result); |
| |
| // Complete with error status, should still transition to DONE when all segments complete |
| Status error_status = Status::InternalError("Test error"); |
| WarmUpState result = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), error_status, 1, 0); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(result, expected_state); |
| |
| // Verify final state is DONE even with error |
| WarmUpState final_state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(final_state, expected_state); |
| } |
| |
| // Test multiple rowsets warmup state management |
| TEST_F(CloudTabletWarmUpStateTest, TestMultipleRowsetsWarmupState) { |
| auto rowset1 = create_rowset(Version(7, 7), 2); |
| auto rowset2 = create_rowset(Version(8, 8), 3); |
| auto rowset3 = create_rowset(Version(9, 9), 1); |
| ASSERT_NE(rowset1, nullptr); |
| ASSERT_NE(rowset2, nullptr); |
| ASSERT_NE(rowset3, nullptr); |
| |
| // Add multiple rowsets |
| EXPECT_TRUE(_tablet->add_rowset_warmup_state(*(rowset1->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN)); |
| EXPECT_TRUE(_tablet->add_rowset_warmup_state(*(rowset2->rowset_meta()), |
| WarmUpTriggerSource::SYNC_ROWSET)); |
| EXPECT_TRUE(_tablet->add_rowset_warmup_state(*(rowset3->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN)); |
| |
| // Verify all states |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(_tablet->get_rowset_warmup_state(rowset1->rowset_id()), expected_state); |
| expected_state = WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(_tablet->get_rowset_warmup_state(rowset2->rowset_id()), expected_state); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(_tablet->get_rowset_warmup_state(rowset3->rowset_id()), expected_state); |
| |
| // Complete rowset1 (2 segments) |
| WarmUpState result = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset1->rowset_id(), Status::OK(), 1, 0); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(result, expected_state); |
| result = _tablet->complete_rowset_segment_warmup(WarmUpTriggerSource::EVENT_DRIVEN, |
| rowset1->rowset_id(), Status::OK(), 1, 0); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(result, expected_state); |
| |
| // Complete rowset3 (1 segment) |
| result = _tablet->complete_rowset_segment_warmup(WarmUpTriggerSource::EVENT_DRIVEN, |
| rowset3->rowset_id(), Status::OK(), 1, 0); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(result, expected_state); |
| |
| // Verify states after completion |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(_tablet->get_rowset_warmup_state(rowset1->rowset_id()), expected_state); |
| expected_state = WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(_tablet->get_rowset_warmup_state(rowset2->rowset_id()), expected_state); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(_tablet->get_rowset_warmup_state(rowset3->rowset_id()), expected_state); |
| } |
| |
| // Test warmup state with zero segments (edge case) |
| TEST_F(CloudTabletWarmUpStateTest, TestWarmupStateWithZeroSegments) { |
| auto rowset = create_rowset(Version(10, 10), 0); |
| ASSERT_NE(rowset, nullptr); |
| |
| // Add rowset with zero segments |
| bool add_result = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(add_result); |
| |
| // State should be immediately ready for completion since there are no segments to warm up |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test concurrent access to warmup state (basic thread safety verification) |
| TEST_F(CloudTabletWarmUpStateTest, TestConcurrentWarmupStateAccess) { |
| auto rowset1 = create_rowset(Version(11, 11), 4); |
| auto rowset2 = create_rowset(Version(12, 12), 3); |
| ASSERT_NE(rowset1, nullptr); |
| ASSERT_NE(rowset2, nullptr); |
| |
| // Add rowsets from different "threads" (simulated by sequential calls) |
| EXPECT_TRUE(_tablet->add_rowset_warmup_state(*(rowset1->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN)); |
| EXPECT_TRUE(_tablet->add_rowset_warmup_state(*(rowset2->rowset_meta()), |
| WarmUpTriggerSource::SYNC_ROWSET)); |
| |
| // Interleaved completion operations |
| WarmUpState result = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset1->rowset_id(), Status::OK(), 1, 0); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(result, expected_state); |
| result = _tablet->complete_rowset_segment_warmup(WarmUpTriggerSource::SYNC_ROWSET, |
| rowset2->rowset_id(), Status::OK(), 1, 0); |
| expected_state = WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(result, expected_state); |
| result = _tablet->complete_rowset_segment_warmup(WarmUpTriggerSource::EVENT_DRIVEN, |
| rowset1->rowset_id(), Status::OK(), 1, 0); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(result, expected_state); |
| |
| // Check states are maintained correctly |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(_tablet->get_rowset_warmup_state(rowset1->rowset_id()), expected_state); |
| expected_state = WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(_tablet->get_rowset_warmup_state(rowset2->rowset_id()), expected_state); |
| } |
| |
| // Test only when the trigger source matches can the state be updated |
| TEST_F(CloudTabletWarmUpStateTest, TestCompleteRowsetSegmentWarmupTriggerSource) { |
| auto rowset = create_rowset(Version(13, 13), 1); |
| ASSERT_NE(rowset, nullptr); |
| // Add rowset warmup state |
| bool add_result = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_TRUE(add_result); |
| |
| // Attempt to update inverted index num with a different trigger source, should fail |
| bool update_result = _tablet->update_rowset_warmup_state_inverted_idx_num( |
| WarmUpTriggerSource::SYNC_ROWSET, rowset->rowset_id(), 1); |
| EXPECT_FALSE(update_result); |
| update_result = _tablet->update_rowset_warmup_state_inverted_idx_num( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), 1); |
| EXPECT_FALSE(update_result); |
| |
| // Attempt to complete with a different trigger source, should not update state |
| WarmUpState result = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::SYNC_ROWSET, rowset->rowset_id(), Status::OK(), 1, 0); |
| WarmUpState expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(result, expected_state); |
| result = _tablet->complete_rowset_segment_warmup(WarmUpTriggerSource::EVENT_DRIVEN, |
| rowset->rowset_id(), Status::OK(), 1, 0); |
| expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(result, expected_state); |
| |
| // Now complete with the correct trigger source |
| result = _tablet->complete_rowset_segment_warmup(WarmUpTriggerSource::JOB, rowset->rowset_id(), |
| Status::OK(), 1, 0); |
| expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DONE}; |
| EXPECT_EQ(result, expected_state); |
| } |
| |
| // Test JOB trigger source can override non-JOB warmup states |
| TEST_F(CloudTabletWarmUpStateTest, TestJobTriggerOverridesNonJobStates) { |
| auto rowset = create_rowset(Version(14, 14), 2); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First, add EVENT_DRIVEN warmup state |
| bool result1 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(result1); |
| |
| // Verify EVENT_DRIVEN state is set |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| |
| // Now add JOB warmup state, should override EVENT_DRIVEN |
| bool result2 = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_TRUE(result2); |
| |
| // Verify JOB state has overridden EVENT_DRIVEN |
| state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test JOB trigger source can override SYNC_ROWSET warmup state |
| TEST_F(CloudTabletWarmUpStateTest, TestJobTriggerOverridesSyncRowsetState) { |
| auto rowset = create_rowset(Version(15, 15), 3); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First, add SYNC_ROWSET warmup state |
| bool result1 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::SYNC_ROWSET); |
| EXPECT_TRUE(result1); |
| |
| // Verify SYNC_ROWSET state is set |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| |
| // Now add JOB warmup state, should override SYNC_ROWSET |
| bool result2 = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_TRUE(result2); |
| |
| // Verify JOB state has overridden SYNC_ROWSET |
| state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test JOB trigger source cannot override another JOB warmup state |
| TEST_F(CloudTabletWarmUpStateTest, TestJobTriggerCannotOverrideAnotherJob) { |
| auto rowset = create_rowset(Version(16, 16), 2); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First, add JOB warmup state |
| bool result1 = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_TRUE(result1); |
| |
| // Verify first JOB state is set |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| |
| // Try to add another JOB warmup state, should fail |
| bool result2 = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_FALSE(result2); |
| |
| // Verify state remains the original JOB state |
| state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test EVENT_DRIVEN cannot override existing JOB state |
| TEST_F(CloudTabletWarmUpStateTest, TestEventDrivenCannotOverrideJobState) { |
| auto rowset = create_rowset(Version(17, 17), 1); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First, add JOB warmup state |
| bool result1 = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_TRUE(result1); |
| |
| // Verify JOB state is set |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| |
| // Try to add EVENT_DRIVEN warmup state, should fail |
| bool result2 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_FALSE(result2); |
| |
| // Verify state remains JOB |
| state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test SYNC_ROWSET cannot override existing JOB state |
| TEST_F(CloudTabletWarmUpStateTest, TestSyncRowsetCannotOverrideJobState) { |
| auto rowset = create_rowset(Version(18, 18), 2); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First, add JOB warmup state |
| bool result1 = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_TRUE(result1); |
| |
| // Verify JOB state is set |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| |
| // Try to add SYNC_ROWSET warmup state, should fail |
| bool result2 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::SYNC_ROWSET); |
| EXPECT_FALSE(result2); |
| |
| // Verify state remains JOB |
| state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test EVENT_DRIVEN cannot override existing SYNC_ROWSET state |
| TEST_F(CloudTabletWarmUpStateTest, TestEventDrivenCannotOverrideSyncRowsetState) { |
| auto rowset = create_rowset(Version(19, 19), 1); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First, add SYNC_ROWSET warmup state |
| bool result1 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::SYNC_ROWSET); |
| EXPECT_TRUE(result1); |
| |
| // Verify SYNC_ROWSET state is set |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| |
| // Try to add EVENT_DRIVEN warmup state, should fail |
| bool result2 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_FALSE(result2); |
| |
| // Verify state remains SYNC_ROWSET |
| state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test SYNC_ROWSET cannot override existing EVENT_DRIVEN state |
| TEST_F(CloudTabletWarmUpStateTest, TestSyncRowsetCannotOverrideEventDrivenState) { |
| auto rowset = create_rowset(Version(20, 20), 3); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First, add EVENT_DRIVEN warmup state |
| bool result1 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(result1); |
| |
| // Verify EVENT_DRIVEN state is set |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| |
| // Try to add SYNC_ROWSET warmup state, should fail |
| bool result2 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::SYNC_ROWSET); |
| EXPECT_FALSE(result2); |
| |
| // Verify state remains EVENT_DRIVEN |
| state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test JOB can override DONE state from non-JOB source |
| TEST_F(CloudTabletWarmUpStateTest, TestJobCanOverrideDoneStateFromNonJob) { |
| auto rowset = create_rowset(Version(21, 21), 1); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First, add EVENT_DRIVEN warmup state |
| bool result1 = _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), |
| WarmUpTriggerSource::EVENT_DRIVEN); |
| EXPECT_TRUE(result1); |
| |
| // Complete the warmup to DONE state |
| WarmUpState result = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::EVENT_DRIVEN, rowset->rowset_id(), Status::OK(), 1, 0); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::EVENT_DRIVEN, WarmUpProgress::DONE}; |
| EXPECT_EQ(result, expected_state); |
| |
| // Now add JOB warmup state, should override the DONE state |
| bool result2 = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_TRUE(result2); |
| |
| // Verify JOB state has overridden the DONE state |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test JOB can override DONE state from JOB source |
| TEST_F(CloudTabletWarmUpStateTest, TestJobCanOverrideDoneStateFromJob) { |
| auto rowset = create_rowset(Version(21, 21), 1); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First, add EVENT_DRIVEN warmup state |
| bool result1 = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_TRUE(result1); |
| |
| // Complete the warmup to DONE state |
| WarmUpState result = _tablet->complete_rowset_segment_warmup( |
| WarmUpTriggerSource::JOB, rowset->rowset_id(), Status::OK(), 1, 0); |
| WarmUpState expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DONE}; |
| EXPECT_EQ(result, expected_state); |
| |
| // Now add JOB warmup state, should override the DONE state |
| bool result2 = |
| _tablet->add_rowset_warmup_state(*(rowset->rowset_meta()), WarmUpTriggerSource::JOB); |
| EXPECT_TRUE(result2); |
| |
| // Verify JOB state has overridden the DONE state |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| expected_state = WarmUpState {WarmUpTriggerSource::JOB, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test class for sync_meta functionality |
| class CloudTabletSyncMetaTest : public testing::Test { |
| public: |
| CloudTabletSyncMetaTest() : _engine(CloudStorageEngine(EngineOptions {})) {} |
| |
| void SetUp() override { |
| config::enable_file_cache = true; |
| |
| // Use incrementing tablet_id to create unique schema cache keys for each test |
| // This avoids cache pollution between tests |
| int64_t unique_tablet_id = 15673 + _test_counter++; |
| |
| // Create tablet meta with a schema that has disable_auto_compaction = false |
| TTabletSchema tablet_schema; |
| tablet_schema.__set_disable_auto_compaction(false); |
| // Add a unique column to ensure unique cache key |
| TColumn col; |
| col.__set_column_name("test_col_" + std::to_string(unique_tablet_id)); |
| col.__set_column_type(TColumnType()); |
| col.column_type.__set_type(TPrimitiveType::INT); |
| col.__set_is_key(true); |
| col.__set_aggregation_type(TAggregationType::NONE); |
| col.__set_col_unique_id(0); |
| tablet_schema.__set_columns({col}); |
| tablet_schema.__set_keys_type(TKeysType::DUP_KEYS); |
| |
| // Use column ordinal 0 -> unique_id 0 mapping |
| _tablet_meta.reset(new TabletMeta(1, 2, unique_tablet_id, 15674, 4, 5, tablet_schema, 1, |
| {{0, 0}}, UniqueId(9, 10), TTabletType::TABLET_TYPE_DISK, |
| TCompressionType::LZ4F)); |
| |
| _tablet = |
| std::make_shared<CloudTablet>(_engine, std::make_shared<TabletMeta>(*_tablet_meta)); |
| |
| _current_tablet_id = unique_tablet_id; |
| } |
| |
| void TearDown() override { config::enable_file_cache = false; } |
| |
| protected: |
| // Helper to create a unique TabletMeta for mock responses |
| TabletMetaSharedPtr createMockTabletMeta(bool disable_auto_compaction, |
| const std::string& compaction_policy = "size_based") { |
| TTabletSchema new_schema; |
| new_schema.__set_disable_auto_compaction(disable_auto_compaction); |
| TColumn col; |
| col.__set_column_name("test_col_" + std::to_string(_current_tablet_id)); |
| col.__set_column_type(TColumnType()); |
| col.column_type.__set_type(TPrimitiveType::INT); |
| col.__set_is_key(true); |
| col.__set_aggregation_type(TAggregationType::NONE); |
| col.__set_col_unique_id(0); |
| new_schema.__set_columns({col}); |
| new_schema.__set_keys_type(TKeysType::DUP_KEYS); |
| |
| TabletMetaSharedPtr meta; |
| meta.reset(new TabletMeta(1, 2, _current_tablet_id, 15674, 4, 5, new_schema, 1, {{0, 0}}, |
| UniqueId(9, 10), TTabletType::TABLET_TYPE_DISK, |
| TCompressionType::LZ4F, 0, false, std::nullopt, |
| compaction_policy)); |
| return meta; |
| } |
| |
| TabletMetaSharedPtr _tablet_meta; |
| std::shared_ptr<CloudTablet> _tablet; |
| CloudStorageEngine _engine; |
| int64_t _current_tablet_id; |
| static inline int _test_counter = 0; |
| }; |
| |
| // Test sync_meta syncs disable_auto_compaction from false to true |
| TEST_F(CloudTabletSyncMetaTest, TestSyncMetaDisableAutoCompactionFalseToTrue) { |
| // Verify initial state: disable_auto_compaction = false |
| EXPECT_FALSE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| |
| auto sp = SyncPoint::get_instance(); |
| sp->clear_all_call_backs(); |
| sp->enable_processing(); |
| |
| // Create mock tablet meta with disable_auto_compaction = true |
| auto mock_tablet_meta = createMockTabletMeta(true); |
| |
| // Mock get_tablet_meta to return tablet_meta with disable_auto_compaction = true |
| sp->set_call_back("CloudMetaMgr::get_tablet_meta", [mock_tablet_meta](auto&& args) { |
| auto* tablet_meta_ptr = try_any_cast<TabletMetaSharedPtr*>(args[1]); |
| *tablet_meta_ptr = mock_tablet_meta; |
| |
| // Tell the sync point to return with Status::OK() |
| try_any_cast_ret<Status>(args)->second = true; |
| }); |
| |
| // Call sync_meta |
| Status st = _tablet->sync_meta(); |
| EXPECT_TRUE(st.ok()); |
| |
| // Verify disable_auto_compaction has been synced to true |
| EXPECT_TRUE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| |
| sp->disable_processing(); |
| sp->clear_all_call_backs(); |
| } |
| |
| // Test sync_meta syncs disable_auto_compaction from true to false |
| TEST_F(CloudTabletSyncMetaTest, TestSyncMetaDisableAutoCompactionTrueToFalse) { |
| // Create a tablet with disable_auto_compaction = true from the start |
| // We need to create a fresh tablet with true to avoid cache pollution issues |
| TTabletSchema initial_schema; |
| initial_schema.__set_disable_auto_compaction(true); |
| TColumn col; |
| col.__set_column_name("test_col_true_to_false_" + std::to_string(_current_tablet_id)); |
| col.__set_column_type(TColumnType()); |
| col.column_type.__set_type(TPrimitiveType::INT); |
| col.__set_is_key(true); |
| col.__set_aggregation_type(TAggregationType::NONE); |
| col.__set_col_unique_id(0); |
| initial_schema.__set_columns({col}); |
| initial_schema.__set_keys_type(TKeysType::DUP_KEYS); |
| |
| TabletMetaSharedPtr tablet_meta_true; |
| tablet_meta_true.reset(new TabletMeta(1, 2, _current_tablet_id + 1000, 15674, 4, 5, |
| initial_schema, 1, {{0, 0}}, UniqueId(9, 10), |
| TTabletType::TABLET_TYPE_DISK, TCompressionType::LZ4F)); |
| _tablet = |
| std::make_shared<CloudTablet>(_engine, std::make_shared<TabletMeta>(*tablet_meta_true)); |
| |
| // Verify initial state: disable_auto_compaction = true |
| EXPECT_TRUE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| |
| auto sp = SyncPoint::get_instance(); |
| sp->clear_all_call_backs(); |
| sp->enable_processing(); |
| |
| // Create mock tablet meta with disable_auto_compaction = false (with matching column name) |
| TTabletSchema mock_schema; |
| mock_schema.__set_disable_auto_compaction(false); |
| TColumn mock_col; |
| mock_col.__set_column_name("test_col_true_to_false_" + std::to_string(_current_tablet_id)); |
| mock_col.__set_column_type(TColumnType()); |
| mock_col.column_type.__set_type(TPrimitiveType::INT); |
| mock_col.__set_is_key(true); |
| mock_col.__set_aggregation_type(TAggregationType::NONE); |
| mock_col.__set_col_unique_id(0); |
| mock_schema.__set_columns({mock_col}); |
| mock_schema.__set_keys_type(TKeysType::DUP_KEYS); |
| |
| TabletMetaSharedPtr mock_tablet_meta; |
| mock_tablet_meta.reset(new TabletMeta(1, 2, _current_tablet_id + 1000, 15674, 4, 5, mock_schema, |
| 1, {{0, 0}}, UniqueId(9, 10), |
| TTabletType::TABLET_TYPE_DISK, TCompressionType::LZ4F)); |
| |
| // Mock get_tablet_meta to return tablet_meta with disable_auto_compaction = false |
| sp->set_call_back("CloudMetaMgr::get_tablet_meta", [mock_tablet_meta](auto&& args) { |
| auto* tablet_meta_ptr = try_any_cast<TabletMetaSharedPtr*>(args[1]); |
| *tablet_meta_ptr = mock_tablet_meta; |
| |
| // Tell the sync point to return with Status::OK() |
| try_any_cast_ret<Status>(args)->second = true; |
| }); |
| |
| // Call sync_meta |
| Status st = _tablet->sync_meta(); |
| EXPECT_TRUE(st.ok()); |
| |
| // Verify disable_auto_compaction has been synced to false |
| EXPECT_FALSE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| |
| sp->disable_processing(); |
| sp->clear_all_call_backs(); |
| } |
| |
| // Test sync_meta when disable_auto_compaction is unchanged |
| TEST_F(CloudTabletSyncMetaTest, TestSyncMetaDisableAutoCompactionUnchanged) { |
| // Verify initial state: disable_auto_compaction = false |
| EXPECT_FALSE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| |
| auto sp = SyncPoint::get_instance(); |
| sp->clear_all_call_backs(); |
| sp->enable_processing(); |
| |
| // Create mock tablet meta with disable_auto_compaction = false (same as current) |
| auto mock_tablet_meta = createMockTabletMeta(false); |
| |
| // Mock get_tablet_meta to return tablet_meta with same disable_auto_compaction = false |
| sp->set_call_back("CloudMetaMgr::get_tablet_meta", [mock_tablet_meta](auto&& args) { |
| auto* tablet_meta_ptr = try_any_cast<TabletMetaSharedPtr*>(args[1]); |
| *tablet_meta_ptr = mock_tablet_meta; |
| |
| // Tell the sync point to return with Status::OK() |
| try_any_cast_ret<Status>(args)->second = true; |
| }); |
| |
| // Call sync_meta |
| Status st = _tablet->sync_meta(); |
| EXPECT_TRUE(st.ok()); |
| |
| // Verify disable_auto_compaction remains false |
| EXPECT_FALSE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| |
| sp->disable_processing(); |
| sp->clear_all_call_backs(); |
| } |
| |
| // Test sync_meta is skipped when enable_file_cache is false |
| TEST_F(CloudTabletSyncMetaTest, TestSyncMetaSkippedWhenFileCacheDisabled) { |
| // Disable file cache |
| config::enable_file_cache = false; |
| |
| // Set initial state: disable_auto_compaction = false |
| EXPECT_FALSE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| |
| auto sp = SyncPoint::get_instance(); |
| sp->clear_all_call_backs(); |
| sp->enable_processing(); |
| |
| bool callback_called = false; |
| sp->set_call_back("CloudMetaMgr::get_tablet_meta", |
| [&callback_called](auto&& args) { callback_called = true; }); |
| |
| // Call sync_meta - should return early without calling get_tablet_meta |
| Status st = _tablet->sync_meta(); |
| EXPECT_TRUE(st.ok()); |
| EXPECT_FALSE(callback_called); |
| |
| // Verify disable_auto_compaction is not changed |
| EXPECT_FALSE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| |
| sp->disable_processing(); |
| sp->clear_all_call_backs(); |
| } |
| |
| // Test sync_meta syncs compaction_policy together with disable_auto_compaction |
| TEST_F(CloudTabletSyncMetaTest, TestSyncMetaMultipleProperties) { |
| // Verify initial states |
| EXPECT_FALSE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| // Default compaction_policy is "size_based" |
| EXPECT_EQ(_tablet->tablet_meta()->compaction_policy(), "size_based"); |
| |
| auto sp = SyncPoint::get_instance(); |
| sp->clear_all_call_backs(); |
| sp->enable_processing(); |
| |
| // Create mock tablet meta with updated properties |
| auto mock_tablet_meta = createMockTabletMeta(true, "time_series"); |
| |
| // Mock get_tablet_meta to return tablet_meta with updated properties |
| sp->set_call_back("CloudMetaMgr::get_tablet_meta", [mock_tablet_meta](auto&& args) { |
| auto* tablet_meta_ptr = try_any_cast<TabletMetaSharedPtr*>(args[1]); |
| *tablet_meta_ptr = mock_tablet_meta; |
| |
| // Tell the sync point to return with Status::OK() |
| try_any_cast_ret<Status>(args)->second = true; |
| }); |
| |
| // Call sync_meta |
| Status st = _tablet->sync_meta(); |
| EXPECT_TRUE(st.ok()); |
| |
| // Verify both properties are synced |
| EXPECT_TRUE(_tablet->tablet_meta()->tablet_schema()->disable_auto_compaction()); |
| EXPECT_EQ(_tablet->tablet_meta()->compaction_policy(), "time_series"); |
| |
| sp->disable_processing(); |
| sp->clear_all_call_backs(); |
| } |
| class CloudTabletApplyVisiblePendingTest : public testing::Test { |
| public: |
| CloudTabletApplyVisiblePendingTest() : _engine(CloudStorageEngine(EngineOptions {})) {} |
| |
| void SetUp() override { |
| _tablet_meta.reset(new TabletMeta(1, 2, 15673, 15674, 4, 5, TTabletSchema(), 6, {{7, 8}}, |
| UniqueId(9, 10), TTabletType::TABLET_TYPE_DISK, |
| TCompressionType::LZ4F)); |
| _tablet = |
| std::make_shared<CloudTablet>(_engine, std::make_shared<TabletMeta>(*_tablet_meta)); |
| } |
| |
| void TearDown() override {} |
| |
| RowsetSharedPtr create_rowset(Version version, int num_segments = 1) { |
| auto rs_meta = std::make_shared<RowsetMeta>(); |
| rs_meta->set_rowset_type(BETA_ROWSET); |
| rs_meta->set_version(version); |
| rs_meta->set_rowset_id(_engine.next_rowset_id()); |
| rs_meta->set_num_segments(num_segments); |
| RowsetSharedPtr rowset; |
| Status st = RowsetFactory::create_rowset(nullptr, "", rs_meta, &rowset); |
| if (!st.ok()) { |
| return nullptr; |
| } |
| return rowset; |
| } |
| |
| RowsetMetaSharedPtr create_pending_rowset_meta(int64_t version) { |
| auto rs_meta = std::make_shared<RowsetMeta>(); |
| rs_meta->set_rowset_type(BETA_ROWSET); |
| rs_meta->set_version(Version(version, version)); |
| rs_meta->set_rowset_id(_engine.next_rowset_id()); |
| rs_meta->set_num_segments(1); |
| return rs_meta; |
| } |
| |
| // Create a rowset whose RowsetMeta carries a valid TabletSchema, |
| // required as template for create_empty_rowset_for_hole. |
| RowsetSharedPtr create_rowset_with_schema(Version version, int num_segments = 1) { |
| auto rs_meta = std::make_shared<RowsetMeta>(); |
| rs_meta->set_rowset_type(BETA_ROWSET); |
| rs_meta->set_version(version); |
| rs_meta->set_rowset_id(_engine.next_rowset_id()); |
| rs_meta->set_num_segments(num_segments); |
| |
| TabletSchemaPB schema_pb; |
| schema_pb.set_keys_type(KeysType::DUP_KEYS); |
| auto* col = schema_pb.add_column(); |
| col->set_unique_id(0); |
| col->set_name("k1"); |
| col->set_type("INT"); |
| col->set_is_key(true); |
| col->set_is_nullable(false); |
| rs_meta->set_tablet_schema(schema_pb); |
| |
| RowsetSharedPtr rowset; |
| Status st = RowsetFactory::create_rowset(nullptr, "", rs_meta, &rowset); |
| if (!st.ok()) { |
| return nullptr; |
| } |
| return rowset; |
| } |
| |
| void add_initial_rowsets(const std::vector<RowsetSharedPtr>& rowsets) { |
| std::unique_lock meta_wlock(_tablet->get_header_lock()); |
| _tablet->add_rowsets(std::vector<RowsetSharedPtr>(rowsets), false, meta_wlock, false); |
| } |
| |
| void add_pending_rowset(int64_t version, RowsetMetaSharedPtr rowset_meta, |
| int64_t expiration_time = INT64_MAX, bool is_empty = false) { |
| std::lock_guard<std::mutex> lock(_tablet->_visible_pending_rs_lock); |
| _tablet->_visible_pending_rs_map.emplace( |
| version, CloudTablet::VisiblePendingRowset {std::move(rowset_meta), expiration_time, |
| is_empty}); |
| } |
| |
| size_t pending_rs_count() const { |
| std::lock_guard<std::mutex> lock(_tablet->_visible_pending_rs_lock); |
| return _tablet->_visible_pending_rs_map.size(); |
| } |
| |
| protected: |
| TabletMetaSharedPtr _tablet_meta; |
| std::shared_ptr<CloudTablet> _tablet; |
| CloudStorageEngine _engine; |
| }; |
| |
| // Test apply with no pending rowsets does nothing |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyNoPendingRowsets) { |
| auto rs = create_rowset(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 1); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| } |
| |
| // Test apply single consecutive non-empty rowset |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplySingleConsecutiveRowset) { |
| auto rs = create_rowset(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| |
| add_pending_rowset(2, create_pending_rowset_meta(2)); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| EXPECT_EQ(_tablet->max_version_unlocked(), 2); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 2); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_TRUE(rowset_map.contains(Version(2, 2))); |
| } |
| |
| // Test apply multiple consecutive non-empty rowsets |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyMultipleConsecutiveRowsets) { |
| auto rs = create_rowset(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| |
| for (int64_t v = 2; v <= 4; ++v) { |
| add_pending_rowset(v, create_pending_rowset_meta(v)); |
| } |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| EXPECT_EQ(_tablet->max_version_unlocked(), 4); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 4); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_TRUE(rowset_map.contains(Version(2, 2))); |
| EXPECT_TRUE(rowset_map.contains(Version(3, 3))); |
| EXPECT_TRUE(rowset_map.contains(Version(4, 4))); |
| } |
| |
| // Test apply with version gap - nothing should be applied |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyWithVersionGap) { |
| auto rs = create_rowset(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| |
| // Add version 3 only, skip version 2 |
| add_pending_rowset(3, create_pending_rowset_meta(3)); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 1); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_FALSE(rowset_map.contains(Version(3, 3))); |
| } |
| |
| // Test apply with partial consecutive versions - only consecutive prefix applied |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyPartialConsecutive) { |
| auto rs = create_rowset(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| |
| // Add versions 2, 3, 5 (version 4 missing) |
| add_pending_rowset(2, create_pending_rowset_meta(2)); |
| add_pending_rowset(3, create_pending_rowset_meta(3)); |
| add_pending_rowset(5, create_pending_rowset_meta(5)); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| // Only versions 2 and 3 should be applied |
| EXPECT_EQ(_tablet->max_version_unlocked(), 3); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 3); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_TRUE(rowset_map.contains(Version(2, 2))); |
| EXPECT_TRUE(rowset_map.contains(Version(3, 3))); |
| EXPECT_FALSE(rowset_map.contains(Version(5, 5))); |
| } |
| |
| // Test apply with pending versions below max_version - nothing applied |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyPendingBelowMaxVersion) { |
| auto rs1 = create_rowset(Version(0, 1)); |
| auto rs2 = create_rowset(Version(2, 5)); |
| ASSERT_NE(rs1, nullptr); |
| ASSERT_NE(rs2, nullptr); |
| add_initial_rowsets({rs1, rs2}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 5); |
| |
| // Add pending versions 3 and 4, both below max_version |
| add_pending_rowset(3, create_pending_rowset_meta(3)); |
| add_pending_rowset(4, create_pending_rowset_meta(4)); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| EXPECT_EQ(_tablet->max_version_unlocked(), 5); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 2); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_TRUE(rowset_map.contains(Version(2, 5))); |
| EXPECT_FALSE(rowset_map.contains(Version(3, 3))); |
| EXPECT_FALSE(rowset_map.contains(Version(4, 4))); |
| } |
| |
| // Test apply with initial max_version = -1 (no initial rowsets) |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyWithNoInitialRowsets) { |
| EXPECT_EQ(_tablet->max_version_unlocked(), -1); |
| |
| add_pending_rowset(0, create_pending_rowset_meta(0)); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| EXPECT_EQ(_tablet->max_version_unlocked(), 0); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 1); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 0))); |
| } |
| |
| // Test apply called multiple times incrementally |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyMultipleCalls) { |
| auto rs = create_rowset(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| |
| // First apply: version 2 |
| add_pending_rowset(2, create_pending_rowset_meta(2)); |
| _tablet->apply_visible_pending_rowsets(); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 2); |
| EXPECT_TRUE(_tablet->rowset_map().contains(Version(2, 2))); |
| |
| // Second apply: version 3 |
| add_pending_rowset(3, create_pending_rowset_meta(3)); |
| _tablet->apply_visible_pending_rowsets(); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 3); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 3); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_TRUE(rowset_map.contains(Version(2, 2))); |
| EXPECT_TRUE(rowset_map.contains(Version(3, 3))); |
| } |
| |
| // Test gap resolved by later apply call |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyGapResolvedLater) { |
| auto rs = create_rowset(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| |
| // Add version 3 first (gap at version 2) |
| add_pending_rowset(3, create_pending_rowset_meta(3)); |
| _tablet->apply_visible_pending_rowsets(); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); // Nothing applied |
| EXPECT_FALSE(_tablet->rowset_map().contains(Version(3, 3))); |
| |
| // Now add version 2 to fill the gap |
| add_pending_rowset(2, create_pending_rowset_meta(2)); |
| _tablet->apply_visible_pending_rowsets(); |
| |
| // Both versions 2 and 3 should now be applied |
| EXPECT_EQ(_tablet->max_version_unlocked(), 3); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 3); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_TRUE(rowset_map.contains(Version(2, 2))); |
| EXPECT_TRUE(rowset_map.contains(Version(3, 3))); |
| } |
| |
| // Test clear_unused_visible_pending_rowsets removes applied entries |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestClearAfterApply) { |
| auto rs = create_rowset(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| |
| add_pending_rowset(2, create_pending_rowset_meta(2)); |
| add_pending_rowset(3, create_pending_rowset_meta(3)); |
| // Version 5 has a gap, won't be applied |
| add_pending_rowset(5, create_pending_rowset_meta(5)); |
| EXPECT_EQ(pending_rs_count(), 3); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| EXPECT_EQ(_tablet->max_version_unlocked(), 3); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 3); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_TRUE(rowset_map.contains(Version(2, 2))); |
| EXPECT_TRUE(rowset_map.contains(Version(3, 3))); |
| EXPECT_FALSE(rowset_map.contains(Version(5, 5))); |
| // Versions 2 and 3 are cleared (applied, version <= max_version) |
| // Version 5 remains (not applied, not expired) |
| EXPECT_EQ(pending_rs_count(), 1); |
| } |
| |
| // Test empty rowset with no existing versions breaks early |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyEmptyRowsetNoExistingVersions) { |
| EXPECT_EQ(_tablet->max_version_unlocked(), -1); |
| |
| // Add empty pending rowset at version 0 |
| add_pending_rowset(0, nullptr, INT64_MAX, true); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| // Cannot create empty rowset without a previous rowset as template |
| EXPECT_EQ(_tablet->max_version_unlocked(), -1); |
| EXPECT_EQ(_tablet->rowset_map().size(), 0); |
| } |
| |
| // Test empty rowset with existing version uses create_empty_rowset_for_hole |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyEmptyRowsetWithExistingVersion) { |
| auto rs = create_rowset_with_schema(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| |
| // Add empty pending rowset at version 2 |
| add_pending_rowset(2, nullptr, INT64_MAX, true); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| EXPECT_EQ(_tablet->max_version_unlocked(), 2); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 2); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_TRUE(rowset_map.contains(Version(2, 2))); |
| } |
| |
| // Test mixed non-empty followed by empty rowset |
| TEST_F(CloudTabletApplyVisiblePendingTest, TestApplyNonEmptyThenEmptyRowset) { |
| auto rs = create_rowset_with_schema(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| add_initial_rowsets({rs}); |
| EXPECT_EQ(_tablet->max_version_unlocked(), 1); |
| |
| // Version 2: non-empty (with schema for empty rowset template), Version 3: empty |
| auto pending_meta = create_pending_rowset_meta(2); |
| TabletSchemaPB schema_pb; |
| schema_pb.set_keys_type(KeysType::DUP_KEYS); |
| auto* col = schema_pb.add_column(); |
| col->set_unique_id(0); |
| col->set_name("k1"); |
| col->set_type("INT"); |
| col->set_is_key(true); |
| col->set_is_nullable(false); |
| pending_meta->set_tablet_schema(schema_pb); |
| add_pending_rowset(2, std::move(pending_meta)); |
| add_pending_rowset(3, nullptr, INT64_MAX, true); |
| |
| _tablet->apply_visible_pending_rowsets(); |
| |
| // Both should be applied; empty rowset uses to_add.back() as prev_rowset |
| EXPECT_EQ(_tablet->max_version_unlocked(), 3); |
| auto& rowset_map = _tablet->rowset_map(); |
| EXPECT_EQ(rowset_map.size(), 3); |
| EXPECT_TRUE(rowset_map.contains(Version(0, 1))); |
| EXPECT_TRUE(rowset_map.contains(Version(2, 2))); |
| EXPECT_TRUE(rowset_map.contains(Version(3, 3))); |
| } |
| |
| // Test is_rowset_warmed_up returns true for rowset NOT in the warmup state map |
| // This is the behavior change: missing warmup state → optimistically warmed up |
| TEST_F(CloudTabletWarmUpStateTest, TestIsRowsetWarmedUpMissingFromMap) { |
| auto rowset = create_rowset(Version(22, 22)); |
| ASSERT_NE(rowset, nullptr); |
| |
| // Rowset is not in the warmup state map at all |
| // Before the fix, this would return false. Now it returns true. |
| EXPECT_TRUE(_tablet->is_rowset_warmed_up(rowset->rowset_id())); |
| } |
| |
| // Test is_rowset_warmed_up returns true for rowset with DONE state |
| TEST_F(CloudTabletWarmUpStateTest, TestIsRowsetWarmedUpWithDoneState) { |
| auto rowset = create_rowset(Version(23, 23)); |
| ASSERT_NE(rowset, nullptr); |
| |
| _tablet->add_warmed_up_rowset(rowset->rowset_id()); |
| EXPECT_TRUE(_tablet->is_rowset_warmed_up(rowset->rowset_id())); |
| } |
| |
| // Test is_rowset_warmed_up returns false for rowset with DOING state (in map but not done) |
| TEST_F(CloudTabletWarmUpStateTest, TestIsRowsetWarmedUpWithDoingState) { |
| auto rowset = create_rowset(Version(24, 24)); |
| ASSERT_NE(rowset, nullptr); |
| |
| _tablet->add_not_warmed_up_rowset(rowset->rowset_id()); |
| EXPECT_FALSE(_tablet->is_rowset_warmed_up(rowset->rowset_id())); |
| } |
| |
| // Test add_not_warmed_up_rowset sets DOING state correctly |
| TEST_F(CloudTabletWarmUpStateTest, TestAddNotWarmedUpRowset) { |
| auto rowset = create_rowset(Version(25, 25)); |
| ASSERT_NE(rowset, nullptr); |
| |
| _tablet->add_not_warmed_up_rowset(rowset->rowset_id()); |
| |
| WarmUpState state = _tablet->get_rowset_warmup_state(rowset->rowset_id()); |
| WarmUpState expected_state = |
| WarmUpState {WarmUpTriggerSource::SYNC_ROWSET, WarmUpProgress::DOING}; |
| EXPECT_EQ(state, expected_state); |
| } |
| |
| // Test that add_warmed_up_rowset can override add_not_warmed_up_rowset |
| TEST_F(CloudTabletWarmUpStateTest, TestWarmedUpOverridesNotWarmedUp) { |
| auto rowset = create_rowset(Version(26, 26)); |
| ASSERT_NE(rowset, nullptr); |
| |
| // First mark as not warmed up |
| _tablet->add_not_warmed_up_rowset(rowset->rowset_id()); |
| EXPECT_FALSE(_tablet->is_rowset_warmed_up(rowset->rowset_id())); |
| |
| // Then mark as warmed up |
| _tablet->add_warmed_up_rowset(rowset->rowset_id()); |
| EXPECT_TRUE(_tablet->is_rowset_warmed_up(rowset->rowset_id())); |
| } |
| |
| class CloudTabletDeleteRowsetsForSchemaChangeTest : public testing::Test { |
| public: |
| CloudTabletDeleteRowsetsForSchemaChangeTest() : _engine(CloudStorageEngine(EngineOptions {})) {} |
| |
| void SetUp() override { |
| _tablet_meta.reset(new TabletMeta(1, 2, 15673, 15674, 4, 5, TTabletSchema(), 6, {{7, 8}}, |
| UniqueId(9, 10), TTabletType::TABLET_TYPE_DISK, |
| TCompressionType::LZ4F)); |
| _tablet = |
| std::make_shared<CloudTablet>(_engine, std::make_shared<TabletMeta>(*_tablet_meta)); |
| } |
| void TearDown() override {} |
| |
| RowsetSharedPtr create_rowset(Version version) { |
| auto rs_meta = std::make_shared<RowsetMeta>(); |
| rs_meta->set_rowset_type(BETA_ROWSET); |
| rs_meta->set_version(version); |
| rs_meta->set_rowset_id(_engine.next_rowset_id()); |
| rs_meta->set_tablet_schema(_tablet->tablet_schema()); |
| RowsetSharedPtr rowset; |
| Status st = RowsetFactory::create_rowset(nullptr, "", rs_meta, &rowset); |
| if (!st.ok()) { |
| return nullptr; |
| } |
| return rowset; |
| } |
| |
| protected: |
| TabletMetaSharedPtr _tablet_meta; |
| std::shared_ptr<CloudTablet> _tablet; |
| CloudStorageEngine _engine; |
| }; |
| |
| // Simulate the DORIS-25014 scenario: |
| // - New tablet has compacted rowset [2-6] from compaction during SC |
| // - SC produces individual output rowsets [2],[3],[4],[5],[6] |
| // - Without the fix, add_rowsets fails to remove [2-6] because |
| // [2].contains([2-6]) = false |
| // - With delete_rowsets_for_schema_change, the stale compaction rowset is |
| // removed from both _rs_version_map and version graph before add_rowsets |
| TEST_F(CloudTabletDeleteRowsetsForSchemaChangeTest, TestSchemaChangeDeletesCompactionRowset) { |
| // Setup: add placeholder [0-1] and compacted rowset [2-6] |
| auto rs_placeholder = create_rowset(Version(0, 1)); |
| auto rs_compacted = create_rowset(Version(2, 6)); |
| ASSERT_NE(rs_placeholder, nullptr); |
| ASSERT_NE(rs_compacted, nullptr); |
| |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| _tablet->add_rowsets({rs_placeholder, rs_compacted}, false, wlock, false); |
| } |
| // Verify initial state |
| ASSERT_EQ(_tablet->rowset_map().size(), 2); |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(2, 6))); |
| |
| // SC produces individual rowsets |
| std::vector<RowsetSharedPtr> sc_output; |
| for (int v = 2; v <= 6; v++) { |
| auto rs = create_rowset(Version(v, v)); |
| ASSERT_NE(rs, nullptr); |
| sc_output.push_back(rs); |
| } |
| |
| // Simulate delete_rowsets_for_schema_change + add_rowsets |
| int64_t alter_version = 6; |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| // Collect rowsets in [2, alter_version] |
| std::vector<RowsetSharedPtr> to_delete; |
| for (auto& [v, rs] : _tablet->rowset_map()) { |
| if (v.first >= 2 && v.second <= alter_version) { |
| to_delete.push_back(rs); |
| } |
| } |
| ASSERT_EQ(to_delete.size(), 1); // only [2-6] |
| ASSERT_EQ(to_delete[0]->version(), Version(2, 6)); |
| |
| _tablet->delete_rowsets_for_schema_change(to_delete, wlock); |
| |
| // [2-6] should be removed from rs_version_map |
| ASSERT_FALSE(_tablet->rowset_map().count(Version(2, 6))); |
| // Should NOT go to stale (to avoid stale path conflicts), but to unused |
| ASSERT_FALSE(_tablet->has_stale_rowsets()); |
| ASSERT_TRUE(_tablet->need_remove_unused_rowsets()); |
| |
| _tablet->add_rowsets(std::move(sc_output), false, wlock, false); |
| } |
| |
| // Verify: individual SC rowsets are now in rs_version_map |
| ASSERT_EQ(_tablet->rowset_map().size(), 6); // [0-1] + 5 individual |
| for (int v = 2; v <= 6; v++) { |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(v, v))) |
| << "Missing version " << v << "-" << v; |
| } |
| ASSERT_FALSE(_tablet->rowset_map().count(Version(2, 6))); |
| |
| // Verify: capture_consistent_versions works correctly (no stale edges) |
| auto versions_result = _tablet->capture_consistent_versions_unlocked(Version(0, 6), {}); |
| ASSERT_TRUE(versions_result.has_value()) << versions_result.error(); |
| auto& versions = versions_result.value(); |
| ASSERT_EQ(versions.size(), 6); // [0-1] + [2],[3],[4],[5],[6] |
| ASSERT_EQ(versions[0], Version(0, 1)); |
| for (int i = 0; i < 5; i++) { |
| ASSERT_EQ(versions[i + 1], Version(2 + i, 2 + i)); |
| } |
| } |
| |
| TEST_F(CloudTabletDeleteRowsetsForSchemaChangeTest, |
| TestReplaceSchemaChangeOutputCleansPollutedTmpGraph) { |
| auto rs_placeholder = create_rowset(Version(0, 1)); |
| auto rs_sc_2 = create_rowset(Version(2, 2)); |
| auto rs_sc_3 = create_rowset(Version(3, 3)); |
| auto rs_compacted = create_rowset(Version(2, 3)); |
| auto rs_post_alter = create_rowset(Version(4, 4)); |
| ASSERT_NE(rs_placeholder, nullptr); |
| ASSERT_NE(rs_sc_2, nullptr); |
| ASSERT_NE(rs_sc_3, nullptr); |
| ASSERT_NE(rs_compacted, nullptr); |
| ASSERT_NE(rs_post_alter, nullptr); |
| |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| _tablet->add_rowsets({rs_placeholder, rs_sc_2, rs_sc_3, rs_compacted, rs_post_alter}, false, |
| wlock, false); |
| } |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(2, 2))); |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(3, 3))); |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(2, 3))); |
| |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| _tablet->replace_rowsets_with_schema_change_output({rs_sc_2, rs_sc_3}, 3, wlock, "test", |
| false); |
| } |
| |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(2, 2))); |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(3, 3))); |
| ASSERT_FALSE(_tablet->rowset_map().count(Version(2, 3))); |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(4, 4))); |
| ASSERT_FALSE(_tablet->need_remove_unused_rowsets()); |
| |
| auto versions_result = _tablet->capture_consistent_versions_unlocked(Version(0, 4), {}); |
| ASSERT_TRUE(versions_result.has_value()) << versions_result.error(); |
| auto& versions = versions_result.value(); |
| ASSERT_EQ(versions.size(), 4); |
| ASSERT_EQ(versions[0], Version(0, 1)); |
| ASSERT_EQ(versions[1], Version(2, 2)); |
| ASSERT_EQ(versions[2], Version(3, 3)); |
| ASSERT_EQ(versions[3], Version(4, 4)); |
| } |
| |
| // Test that delete_rowsets_for_schema_change with empty input is a no-op |
| TEST_F(CloudTabletDeleteRowsetsForSchemaChangeTest, TestEmptyDeleteIsNoop) { |
| auto rs = create_rowset(Version(0, 1)); |
| ASSERT_NE(rs, nullptr); |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| _tablet->add_rowsets({rs}, false, wlock, false); |
| } |
| ASSERT_EQ(_tablet->rowset_map().size(), 1); |
| |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| _tablet->delete_rowsets_for_schema_change({}, wlock); |
| } |
| ASSERT_EQ(_tablet->rowset_map().size(), 1); |
| ASSERT_FALSE(_tablet->has_stale_rowsets()); |
| } |
| |
| // Test with multiple compaction rowsets spanning different version ranges |
| TEST_F(CloudTabletDeleteRowsetsForSchemaChangeTest, TestMultipleCompactionRowsets) { |
| auto rs_placeholder = create_rowset(Version(0, 1)); |
| auto rs_comp1 = create_rowset(Version(2, 5)); |
| auto rs_comp2 = create_rowset(Version(6, 10)); |
| auto rs_post = create_rowset(Version(11, 11)); // after alter_version, should NOT be deleted |
| ASSERT_NE(rs_placeholder, nullptr); |
| ASSERT_NE(rs_comp1, nullptr); |
| ASSERT_NE(rs_comp2, nullptr); |
| ASSERT_NE(rs_post, nullptr); |
| |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| _tablet->add_rowsets({rs_placeholder, rs_comp1, rs_comp2, rs_post}, false, wlock, false); |
| } |
| ASSERT_EQ(_tablet->rowset_map().size(), 4); |
| |
| // SC output: individual rowsets for versions 2-10 |
| std::vector<RowsetSharedPtr> sc_output; |
| for (int v = 2; v <= 10; v++) { |
| auto rs = create_rowset(Version(v, v)); |
| ASSERT_NE(rs, nullptr); |
| sc_output.push_back(rs); |
| } |
| |
| int64_t alter_version = 10; |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| std::vector<RowsetSharedPtr> to_delete; |
| for (auto& [v, rs] : _tablet->rowset_map()) { |
| if (v.first >= 2 && v.second <= alter_version) { |
| to_delete.push_back(rs); |
| } |
| } |
| ASSERT_EQ(to_delete.size(), 2); // [2-5] and [6-10] |
| |
| _tablet->delete_rowsets_for_schema_change(to_delete, wlock); |
| |
| // Post-alter rowset should survive |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(11, 11))); |
| ASSERT_FALSE(_tablet->rowset_map().count(Version(2, 5))); |
| ASSERT_FALSE(_tablet->rowset_map().count(Version(6, 10))); |
| |
| _tablet->add_rowsets(std::move(sc_output), false, wlock, false); |
| } |
| |
| // Verify: [0-1], [2],[3],...,[10], [11-11] |
| ASSERT_EQ(_tablet->rowset_map().size(), 11); |
| |
| // Verify capture |
| auto versions_result = _tablet->capture_consistent_versions_unlocked(Version(0, 11), {}); |
| ASSERT_TRUE(versions_result.has_value()) << versions_result.error(); |
| auto& versions = versions_result.value(); |
| ASSERT_EQ(versions.size(), 11); |
| ASSERT_EQ(versions[0], Version(0, 1)); |
| for (int i = 0; i < 9; i++) { |
| ASSERT_EQ(versions[i + 1], Version(2 + i, 2 + i)); |
| } |
| ASSERT_EQ(versions[10], Version(11, 11)); |
| } |
| |
| TEST_F(CloudTabletDeleteRowsetsForSchemaChangeTest, TestFillVersionHolesBeforeSchemaChangeRunning) { |
| static_cast<void>(_tablet->set_tablet_state(TABLET_NOTREADY)); |
| _tablet->set_alter_version(10); |
| |
| auto rs_placeholder = create_rowset(Version(0, 1)); |
| auto rs_historical = create_rowset(Version(2, 10)); |
| auto rs_after_first_hole = create_rowset(Version(12, 12)); |
| auto rs_after_second_hole = create_rowset(Version(14, 14)); |
| ASSERT_NE(rs_placeholder, nullptr); |
| ASSERT_NE(rs_historical, nullptr); |
| ASSERT_NE(rs_after_first_hole, nullptr); |
| ASSERT_NE(rs_after_second_hole, nullptr); |
| |
| cloud::CloudMetaMgr meta_mgr; |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| _tablet->add_rowsets( |
| {rs_placeholder, rs_historical, rs_after_first_hole, rs_after_second_hole}, false, |
| wlock, false); |
| ASSERT_FALSE(_tablet->rowset_map().count(Version(11, 11))); |
| ASSERT_FALSE(_tablet->rowset_map().count(Version(13, 13))); |
| ASSERT_FALSE(_tablet->capture_consistent_versions_unlocked(Version(0, 14), {}).has_value()); |
| |
| auto status = |
| meta_mgr.fill_version_holes(_tablet.get(), _tablet->max_version_unlocked(), wlock); |
| ASSERT_TRUE(status.ok()) << status.to_string(); |
| ASSERT_TRUE(_tablet->set_tablet_state(TABLET_RUNNING).ok()); |
| } |
| |
| for (const Version& version : {Version(11, 11), Version(13, 13)}) { |
| ASSERT_TRUE(_tablet->rowset_map().count(version)); |
| auto hole_rowset = _tablet->rowset_map().at(version); |
| ASSERT_TRUE(hole_rowset->empty()); |
| ASSERT_TRUE(hole_rowset->is_hole_rowset()); |
| } |
| |
| auto versions_result = _tablet->capture_consistent_versions_unlocked(Version(0, 14), {}); |
| ASSERT_TRUE(versions_result.has_value()) << versions_result.error(); |
| const auto& versions = versions_result.value(); |
| ASSERT_EQ(versions.size(), 6); |
| ASSERT_EQ(versions[0], Version(0, 1)); |
| ASSERT_EQ(versions[1], Version(2, 10)); |
| ASSERT_EQ(versions[2], Version(11, 11)); |
| ASSERT_EQ(versions[3], Version(12, 12)); |
| ASSERT_EQ(versions[4], Version(13, 13)); |
| ASSERT_EQ(versions[5], Version(14, 14)); |
| } |
| |
| // Reproduce the CI crash scenario: SC delete puts rowsets to stale, then |
| // compaction creates a new stale path with overlapping version keys. When |
| // one stale path is cleaned, the other hits DCHECK(false) because the |
| // version is already removed from _stale_rs_version_map. |
| // With the fix (bypassing stale tracking), this should not happen. |
| TEST_F(CloudTabletDeleteRowsetsForSchemaChangeTest, TestNoStalePathConflictWithCompaction) { |
| // Setup: [0-1] placeholder, [2-6] compaction product during SC |
| auto rs_placeholder = create_rowset(Version(0, 1)); |
| auto rs_compacted = create_rowset(Version(2, 6)); |
| ASSERT_NE(rs_placeholder, nullptr); |
| ASSERT_NE(rs_compacted, nullptr); |
| |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| _tablet->add_rowsets({rs_placeholder, rs_compacted}, false, wlock, false); |
| } |
| |
| // SC output: individual rowsets [2],[3],[4],[5],[6] |
| std::vector<RowsetSharedPtr> sc_output; |
| for (int v = 2; v <= 6; v++) { |
| sc_output.push_back(create_rowset(Version(v, v))); |
| } |
| |
| // Step 1: delete_rowsets_for_schema_change + add SC output |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| _tablet->delete_rowsets_for_schema_change({rs_compacted}, wlock); |
| _tablet->add_rowsets(std::move(sc_output), false, wlock, false); |
| } |
| // Stale should be empty — SC delete bypasses stale tracking |
| ASSERT_FALSE(_tablet->has_stale_rowsets()); |
| |
| // Step 2: compaction merges SC output [2],[3],[4],[5],[6] -> [2-6] |
| auto rs_new_compacted = create_rowset(Version(2, 6)); |
| std::vector<RowsetSharedPtr> compaction_input; |
| { |
| std::unique_lock wlock(_tablet->get_header_lock()); |
| for (auto& [v, rs] : _tablet->rowset_map()) { |
| if (v.first >= 2 && v.second <= 6) { |
| compaction_input.push_back(rs); |
| } |
| } |
| ASSERT_EQ(compaction_input.size(), 5); |
| // Normal compaction delete_rowsets — this WILL use stale tracking |
| _tablet->delete_rowsets(compaction_input, wlock); |
| _tablet->add_rowsets({rs_new_compacted}, false, wlock, false); |
| } |
| // Now stale has the compaction inputs |
| ASSERT_TRUE(_tablet->has_stale_rowsets()); |
| |
| // Step 3: delete_expired_stale_rowsets — this is where CI crashed |
| // With old code: stale path from SC and compaction both reference [2-6] key, |
| // causing DCHECK(false). With fix: only compaction stale path exists, no conflict. |
| config::tablet_rowset_stale_sweep_time_sec = 0; // expire immediately |
| ASSERT_NO_FATAL_FAILURE(_tablet->delete_expired_stale_rowsets()); |
| |
| // Verify final state: [0-1] and [2-6] active, no stale left |
| ASSERT_EQ(_tablet->rowset_map().size(), 2); |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(0, 1))); |
| ASSERT_TRUE(_tablet->rowset_map().count(Version(2, 6))); |
| ASSERT_FALSE(_tablet->has_stale_rowsets()); |
| } |
| |
| } // namespace doris |