This is an automated email from the ASF dual-hosted git repository.

yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new e71218f7476 branch-4.1: [fix](cloud) Avoid conflicts when recycling 
transactions with the same label #65746 (#67478)
e71218f7476 is described below

commit e71218f74762a81b2799cb0f48698cda8f8c3708
Author: github-actions[bot] 
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Sun Sep 6 06:01:10 2026 +0800

    branch-4.1: [fix](cloud) Avoid conflicts when recycling transactions with 
the same label #65746 (#67478)
    
    Cherry-picked from #65746
    
    Co-authored-by: Yixuan Wang <[email protected]>
---
 cloud/src/recycler/recycler.cpp | 104 +++++++++++-------
 cloud/test/recycler_test.cpp    | 231 +++++++++++++++++++++++++++++++++++++---
 2 files changed, 280 insertions(+), 55 deletions(-)

diff --git a/cloud/src/recycler/recycler.cpp b/cloud/src/recycler/recycler.cpp
index 4023cd1489f..a7c551a74d6 100644
--- a/cloud/src/recycler/recycler.cpp
+++ b/cloud/src/recycler/recycler.cpp
@@ -6706,7 +6706,7 @@ int InstanceRecycler::recycle_expired_txn_label() {
     std::string end_recycle_txn_key;
     recycle_txn_key(recycle_txn_key_info0, &begin_recycle_txn_key);
     recycle_txn_key(recycle_txn_key_info1, &end_recycle_txn_key);
-    std::vector<std::string> recycle_txn_info_keys;
+    std::unordered_map<std::string, std::vector<std::string>> 
recycle_txn_keys_by_label;
 
     LOG_WARNING("begin to recycle expired txn").tag("instance_id", 
instance_id_);
 
@@ -6747,7 +6747,17 @@ int InstanceRecycler::recycle_expired_txn_label() {
              current_time_ms)) {
             VLOG_DEBUG << "found recycle txn, key=" << hex(k);
             num_expired++;
-            recycle_txn_info_keys.emplace_back(k);
+
+            std::string_view k1 = k;
+            k1.remove_prefix(1); // Remove key space
+            std::vector<std::tuple<std::variant<int64_t, std::string>, int, 
int>> out;
+            if (decode_key(&k1, &out) != 0) {
+                LOG_ERROR("failed to decode key").tag("key", hex(k));
+                return -1;
+            }
+            int64_t db_id = std::get<int64_t>(std::get<0>(out[3]));
+            auto label_key = txn_label_key({instance_id_, db_id, 
recycle_txn_pb.label()});
+            recycle_txn_keys_by_label[label_key].emplace_back(k);
         }
         return 0;
     };
@@ -6856,35 +6866,48 @@ int InstanceRecycler::recycle_expired_txn_label() {
 
     auto loop_done = [&]() -> int {
         DORIS_CLOUD_DEFER {
-            recycle_txn_info_keys.clear();
+            recycle_txn_keys_by_label.clear();
         };
         TEST_SYNC_POINT_CALLBACK(
-                
"InstanceRecycler::recycle_expired_txn_label.check_recycle_txn_info_keys",
-                &recycle_txn_info_keys);
-        for (const auto& k : recycle_txn_info_keys) {
-            concurrent_delete_executor.add([&]() {
-                int ret = delete_recycle_txn_kv(k);
-                if (ret == 1) {
-                    const int max_retry = std::max(1, 
config::recycle_txn_delete_max_retry_times);
-                    for (int i = 1; i <= max_retry; ++i) {
-                        LOG(WARNING) << "txn conflict, retry times=" << i << " 
key=" << hex(k);
-                        ret = delete_recycle_txn_kv(k);
-                        // clang-format off
-                        TEST_SYNC_POINT_CALLBACK(
-                                
"InstanceRecycler::recycle_expired_txn_label.delete_recycle_txn_kv_error", 
&ret);
-                        // clang-format off
-                        if (ret != 1) {
-                            break;
-                        }
-                        // random sleep 0-100 ms to retry
-                        
std::this_thread::sleep_for(std::chrono::milliseconds(rand() % 100));
+                
"InstanceRecycler::recycle_expired_txn_label.check_recycle_txn_keys_by_label",
+                &recycle_txn_keys_by_label);
+        auto delete_recycle_txn_kv_with_retry = [&](const std::string& k) -> 
int {
+            int ret = delete_recycle_txn_kv(k);
+            TEST_SYNC_POINT_CALLBACK(
+                    
"InstanceRecycler::recycle_expired_txn_label.delete_recycle_txn_kv_error",
+                    &ret);
+            if (ret == 1) {
+                const int max_retry = std::max(1, 
config::recycle_txn_delete_max_retry_times);
+                for (int i = 1; i <= max_retry; ++i) {
+                    LOG(WARNING) << "txn conflict, retry times=" << i << " 
key=" << hex(k);
+                    ret = delete_recycle_txn_kv(k);
+                    TEST_SYNC_POINT_CALLBACK(
+                            
"InstanceRecycler::recycle_expired_txn_label.delete_recycle_txn_kv_"
+                            "error",
+                            &ret);
+                    if (ret != 1) {
+                        break;
                     }
+                    // random sleep 0-100 ms to retry
+                    
std::this_thread::sleep_for(std::chrono::milliseconds(rand() % 100));
                 }
-                if (ret != 0) {
-                    LOG_WARNING("failed to delete recycle txn kv")
-                            .tag("instance id", instance_id_)
-                            .tag("key", hex(k));
-                    return -1;
+            }
+            return ret;
+        };
+
+        for (auto& [label_key, txn_keys] : recycle_txn_keys_by_label) {
+            concurrent_delete_executor.add([&, txn_keys = std::move(txn_keys), 
label_key]() {
+                VLOG_DEBUG << "recycle txn label group, key=" << hex(label_key)
+                           << " txn_count=" << txn_keys.size();
+                for (const auto& k : txn_keys) {
+                    int ret = delete_recycle_txn_kv_with_retry(k);
+                    if (ret != 0) {
+                        LOG_WARNING("failed to delete recycle txn kv")
+                                .tag("instance id", instance_id_)
+                                .tag("key", hex(k))
+                                .tag("label_key", hex(label_key));
+                        return -1;
+                    }
                 }
                 return 0;
             });
@@ -7646,8 +7669,8 @@ int InstanceRecycler::scan_and_statistics_rowsets() {
     std::string recyc_rs_key0;
     std::string recyc_rs_key1;
     recycle_rowset_key(recyc_rs_key_info0, &recyc_rs_key0);
-                recycle_rowset_key(recyc_rs_key_info1, &recyc_rs_key1);
-       int64_t earlest_ts = std::numeric_limits<int64_t>::max();
+    recycle_rowset_key(recyc_rs_key_info1, &recyc_rs_key1);
+    int64_t earlest_ts = std::numeric_limits<int64_t>::max();
 
     auto handle_rowset_kv = [&, this](std::string_view k, std::string_view v) 
-> int {
         RecycleRowsetPB rowset;
@@ -7671,7 +7694,8 @@ int InstanceRecycler::scan_and_statistics_rowsets() {
             metrics_context.total_need_recycle_num++;
             metrics_context.total_need_recycle_data_size += 
rowset.rowset_meta().total_disk_size();
             segment_metrics_context_.total_need_recycle_num += 
rowset.rowset_meta().num_segments();
-            segment_metrics_context_.total_need_recycle_data_size += 
rowset.rowset_meta().total_disk_size();
+            segment_metrics_context_.total_need_recycle_data_size +=
+                    rowset.rowset_meta().total_disk_size();
             return 0;
         }
 
@@ -7724,7 +7748,7 @@ int InstanceRecycler::scan_and_statistics_tmp_rowsets() {
         DCHECK_GT(rowset.txn_id(), 0)
                 << "txn_id=" << rowset.txn_id() << " rowset=" << 
rowset.ShortDebugString();
 
-        if(!rowset.has_is_recycled() || !rowset.is_recycled()) {
+        if (!rowset.has_is_recycled() || !rowset.is_recycled()) {
             return 0;
         }
 
@@ -7801,7 +7825,8 @@ int 
InstanceRecycler::scan_and_statistics_abort_timeout_txn() {
         return 0;
     };
 
-    int ret = scan_and_recycle(begin_txn_running_key, end_txn_running_key, 
std::move(handle_abort_timeout_txn_kv));
+    int ret = scan_and_recycle(begin_txn_running_key, end_txn_running_key,
+                               std::move(handle_abort_timeout_txn_kv));
     metrics_context.report(true);
     return ret;
 }
@@ -7835,7 +7860,8 @@ int 
InstanceRecycler::scan_and_statistics_expired_txn_label() {
         return 0;
     };
 
-    int ret = scan_and_recycle(begin_recycle_txn_key, end_recycle_txn_key, 
std::move(handle_expired_txn_label_kv));
+    int ret = scan_and_recycle(begin_recycle_txn_key, end_recycle_txn_key,
+                               std::move(handle_expired_txn_label_kv));
     metrics_context.report(true);
     return ret;
 }
@@ -8066,7 +8092,7 @@ int InstanceRecycler::scan_and_statistics_restore_jobs() {
             return 0;
         }
         metrics_context.total_need_recycle_num++;
-        if(restore_job_pb.need_recycle_data()) {
+        if (restore_job_pb.need_recycle_data()) {
             scan_tablet_and_statistics(restore_job_pb.tablet_id(), 
metrics_context);
         }
         return 0;
@@ -8108,7 +8134,7 @@ void 
InstanceRecycler::scan_and_statistics_operation_logs() {
 
         OperationLogReferenceInfo ref_info;
         if (recycle_checker.can_recycle(log_versionstamp, 
operation_log.min_timestamp(),
-                                         &ref_info)) {
+                                        &ref_info)) {
             metrics_context.total_need_recycle_num++;
             metrics_context.total_need_recycle_data_size += 
operation_log.ByteSizeLong();
         }
@@ -8264,8 +8290,8 @@ int InstanceRecycler::cleanup_rowset_metadata(const 
std::vector<RowsetDeleteTask
         std::string dbm_start_key =
                 meta_delete_bitmap_key({reference_instance_id, tablet_id, 
rowset_id, 0, 0});
         std::string dbm_end_key = meta_delete_bitmap_key(
-                {reference_instance_id, tablet_id, rowset_id,
-                 std::numeric_limits<int64_t>::max(), 
std::numeric_limits<int64_t>::max()});
+                {reference_instance_id, tablet_id, rowset_id, 
std::numeric_limits<int64_t>::max(),
+                 std::numeric_limits<int64_t>::max()});
         txn->remove(dbm_start_key, dbm_end_key);
         LOG_INFO("remove delete bitmap kv in cleanup phase")
                 .tag("instance_id", instance_id_)
@@ -8288,8 +8314,8 @@ int InstanceRecycler::cleanup_rowset_metadata(const 
std::vector<RowsetDeleteTask
 
         // Remove versioned meta rowset key
         if (!task.versioned_rowset_key.empty()) {
-            versioned::document_remove<RowsetMetaCloudPB>(
-                txn.get(), task.versioned_rowset_key, task.versionstamp);
+            versioned::document_remove<RowsetMetaCloudPB>(txn.get(), 
task.versioned_rowset_key,
+                                                          task.versionstamp);
             LOG_INFO("remove versioned meta rowset key in cleanup phase")
                     .tag("instance_id", instance_id_)
                     .tag("tablet_id", tablet_id)
diff --git a/cloud/test/recycler_test.cpp b/cloud/test/recycler_test.cpp
index d23264c0fbe..56a736f048d 100644
--- a/cloud/test/recycler_test.cpp
+++ b/cloud/test/recycler_test.cpp
@@ -9233,7 +9233,9 @@ void 
make_single_txn_related_kvs(std::shared_ptr<cloud::TxnKv> txn_kv, int64_t i
     } else {
         recycle_txn_pb.set_creation_time(current_time);
     }
-    recycle_txn_pb.set_label("recycle_txn_key_info_label_" + 
std::to_string(i));
+    // Production writes RecycleTxnPB.label and TxnInfoPB.label from the same 
transaction label.
+    const std::string label = "txn_label_" + std::to_string(i);
+    recycle_txn_pb.set_label(label);
     if (!recycle_txn_pb.SerializeToString(&recycle_txn_info_val)) {
         LOG_WARNING("failed to serialize recycle txn info")
                 .tag("key", hex(recycle_txn_info_key))
@@ -9263,7 +9265,7 @@ void 
make_single_txn_related_kvs(std::shared_ptr<cloud::TxnKv> txn_kv, int64_t i
     std::string info_val;
     TxnInfoPB txn_info_pb;
     txn_info_pb.add_sub_txn_ids(sub_txn_id);
-    txn_info_pb.set_label("txn_info_label_" + std::to_string(i));
+    txn_info_pb.set_label(label);
     if (!txn_info_pb.SerializeToString(&info_val)) {
         LOG_WARNING("failed to serialize txn info")
                 .tag("key", hex(info_key))
@@ -9489,20 +9491,24 @@ TEST(RecyclerTest, 
concurrent_recycle_txn_label_failure_test) {
 
     auto txn_kv = mem_txn_kv;
     ASSERT_TRUE(txn_kv.get()) << "exit get MemTxnKv error" << std::endl;
-    make_multiple_txn_info_kvs(txn_kv, 20000, 15000);
-    check_multiple_txn_info_kvs(txn_kv, 20000);
+    make_multiple_txn_info_kvs(txn_kv, 40000, 30000);
+    check_multiple_txn_info_kvs(txn_kv, 40000);
 
     auto* sp = SyncPoint::get_instance();
     DORIS_CLOUD_DEFER {
         SyncPoint::get_instance()->clear_all_call_backs();
     };
-    size_t recycle_txn_info_keys_cnt = 0;
-    
sp->set_call_back("InstanceRecycler::recycle_expired_txn_label.check_recycle_txn_info_keys",
-                      [&](auto&& args) {
-                          auto* recycle_txn_info_keys =
-                                  
try_any_cast<std::vector<std::string>*>(args[0]);
-                          recycle_txn_info_keys_cnt += 
recycle_txn_info_keys->size();
-                      });
+    size_t recycle_txn_keys_cnt = 0;
+    sp->set_call_back(
+            
"InstanceRecycler::recycle_expired_txn_label.check_recycle_txn_keys_by_label",
+            [&](auto&& args) {
+                auto* recycle_txn_keys_by_label =
+                        try_any_cast<std::unordered_map<std::string, 
std::vector<std::string>>*>(
+                                args[0]);
+                for (const auto& entry : *recycle_txn_keys_by_label) {
+                    recycle_txn_keys_cnt += entry.second.size();
+                }
+            });
     sp->set_call_back("InstanceRecycler::recycle_expired_txn_label.failure", 
[](auto&& args) {
         auto* ret = try_any_cast<int*>(args[0]);
         *ret = -1;
@@ -9520,7 +9526,7 @@ TEST(RecyclerTest, 
concurrent_recycle_txn_label_failure_test) {
     std::cout << "recycle expired txn label cost="
               << std::chrono::duration_cast<std::chrono::milliseconds>(finish 
- start).count()
               << "ms" << std::endl;
-    check_multiple_txn_info_kvs(txn_kv, (20000 - recycle_txn_info_keys_cnt));
+    check_multiple_txn_info_kvs(txn_kv, (40000 - recycle_txn_keys_cnt));
 }
 TEST(RecyclerTest, concurrent_recycle_txn_label_conflict_test) {
     config::label_keep_max_second = 0;
@@ -9671,7 +9677,7 @@ TEST(RecyclerTest, 
concurrent_recycle_txn_label_conflict_test) {
     std::cout << "Update label after count: " << update_label_after_count << 
std::endl;
     std::cout << "Transaction conflict count: " << txn_conflict_count << 
std::endl;
 
-    EXPECT_GT(txn_conflict_count, 0) << "txn_conflict sync point should be 
triggered";
+    EXPECT_EQ(txn_conflict_count, 0) << "txn conflicts should not occur within 
one label group";
 
     std::unique_ptr<Transaction> verify_txn;
     ASSERT_EQ(mem_txn_kv->create_txn(&verify_txn), TxnErrorCode::TXN_OK);
@@ -9699,7 +9705,7 @@ TEST(RecyclerTest, 
concurrent_recycle_txn_label_conflict_test) {
     }
 }
 
-TEST(RecyclerTest, recycle_txn_label_deal_with_conflict_error_test) {
+TEST(RecyclerTest, recycle_txn_label_propagate_delete_error_test) {
     config::label_keep_max_second = 0;
     config::recycle_pool_parallelism = 20;
 
@@ -9845,10 +9851,203 @@ TEST(RecyclerTest, 
recycle_txn_label_deal_with_conflict_error_test) {
                               std::make_shared<TxnLazyCommitter>(mem_txn_kv));
     ASSERT_EQ(recycler.init(), 0);
 
-    // deal with conflict but error during recycle
+    // Propagate a recycle error without relying on an internal label conflict.
+    ASSERT_EQ(recycler.recycle_expired_txn_label(), -1);
+
+    EXPECT_EQ(txn_conflict_count, 0) << "txn conflicts should not occur within 
one label group";
+}
+
+TEST(RecyclerTest, recycle_txn_label_retry_after_conflict_test) {
+    config::label_keep_max_second = 0;
+
+    auto txn_kv = 
std::dynamic_pointer_cast<TxnKv>(std::make_shared<MemTxnKv>());
+    ASSERT_NE(txn_kv.get(), nullptr);
+    ASSERT_EQ(txn_kv->init(), 0);
+    auto resource_mgr = std::make_shared<MockResourceManager>(txn_kv);
+    auto rate_limiter = std::make_shared<RateLimiter>();
+    auto snapshot = std::make_shared<SnapshotManager>(txn_kv);
+    auto meta_service =
+            std::make_unique<MetaServiceImpl>(txn_kv, resource_mgr, 
rate_limiter, snapshot);
+
+    constexpr int64_t db_id = 10001;
+    constexpr int64_t table_id = 20001;
+    const std::string cloud_unique_id = 
"recycle_txn_label_retry_after_conflict_test";
+    const std::string label = "recycle_txn_label_retry_after_conflict_test";
+
+    int64_t recycled_txn_id = -1;
+    {
+        brpc::Controller cntl;
+        BeginTxnRequest req;
+        BeginTxnResponse res;
+        req.set_cloud_unique_id(cloud_unique_id);
+        auto* txn_info = req.mutable_txn_info();
+        txn_info->set_db_id(db_id);
+        txn_info->set_label(label);
+        txn_info->add_table_ids(table_id);
+        txn_info->set_timeout_ms(36000);
+        
meta_service->begin_txn(reinterpret_cast<::google::protobuf::RpcController*>(&cntl),
 &req,
+                                &res, nullptr);
+        ASSERT_EQ(res.status().code(), MetaServiceCode::OK) << 
res.ShortDebugString();
+        ASSERT_TRUE(res.has_txn_id());
+        recycled_txn_id = res.txn_id();
+    }
+    {
+        brpc::Controller cntl;
+        AbortTxnRequest req;
+        AbortTxnResponse res;
+        req.set_cloud_unique_id(cloud_unique_id);
+        req.set_db_id(db_id);
+        req.set_txn_id(recycled_txn_id);
+        req.set_reason("test");
+        
meta_service->abort_txn(reinterpret_cast<::google::protobuf::RpcController*>(&cntl),
 &req,
+                                &res, nullptr);
+        ASSERT_EQ(res.status().code(), MetaServiceCode::OK) << 
res.ShortDebugString();
+    }
+
+    auto* sp = SyncPoint::get_instance();
+    DORIS_CLOUD_DEFER {
+        SyncPoint::get_instance()->clear_all_call_backs();
+        SyncPoint::get_instance()->disable_processing();
+    };
+
+    std::atomic<int> before_commit_count {0};
+    std::atomic<int> txn_conflict_count {0};
+    std::atomic<int> external_begin_code {-1};
+    std::atomic<int64_t> new_txn_id {-1};
+    
sp->set_call_back("InstanceRecycler::recycle_expired_txn_label.before_commit", 
[&](auto&&) {
+        if (before_commit_count.fetch_add(1) != 0) {
+            return;
+        }
+
+        brpc::Controller cntl;
+        BeginTxnRequest req;
+        BeginTxnResponse res;
+        req.set_cloud_unique_id(cloud_unique_id);
+        auto* txn_info = req.mutable_txn_info();
+        txn_info->set_db_id(db_id);
+        txn_info->set_label(label);
+        txn_info->add_table_ids(table_id);
+        txn_info->set_timeout_ms(36000);
+        
meta_service->begin_txn(reinterpret_cast<::google::protobuf::RpcController*>(&cntl),
 &req,
+                                &res, nullptr);
+        external_begin_code.store(static_cast<int>(res.status().code()));
+        if (res.has_txn_id()) {
+            new_txn_id.store(res.txn_id());
+        }
+    });
+    
sp->set_call_back("InstanceRecycler::recycle_expired_txn_label.txn_conflict",
+                      [&](auto&&) { txn_conflict_count.fetch_add(1); });
+    sp->enable_processing();
+
+    InstanceInfoPB instance;
+    instance.set_instance_id(mock_instance);
+    InstanceRecycler recycler(txn_kv, instance, thread_group,
+                              std::make_shared<TxnLazyCommitter>(txn_kv));
+    ASSERT_EQ(recycler.init(), 0);
+
+    ASSERT_EQ(recycler.recycle_expired_txn_label(), 0);
+    EXPECT_EQ(external_begin_code.load(), 
static_cast<int>(MetaServiceCode::OK));
+    EXPECT_GT(new_txn_id.load(), 0);
+    EXPECT_EQ(txn_conflict_count.load(), 1);
+    EXPECT_EQ(before_commit_count.load(), 2);
+
+    std::unique_ptr<Transaction> verify_txn;
+    ASSERT_EQ(txn_kv->create_txn(&verify_txn), TxnErrorCode::TXN_OK);
+    const std::string recycle_key = recycle_txn_key({mock_instance, db_id, 
recycled_txn_id});
+    std::string recycle_value;
+    EXPECT_EQ(verify_txn->get(recycle_key, &recycle_value), 
TxnErrorCode::TXN_KEY_NOT_FOUND);
+
+    const std::string label_key = txn_label_key({mock_instance, db_id, label});
+    std::string label_value;
+    ASSERT_EQ(verify_txn->get(label_key, &label_value), TxnErrorCode::TXN_OK);
+    TxnLabelPB txn_label;
+    ASSERT_TRUE(
+            txn_label.ParseFromArray(label_value.data(), label_value.size() - 
VERSION_STAMP_LEN));
+    ASSERT_EQ(txn_label.txn_ids_size(), 1);
+    EXPECT_EQ(txn_label.txn_ids(0), new_txn_id.load());
+
+    std::string new_info_value;
+    ASSERT_EQ(verify_txn->get(txn_info_key({mock_instance, db_id, 
new_txn_id.load()}),
+                              &new_info_value),
+              TxnErrorCode::TXN_OK);
+    TxnInfoPB new_txn_info;
+    ASSERT_TRUE(new_txn_info.ParseFromString(new_info_value));
+    EXPECT_EQ(new_txn_info.label(), label);
+}
+
+TEST(RecyclerTest, recycle_txn_label_retry_exhausted_then_recover_test) {
+    const int old_max_retry_times = config::recycle_txn_delete_max_retry_times;
+    DORIS_CLOUD_DEFER {
+        config::recycle_txn_delete_max_retry_times = old_max_retry_times;
+    };
+    config::label_keep_max_second = 0;
+    config::recycle_txn_delete_max_retry_times = 2;
+
+    auto mem_txn_kv = std::make_shared<MemTxnKv>();
+    ASSERT_EQ(mem_txn_kv->init(), 0);
+    make_single_txn_related_kvs(mem_txn_kv, 0, 1);
+
+    const std::string recycle_key = recycle_txn_key({instance_id, 0, 1000000});
+    const std::string label_key = txn_label_key({instance_id, 0, 
"txn_label_0"});
+
+    auto* sp = SyncPoint::get_instance();
+    DORIS_CLOUD_DEFER {
+        SyncPoint::get_instance()->clear_all_call_backs();
+        SyncPoint::get_instance()->disable_processing();
+    };
+
+    std::atomic<int> external_write_count {0};
+    std::atomic<int> external_write_error_count {0};
+    std::atomic<int> txn_conflict_count {0};
+    
sp->set_call_back("InstanceRecycler::recycle_expired_txn_label.before_commit", 
[&](auto&&) {
+        std::unique_ptr<Transaction> txn;
+        if (mem_txn_kv->create_txn(&txn) != TxnErrorCode::TXN_OK) {
+            external_write_error_count.fetch_add(1);
+            return;
+        }
+        std::string label_value;
+        if (txn->get(label_key, &label_value) != TxnErrorCode::TXN_OK) {
+            external_write_error_count.fetch_add(1);
+            return;
+        }
+        txn->put(label_key, label_value);
+        if (txn->commit() != TxnErrorCode::TXN_OK) {
+            external_write_error_count.fetch_add(1);
+            return;
+        }
+        external_write_count.fetch_add(1);
+    });
+    
sp->set_call_back("InstanceRecycler::recycle_expired_txn_label.txn_conflict",
+                      [&](auto&&) { txn_conflict_count.fetch_add(1); });
+    sp->enable_processing();
+
+    InstanceInfoPB instance;
+    instance.set_instance_id(instance_id);
+    InstanceRecycler recycler(mem_txn_kv, instance, thread_group,
+                              std::make_shared<TxnLazyCommitter>(mem_txn_kv));
+    ASSERT_EQ(recycler.init(), 0);
+
     ASSERT_EQ(recycler.recycle_expired_txn_label(), -1);
+    EXPECT_EQ(external_write_error_count.load(), 0);
+    EXPECT_EQ(external_write_count.load(), 3);
+    EXPECT_EQ(txn_conflict_count.load(), 3);
 
-    EXPECT_GT(txn_conflict_count, 0) << "txn_conflict sync point should be 
triggered";
+    {
+        std::unique_ptr<Transaction> verify_txn;
+        ASSERT_EQ(mem_txn_kv->create_txn(&verify_txn), TxnErrorCode::TXN_OK);
+        std::string recycle_value;
+        EXPECT_EQ(verify_txn->get(recycle_key, &recycle_value), 
TxnErrorCode::TXN_OK);
+    }
+
+    sp->clear_all_call_backs();
+    sp->disable_processing();
+    ASSERT_EQ(recycler.recycle_expired_txn_label(), 0);
+
+    std::unique_ptr<Transaction> verify_txn;
+    ASSERT_EQ(mem_txn_kv->create_txn(&verify_txn), TxnErrorCode::TXN_OK);
+    std::string value;
+    EXPECT_EQ(verify_txn->get(recycle_key, &value), 
TxnErrorCode::TXN_KEY_NOT_FOUND);
+    EXPECT_EQ(verify_txn->get(label_key, &value), 
TxnErrorCode::TXN_KEY_NOT_FOUND);
 }
 
 TEST(RecyclerTest, recycle_restore_job_complete_state) {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to