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

gavinchou pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new def2e00ff63 [fix](cloud) Track deleted instance recycle status (#66519)
def2e00ff63 is described below

commit def2e00ff6390d418ae34facf892a59afc8e3e2f
Author: Yixuan Wang <[email protected]>
AuthorDate: Tue Aug 11 21:14:56 2026 +0800

    [fix](cloud) Track deleted instance recycle status (#66519)
    
    Introduce a persistent recycle state for deleted instances to split
    instance recycling into separate data cleanup and metadata cleanup
    stages.
    
    The recycle process now transitions through:
    ```
    DATA_CLEANUP_PENDING
    -> METADATA_CLEANUP_PENDING
    -> CLEANUP_COMPLETED
    ```
    
    Each transition is persisted so that recycling can resume from the
    latest state after failures or retries.
    
    Also add an API to skip instance data cleanup by advancing the recycle
    state directly from DATA_CLEANUP_PENDING to METADATA_CLEANUP_PENDING.
    Subsequent recycle rounds will skip object data cleanup and continue
    with metadata cleanup.
    
    Repeated DROP operations preserve the existing recycle state instead of
    resetting the recycle process.
---
 cloud/src/common/http_helper.cpp                 |  50 ++++-
 cloud/src/common/http_helper.h                   |   6 +
 cloud/src/meta-service/meta_service.h            |   3 +
 cloud/src/meta-service/meta_service_resource.cpp |  70 +++++++
 cloud/src/recycler/recycler.cpp                  | 234 +++++++++++++++++++----
 cloud/src/recycler/recycler.h                    |  13 ++
 cloud/src/recycler/recycler_service.cpp          |  80 ++++++++
 cloud/src/recycler/recycler_service.h            |   3 +
 cloud/src/recycler/recycler_snapshot.cpp         |   4 +
 cloud/test/meta_service_test.cpp                 |  17 ++
 cloud/test/recycle_versioned_keys_test.cpp       |   7 +
 cloud/test/recycler_operation_log_test.cpp       |   4 +-
 cloud/test/recycler_test.cpp                     | 206 +++++++++++++++++++-
 gensrc/proto/cloud.proto                         |  11 ++
 14 files changed, 665 insertions(+), 43 deletions(-)

diff --git a/cloud/src/common/http_helper.cpp b/cloud/src/common/http_helper.cpp
index e0ea14eb268..4e231894ca0 100644
--- a/cloud/src/common/http_helper.cpp
+++ b/cloud/src/common/http_helper.cpp
@@ -370,7 +370,12 @@ const std::unordered_map<std::string_view, 
HttpHandlerInfo>& get_http_handlers()
                               return process_query_rate_limit((MS*)s, c);
                           },
                   .role = HttpRole::META_SERVICE}},
-
+                {"check_instance_recycle_completed",
+                 {.handler =
+                          [](void* s, brpc::Controller* c) {
+                              return 
process_check_instance_recycle_completed((MS*)s, c);
+                          },
+                  .role = HttpRole::META_SERVICE}},
                 // Recycler APIs
                 {"recycle_instance",
                  {.handler =
@@ -400,6 +405,12 @@ const std::unordered_map<std::string_view, 
HttpHandlerInfo>& get_http_handlers()
                  {.handler = [](void* s,
                                 brpc::Controller* c) { return 
process_check_instance((RS*)s, c); },
                   .role = HttpRole::RECYCLER}},
+                {"skip_instance_data_cleanup",
+                 {.handler =
+                          [](void* s, brpc::Controller* c) {
+                              return 
process_skip_instance_data_cleanup((RS*)s, c);
+                          },
+                  .role = HttpRole::RECYCLER}},
                 {"check_job_info",
                  {.handler = [](void* s,
                                 brpc::Controller* c) { return 
process_check_job_info((RS*)s, c); },
@@ -561,6 +572,17 @@ HttpResponse process_alter_instance(MetaServiceImpl* 
service, brpc::Controller*
     return http_json_reply(resp.status());
 }
 
+HttpResponse process_skip_instance_data_cleanup(RecyclerServiceImpl* service,
+                                                brpc::Controller* ctrl) {
+    auto& uri = ctrl->http_request().uri();
+    std::string instance_id(http_query(uri, "instance_id"));
+    if (instance_id.empty()) {
+        return http_json_reply(MetaServiceCode::INVALID_ARGUMENT, "instance_id 
is empty");
+    }
+    auto [code, msg] = service->skip_instance_data_cleanup(instance_id);
+    return http_json_reply(code, msg);
+}
+
 HttpResponse process_abort_txn(MetaServiceImpl* service, brpc::Controller* 
ctrl) {
     AbortTxnRequest req;
     PARSE_MESSAGE_OR_RETURN(ctrl, req);
@@ -701,6 +723,32 @@ HttpResponse process_query_rate_limit(MetaServiceImpl* 
service, brpc::Controller
     return http_json_reply(MetaServiceCode::OK, "", sb.GetString());
 }
 
+HttpResponse process_check_instance_recycle_completed(MetaServiceImpl* service,
+                                                      brpc::Controller* cntl) {
+    const auto* instance_id = 
cntl->http_request().uri().GetQuery("instance_id");
+    if (!instance_id || instance_id->empty()) {
+        return http_json_reply(MetaServiceCode::INVALID_ARGUMENT, "no instance 
id");
+    }
+
+    bool finished = false;
+    std::string reason;
+    auto [code, msg] = service->check_instance_recycle_completed(*instance_id, 
finished, reason);
+    if (code != MetaServiceCode::OK) {
+        return http_json_reply(code, msg);
+    }
+
+    rapidjson::Document result;
+    result.SetObject();
+    result.AddMember("finished", finished, result.GetAllocator());
+    result.AddMember("reason",
+                     rapidjson::Value(reason.data(), reason.size(), 
result.GetAllocator()),
+                     result.GetAllocator());
+    rapidjson::StringBuffer buffer;
+    rapidjson::Writer<rapidjson::StringBuffer> writer(buffer);
+    result.Accept(writer);
+    return http_json_reply(code, msg, buffer.GetString());
+}
+
 // Recycler HTTP handlers
 HttpResponse process_recycle_instance(RecyclerServiceImpl* service, 
brpc::Controller* cntl) {
     std::string request_body = cntl->request_attachment().to_string();
diff --git a/cloud/src/common/http_helper.h b/cloud/src/common/http_helper.h
index 7224510bc1d..a0630363749 100644
--- a/cloud/src/common/http_helper.h
+++ b/cloud/src/common/http_helper.h
@@ -124,6 +124,9 @@ const std::unordered_map<std::string_view, 
HttpHandlerInfo>& get_http_handlers()
 [[maybe_unused]] HttpResponse process_query_rate_limit(MetaServiceImpl* 
service,
                                                        brpc::Controller* cntl);
 
+[[maybe_unused]] HttpResponse 
process_check_instance_recycle_completed(MetaServiceImpl* service,
+                                                                       
brpc::Controller* cntl);
+
 [[maybe_unused]] HttpResponse process_decode_key(MetaServiceImpl*, 
brpc::Controller* ctrl);
 
 [[maybe_unused]] HttpResponse process_encode_key(MetaServiceImpl*, 
brpc::Controller* ctrl);
@@ -200,6 +203,9 @@ const std::unordered_map<std::string_view, 
HttpHandlerInfo>& get_http_handlers()
 [[maybe_unused]] HttpResponse process_check_instance(RecyclerServiceImpl* 
service,
                                                      brpc::Controller* cntl);
 
+[[maybe_unused]] HttpResponse 
process_skip_instance_data_cleanup(RecyclerServiceImpl* service,
+                                                                 
brpc::Controller* ctrl);
+
 [[maybe_unused]] HttpResponse process_check_job_info(RecyclerServiceImpl* 
service,
                                                      brpc::Controller* cntl);
 
diff --git a/cloud/src/meta-service/meta_service.h 
b/cloud/src/meta-service/meta_service.h
index a207a71cd0f..26b265f1d05 100644
--- a/cloud/src/meta-service/meta_service.h
+++ b/cloud/src/meta-service/meta_service.h
@@ -435,6 +435,9 @@ public:
                           const CompactSnapshotRequest* request, 
CompactSnapshotResponse* response,
                           ::google::protobuf::Closure* done) override;
 
+    std::pair<MetaServiceCode, std::string> check_instance_recycle_completed(
+            const std::string& instance_id, bool& finished, std::string& 
reason);
+
 private:
     std::pair<MetaServiceCode, std::string> alter_instance(
             const AlterInstanceRequest* request,
diff --git a/cloud/src/meta-service/meta_service_resource.cpp 
b/cloud/src/meta-service/meta_service_resource.cpp
index 2e19ebbb8a3..ecdb13a63c1 100644
--- a/cloud/src/meta-service/meta_service_resource.cpp
+++ b/cloud/src/meta-service/meta_service_resource.cpp
@@ -2562,6 +2562,11 @@ static std::pair<MetaServiceCode, std::string> 
drop_single_instance(const std::s
 
     instance->set_status(InstanceInfoPB::DELETED);
     
instance->set_mtime(duration_cast<seconds>(system_clock::now().time_since_epoch()).count());
+    if (!instance->has_recycle_state()) {
+        
instance->set_recycle_state(INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING);
+        instance->set_recycle_state_update_time_ms(
+                
duration_cast<milliseconds>(system_clock::now().time_since_epoch()).count());
+    }
 
     std::string serialized = instance->SerializeAsString();
     if (serialized.empty()) {
@@ -2634,6 +2639,11 @@ static std::pair<MetaServiceCode, std::string> 
drop_instance_chain(
     for (auto& instance : predecessors) {
         instance.set_status(InstanceInfoPB::DELETED);
         instance.set_mtime(now);
+        if (!instance.has_recycle_state()) {
+            
instance.set_recycle_state(INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING);
+            instance.set_recycle_state_update_time_ms(
+                    
duration_cast<milliseconds>(system_clock::now().time_since_epoch()).count());
+        }
         std::string serialized;
         if (!instance.SerializeToString(&serialized)) {
             std::string msg =
@@ -2647,6 +2657,11 @@ static std::pair<MetaServiceCode, std::string> 
drop_instance_chain(
 
     tail_instance->set_status(InstanceInfoPB::DELETED);
     tail_instance->set_mtime(now);
+    if (!tail_instance->has_recycle_state()) {
+        
tail_instance->set_recycle_state(INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING);
+        tail_instance->set_recycle_state_update_time_ms(
+                
duration_cast<milliseconds>(system_clock::now().time_since_epoch()).count());
+    }
     std::string serialized = tail_instance->SerializeAsString();
     if (serialized.empty()) {
         std::string msg = "failed to serialize";
@@ -2658,6 +2673,55 @@ static std::pair<MetaServiceCode, std::string> 
drop_instance_chain(
     return {MetaServiceCode::OK, std::move(serialized)};
 }
 
+std::pair<MetaServiceCode, std::string> 
MetaServiceImpl::check_instance_recycle_completed(
+        const std::string& instance_id, bool& finished, std::string& reason) {
+    reason.clear();
+    std::unique_ptr<Transaction> txn;
+    TxnErrorCode err = txn_kv_->create_txn(&txn);
+    if (err != TxnErrorCode::TXN_OK) {
+        std::string msg = fmt::format("failed to create txn, err={}", err);
+        LOG(WARNING) << msg << " instance_id=" << instance_id;
+        return {MetaServiceCode::KV_TXN_CREATE_ERR, std::move(msg)};
+    }
+
+    std::string key = instance_key({instance_id});
+    std::string value;
+    err = txn->get(key, &value);
+    if (err == TxnErrorCode::TXN_KEY_NOT_FOUND) {
+        // The recycler only removes keys after cleanup is complete, so do not 
treat this as an
+        // incomplete or unknown state.
+        finished = true;
+        reason = fmt::format(
+                "instance recycling is considered completed because the 
instance key does not "
+                "exist, instance_id={}",
+                instance_id);
+        return {MetaServiceCode::OK, "OK"};
+    }
+    if (err != TxnErrorCode::TXN_OK) {
+        std::string msg =
+                fmt::format("failed to get instance, instance_id={}, err={}", 
instance_id, err);
+        LOG(WARNING) << msg;
+        return {MetaServiceCode::KV_TXN_GET_ERR, std::move(msg)};
+    }
+
+    InstanceInfoPB instance;
+    if (!instance.ParseFromString(value)) {
+        std::string msg = fmt::format("malformed instance info, key={}", 
hex(key));
+        LOG(WARNING) << msg;
+        return {MetaServiceCode::PROTOBUF_PARSE_ERR, std::move(msg)};
+    }
+
+    finished = instance.recycle_state() ==
+               InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED;
+    if (!finished) {
+        reason = fmt::format(
+                "instance has not completed recycling, instance_id={}, 
recycle_state={}",
+                instance_id, 
InstanceRecycleState_Name(instance.recycle_state()));
+        return {MetaServiceCode::OK, "OK"};
+    }
+    return {MetaServiceCode::OK, "OK"};
+}
+
 void MetaServiceImpl::alter_instance(google::protobuf::RpcController* 
controller,
                                      const AlterInstanceRequest* request,
                                      AlterInstanceResponse* response,
@@ -2686,6 +2750,12 @@ void 
MetaServiceImpl::alter_instance(google::protobuf::RpcController* controller
     switch (request->op()) {
     case AlterInstanceRequest::DROP: {
         ret = alter_instance(request, [&instance_id](Transaction* txn, 
InstanceInfoPB* instance) {
+            if (instance->status() == InstanceInfoPB::DELETED) {
+                std::string msg = "instance has already been recycled";
+                LOG(WARNING) << msg << " instance_id=" << instance_id;
+                return std::make_pair(MetaServiceCode::OK, 
instance->SerializeAsString());
+            }
+
             // check instance doesn't have any cluster.
             if (instance->clusters_size() != 0) {
                 std::string msg = "failed to drop instance, instance has 
clusters";
diff --git a/cloud/src/recycler/recycler.cpp b/cloud/src/recycler/recycler.cpp
index fa4cf2538ee..011062b1159 100644
--- a/cloud/src/recycler/recycler.cpp
+++ b/cloud/src/recycler/recycler.cpp
@@ -757,6 +757,12 @@ int InstanceRecycler::init_storage_vault_accessors() {
 }
 
 int InstanceRecycler::init() {
+    if (instance_info_.status() == InstanceInfoPB::DELETED &&
+        (instance_info_.recycle_state() == 
INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING ||
+         instance_info_.recycle_state() == 
INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED)) {
+        return 0;
+    }
+
     int ret = init_obj_store_accessors();
     if (ret != 0) {
         return ret;
@@ -847,14 +853,49 @@ int InstanceRecycler::recycle_deleted_instance() {
 
     int ret = 0;
     auto start_time = steady_clock::now();
+    const auto recycle_state = instance_info_.recycle_state();
 
     DORIS_CLOUD_DEFER {
         auto cost = duration<float>(steady_clock::now() - start_time).count();
-        LOG(WARNING) << (ret == 0 ? "successfully" : "failed to")
-                     << " recycle deleted instance, cost=" << cost
-                     << "s, instance_id=" << instance_id_;
+        if (ret != 0) {
+            LOG(WARNING) << "failed to recycle deleted instance, 
recycle_state="
+                         << InstanceRecycleState_Name(recycle_state) << ", 
cost=" << cost
+                         << "s, instance_id=" << instance_id_;
+        } else if (recycle_state ==
+                   
InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED) {
+            LOG(INFO) << "successfully removed recycled instance key, cost=" 
<< cost
+                      << "s, instance_id=" << instance_id_;
+        } else {
+            LOG(INFO) << "finished recycle deleted instance step, 
recycle_state="
+                      << InstanceRecycleState_Name(recycle_state) << ", cost=" 
<< cost
+                      << "s, instance_id=" << instance_id_;
+        }
     };
 
+    switch (recycle_state) {
+    case InstanceRecycleState::INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING:
+        ret = recycle_deleted_instance_data();
+        break;
+    case InstanceRecycleState::INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING:
+        ret = recycle_deleted_instance_metadata();
+        break;
+    case InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED:
+        ret = remove_instance_key();
+        break;
+    default:
+        LOG_WARNING("invalid instance recycle state")
+                .tag("instance_id", instance_id_)
+                .tag("recycle_state", instance_info_.recycle_state());
+        ret = -1;
+        break;
+    }
+
+    return ret;
+}
+
+int InstanceRecycler::recycle_deleted_instance_data() {
+    int ret = 0;
+
     // Step 1: Recycle tmp rowsets (contains ref count but txn is not 
committed)
     auto recycle_tmp_rowsets_with_mark_delete_enabled = [&]() -> int {
         int res = recycle_tmp_rowsets();
@@ -867,23 +908,21 @@ int InstanceRecycler::recycle_deleted_instance() {
         }
         return res;
     };
+
     if (recycle_tmp_rowsets_with_mark_delete_enabled() != 0) {
         LOG_WARNING("failed to recycle tmp rowsets").tag("instance_id", 
instance_id_);
-        ret = -1;
         return -1;
     }
 
     // Step 2: Recycle versioned rowsets in recycle space (already marked for 
deletion)
     if (recycle_versioned_rowsets() != 0) {
         LOG_WARNING("failed to recycle versioned rowsets").tag("instance_id", 
instance_id_);
-        ret = -1;
         return -1;
     }
 
     // Step 3: Recycle operation logs (can recycle logs not referenced by 
snapshots)
     if (recycle_operation_logs() != 0) {
         LOG_WARNING("failed to recycle operation logs").tag("instance_id", 
instance_id_);
-        ret = -1;
         return -1;
     }
 
@@ -891,7 +930,6 @@ int InstanceRecycler::recycle_deleted_instance() {
     bool has_snapshots = false;
     if (has_cluster_snapshots(&has_snapshots) != 0) {
         LOG(WARNING) << "check instance cluster snapshots failed, 
instance_id=" << instance_id_;
-        ret = -1;
         return -1;
     } else if (has_snapshots) {
         LOG(INFO) << "instance has cluster snapshots, skip recycling, 
instance_id=" << instance_id_;
@@ -905,7 +943,6 @@ int InstanceRecycler::recycle_deleted_instance() {
         bool has_unrecycled_rowsets = false;
         if (recycle_ref_rowsets(&has_unrecycled_rowsets) != 0) {
             LOG_WARNING("failed to recycle ref rowsets").tag("instance_id", 
instance_id_);
-            ret = -1;
             return -1;
         } else if (has_unrecycled_rowsets) {
             LOG_INFO("instance has referenced rowsets, skip recycling")
@@ -935,6 +972,16 @@ int InstanceRecycler::recycle_deleted_instance() {
         }
     }
 
+    if (update_instance_recycle_state(
+                
InstanceRecycleState::INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING,
+                
InstanceRecycleState::INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING) != 0) {
+        return -1;
+    }
+
+    return 0;
+}
+
+int InstanceRecycler::recycle_deleted_instance_metadata() {
     // Check successor instance, if exists, skip deleting kv because successor 
instance may still need the data in kv
     if (instance_info_.has_successor_instance_id() &&
         !instance_info_.successor_instance_id().empty()) {
@@ -945,7 +992,6 @@ int InstanceRecycler::recycle_deleted_instance() {
             LOG(WARNING) << "failed to create txn, instance_id=" << 
instance_id_
                          << " successor_instance_id=" << 
instance_info_.successor_instance_id()
                          << " err=" << err;
-            ret = -1;
             return -1;
         }
 
@@ -960,7 +1006,6 @@ int InstanceRecycler::recycle_deleted_instance() {
             LOG(WARNING) << "failed to get successor instance, instance_id=" 
<< instance_id_
                          << " successor_instance_id=" << 
instance_info_.successor_instance_id()
                          << " err=" << err;
-            ret = -1;
             return -1;
         }
     }
@@ -970,7 +1015,6 @@ int InstanceRecycler::recycle_deleted_instance() {
     TxnErrorCode err = txn_kv_->create_txn(&txn);
     if (err != TxnErrorCode::TXN_OK) {
         LOG(WARNING) << "failed to create txn";
-        ret = -1;
         return -1;
     }
     LOG(INFO) << "begin to delete all kv, instance_id=" << instance_id_;
@@ -1021,33 +1065,157 @@ int InstanceRecycler::recycle_deleted_instance() {
     std::string versioned_log_key_start = 
versioned::log_key_prefix(instance_id_);
     std::string versioned_log_key_end = versioned::log_key_prefix(instance_id_ 
+ '\x00');
     txn->remove(versioned_log_key_start, versioned_log_key_end);
+
+    // Updating the recycle state also commits this transaction, making the 
metadata deletions
+    // and state transition atomic.
+    if (update_instance_recycle_state(
+                
InstanceRecycleState::INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING,
+                
InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED, txn.get()) != 
0) {
+        return -1;
+    }
+
+    return 0;
+}
+
+int InstanceRecycler::remove_instance_key() {
+    std::unique_ptr<Transaction> txn;
+    std::string key = instance_key(instance_info_.instance_id());
+    std::string value;
+    TxnErrorCode err = txn_kv_->create_txn(&txn);
+    if (err != TxnErrorCode::TXN_OK) {
+        LOG(WARNING) << "failed to create txn";
+        return -1;
+    }
+
+    err = txn->get(key, &value);
+    if (err != TxnErrorCode::TXN_OK) {
+        LOG(WARNING) << "failed to get instance, instance_id=" << 
instance_info_.instance_id()
+                     << ", err=" << err;
+        return -1;
+    }
+
+    InstanceInfoPB instance;
+    if (!instance.ParseFromString(value)) {
+        LOG(WARNING) << "malformed instance info, key=" << key;
+        return -1;
+    }
+
+    if (instance.status() != InstanceInfoPB::DELETED) {
+        LOG(WARNING) << "failed to remove instance key, instance is not 
deleted, instance_id="
+                     << instance_id_ << ", status=" << instance.status();
+        return -1;
+    }
+
+    if (instance.recycle_state() != INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED) {
+        LOG(WARNING) << "failed to remove instance key, invalid recycle state, 
instance_id="
+                     << instance_id_
+                     << ", current_state=" << 
InstanceRecycleState_Name(instance.recycle_state())
+                     << ", expected_state="
+                     << 
InstanceRecycleState_Name(INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED);
+        return -1;
+    }
+
+    txn->atomic_add(system_meta_service_instance_update_key(), 1);
+    txn->remove(key);
     err = txn->commit();
     if (err != TxnErrorCode::TXN_OK) {
-        LOG(WARNING) << "failed to delete all kv, instance_id=" << 
instance_id_ << ", err=" << err;
-        ret = -1;
+        LOG(WARNING) << "failed to delete instance kv, instance_id=" << 
instance_id_
+                     << " err=" << err;
+        return -1;
     }
+    return 0;
+}
 
-    if (ret == 0) {
-        // remove instance kv
-        // ATTN: MUST ensure that cloud platform won't regenerate the same 
instance id
-        err = txn_kv_->create_txn(&txn);
-        if (err != TxnErrorCode::TXN_OK) {
-            LOG(WARNING) << "failed to create txn";
-            ret = -1;
-            return ret;
-        }
-        std::string key;
-        instance_key({instance_id_}, &key);
-        txn->atomic_add(system_meta_service_instance_update_key(), 1);
-        txn->remove(key);
-        err = txn->commit();
-        if (err != TxnErrorCode::TXN_OK) {
-            LOG(WARNING) << "failed to delete instance kv, instance_id=" << 
instance_id_
-                         << " err=" << err;
-            ret = -1;
-        }
+int InstanceRecycler::update_instance_recycle_state(InstanceRecycleState 
expected_state,
+                                                    InstanceRecycleState 
target_state) {
+    std::unique_ptr<Transaction> txn;
+    TxnErrorCode err = txn_kv_->create_txn(&txn);
+    if (err != TxnErrorCode::TXN_OK) {
+        LOG(WARNING) << "failed to create txn";
+        return -1;
     }
-    return ret;
+    return update_instance_recycle_state(expected_state, target_state, 
txn.get());
+}
+
+int InstanceRecycler::update_instance_recycle_state(InstanceRecycleState 
current_state,
+                                                    InstanceRecycleState 
target_state,
+                                                    Transaction* txn) {
+    const bool valid_transition =
+            (current_state == INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING &&
+             target_state == INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING) 
||
+            (current_state == INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING 
&&
+             target_state == INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED);
+
+    if (!valid_transition) {
+        LOG_WARNING("invalid instance recycled state transition")
+                .tag("instance_id", instance_id_)
+                .tag("current_state", InstanceRecycleState_Name(current_state))
+                .tag("target_state", InstanceRecycleState_Name(target_state));
+        return -1;
+    }
+
+    std::string key = instance_key({instance_id_});
+    std::string value;
+    TxnErrorCode err = txn->get(key, &value);
+    if (err != TxnErrorCode::TXN_OK) {
+        LOG_WARNING("failed to get instance when updating instance recycled 
state")
+                .tag("instance_id", instance_id_)
+                .tag("current_state", InstanceRecycleState_Name(current_state))
+                .tag("target_state", InstanceRecycleState_Name(target_state))
+                .tag("err", err);
+        return -1;
+    }
+
+    InstanceInfoPB instance;
+    if (!instance.ParseFromString(value)) {
+        LOG_WARNING("failed to parse InstanceInfoPB when updating instance 
recycled state")
+                .tag("instance_id", instance_id_)
+                .tag("current_state", InstanceRecycleState_Name(current_state))
+                .tag("target_state", InstanceRecycleState_Name(target_state));
+        return -1;
+    }
+    if (instance.status() != InstanceInfoPB::DELETED) {
+        LOG_WARNING("instance is not deleted when updating instance recycled 
state")
+                .tag("instance_id", instance_id_)
+                .tag("status", instance.status())
+                .tag("current_state", InstanceRecycleState_Name(current_state))
+                .tag("target_state", InstanceRecycleState_Name(target_state));
+        return -1;
+    }
+    if (instance.recycle_state() != current_state) {
+        LOG_WARNING("instance recycled state changed before update")
+                .tag("instance_id", instance_id_)
+                .tag("current_state", 
InstanceRecycleState_Name(instance.recycle_state()))
+                .tag("expected_state", 
InstanceRecycleState_Name(current_state))
+                .tag("target_state", InstanceRecycleState_Name(target_state));
+        return -1;
+    }
+
+    instance.set_recycle_state(target_state);
+    instance.set_recycle_state_update_time_ms(
+            
duration_cast<milliseconds>(system_clock::now().time_since_epoch()).count());
+    if (!instance.SerializeToString(&value)) {
+        LOG_WARNING("failed to serialize InstanceInfoPB when updating instance 
recycled state")
+                .tag("instance_id", instance_id_)
+                .tag("expected_state", 
InstanceRecycleState_Name(current_state))
+                .tag("target_state", InstanceRecycleState_Name(target_state));
+        return -1;
+    }
+
+    txn->atomic_add(system_meta_service_instance_update_key(), 1);
+    txn->put(key, value);
+    err = txn->commit();
+    if (err != TxnErrorCode::TXN_OK) {
+        LOG(WARNING) << "failed to commit fdb txn when updating instance 
recycled state, "
+                     << "instance_id=" << instance_id_ << ", err=" << err;
+        return -1;
+    }
+
+    instance_info_.Swap(&instance);
+    LOG_INFO("updated instance recycled state")
+            .tag("instance_id", instance_id_)
+            .tag("recycle_state", InstanceRecycleState_Name(target_state));
+    return 0;
 }
 
 int InstanceRecycler::check_rowset_exists(int64_t tablet_id, const 
std::string& rowset_id,
diff --git a/cloud/src/recycler/recycler.h b/cloud/src/recycler/recycler.h
index 32f5c638b5e..8d31f69151a 100644
--- a/cloud/src/recycler/recycler.h
+++ b/cloud/src/recycler/recycler.h
@@ -288,6 +288,16 @@ public:
     // returns 0 for success otherwise error
     int recycle_deleted_instance();
 
+    int recycle_deleted_instance_data();
+
+    int recycle_deleted_instance_metadata();
+
+    int update_instance_recycle_state(InstanceRecycleState expected_state,
+                                      InstanceRecycleState target_state);
+
+    int update_instance_recycle_state(InstanceRecycleState expected_state,
+                                      InstanceRecycleState target_state, 
Transaction* txn);
+
     // scan and recycle expired indexes:
     // 1. dropped table, dropped mv
     // 2. half-successtable/index when create
@@ -430,6 +440,9 @@ public:
                                        const SnapshotPB& snapshot_pb);
 
 private:
+    // returns 0 for success otherwise error
+    int remove_instance_key();
+
     // returns 0 for success otherwise error
     int init_obj_store_accessors();
 
diff --git a/cloud/src/recycler/recycler_service.cpp 
b/cloud/src/recycler/recycler_service.cpp
index b916d751925..c9c68004f8a 100644
--- a/cloud/src/recycler/recycler_service.cpp
+++ b/cloud/src/recycler/recycler_service.cpp
@@ -27,6 +27,7 @@
 #include <rapidjson/stringbuffer.h>
 
 #include <algorithm>
+#include <chrono>
 #include <functional>
 #include <memory>
 #include <numeric>
@@ -398,6 +399,85 @@ void RecyclerServiceImpl::check_instance(const 
std::string& instance_id, MetaSer
     }
 }
 
+std::pair<MetaServiceCode, std::string> 
RecyclerServiceImpl::skip_instance_data_cleanup(
+        const std::string& instance_id) {
+    std::unique_ptr<Transaction> txn;
+    TxnErrorCode err = txn_kv_->create_txn(&txn);
+    if (err != TxnErrorCode::TXN_OK) {
+        std::string msg = fmt::format("failed to create txn, err={}", err);
+        LOG(WARNING) << msg << " instance_id=" << instance_id;
+        return {MetaServiceCode::KV_TXN_CREATE_ERR, std::move(msg)};
+    }
+
+    std::string key = instance_key({instance_id});
+    std::string value;
+    err = txn->get(key, &value);
+    if (err != TxnErrorCode::TXN_OK) {
+        std::string msg =
+                fmt::format("failed to get instance, instance_id={}, err={}", 
instance_id, err);
+        LOG(WARNING) << msg;
+        return {MetaServiceCode::KV_TXN_GET_ERR, std::move(msg)};
+    }
+
+    InstanceInfoPB instance;
+    if (!instance.ParseFromString(value)) {
+        std::string msg = fmt::format("malformed instance info, key={}", 
hex(key));
+        LOG(WARNING) << msg;
+        return {MetaServiceCode::PROTOBUF_PARSE_ERR, std::move(msg)};
+    }
+    auto current_state = instance.recycle_state();
+    if (instance.status() != InstanceInfoPB::DELETED) {
+        std::string msg = fmt::format(
+                "failed to set instance recycle state, instance is not 
deleted, instance_id={}",
+                instance_id);
+        LOG(WARNING) << msg;
+        return {MetaServiceCode::INVALID_ARGUMENT, std::move(msg)};
+    }
+    if (current_state != INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING) {
+        std::string msg = fmt::format(
+                "failed to set instance recycle state, instance state should 
be {}"
+                ", current_state={}"
+                ", instance_id={}",
+                INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING, current_state, 
instance_id);
+        LOG(WARNING) << msg;
+        return {MetaServiceCode::INVALID_ARGUMENT, std::move(msg)};
+    }
+    if (instance.has_multi_version_status() &&
+        instance.multi_version_status() != 
MultiVersionStatus::MULTI_VERSION_DISABLED) {
+        std::string msg = fmt::format(
+                "cannot skip instance data cleanup for a multi-version 
instance, instance_id={}, "
+                "multi_version_status={}",
+                instance_id, 
MultiVersionStatus_Name(instance.multi_version_status()));
+        LOG(WARNING) << msg;
+        return {MetaServiceCode::INVALID_ARGUMENT, std::move(msg)};
+    }
+
+    
instance.set_recycle_state(INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING);
+    instance.set_recycle_state_update_time_ms(
+            std::chrono::duration_cast<std::chrono::milliseconds>(
+                    std::chrono::system_clock::now().time_since_epoch())
+                    .count());
+    if (!instance.SerializeToString(&value)) {
+        std::string msg = "failed to serialize InstanceInfoPB";
+        LOG(WARNING) << msg << " instance_id=" << instance_id;
+        return {MetaServiceCode::PROTOBUF_SERIALIZE_ERR, std::move(msg)};
+    }
+
+    txn->atomic_add(system_meta_service_instance_update_key(), 1);
+    txn->put(key, value);
+    err = txn->commit();
+    if (err != TxnErrorCode::TXN_OK) {
+        std::string msg = fmt::format("failed to commit kv txn, err={}", err);
+        LOG(WARNING) << msg << " instance_id=" << instance_id;
+        return {MetaServiceCode::KV_TXN_COMMIT_ERR, std::move(msg)};
+    }
+    LOG(WARNING) << "successfully skipped instance data cleanup, instance_id=" 
<< instance_id
+                 << " current_state=" << 
InstanceRecycleState_Name(current_state)
+                 << " target_state=" << 
InstanceRecycleState_Name(instance.recycle_state());
+
+    return {MetaServiceCode::OK, "OK"};
+}
+
 void recycle_copy_jobs(const std::shared_ptr<TxnKv>& txn_kv, const 
std::string& instance_id,
                        MetaServiceCode& code, std::string& msg,
                        RecyclerThreadPoolGroup thread_pool_group,
diff --git a/cloud/src/recycler/recycler_service.h 
b/cloud/src/recycler/recycler_service.h
index 3441fa7303c..01849dd8b08 100644
--- a/cloud/src/recycler/recycler_service.h
+++ b/cloud/src/recycler/recycler_service.h
@@ -50,6 +50,9 @@ public:
 
     void check_instance(const std::string& instance_id, MetaServiceCode& code, 
std::string& msg);
 
+    std::pair<MetaServiceCode, std::string> skip_instance_data_cleanup(
+            const std::string& instance_id);
+
     std::shared_ptr<TxnKv> txn_kv() { return txn_kv_; }
     Recycler* recycler() { return recycler_; }
     Checker* checker() { return checker_; }
diff --git a/cloud/src/recycler/recycler_snapshot.cpp 
b/cloud/src/recycler/recycler_snapshot.cpp
index 6b8203dbc10..856fdfc9797 100644
--- a/cloud/src/recycler/recycler_snapshot.cpp
+++ b/cloud/src/recycler/recycler_snapshot.cpp
@@ -40,6 +40,10 @@
 namespace doris::cloud {
 
 int InstanceRecycler::recycle_cluster_snapshots() {
+    if (instance_info_.recycle_state() == 
INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING ||
+        instance_info_.recycle_state() == 
INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED) {
+        return 0;
+    }
     return snapshot_manager_->recycle_snapshots(this);
 }
 
diff --git a/cloud/test/meta_service_test.cpp b/cloud/test/meta_service_test.cpp
index c74dd1655bf..e6a1832fa3b 100644
--- a/cloud/test/meta_service_test.cpp
+++ b/cloud/test/meta_service_test.cpp
@@ -591,6 +591,23 @@ TEST(MetaServiceTest, CreateInstanceTest) {
         instance.ParseFromString(val);
         ASSERT_EQ(instance.status(), InstanceInfoPB::DELETED);
         ASSERT_EQ(res.status().code(), MetaServiceCode::OK);
+
+        
instance.set_recycle_state(InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED);
+        txn->put(key, instance.SerializeAsString());
+        ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+
+        AlterInstanceResponse retry_res;
+        
meta_service->alter_instance(reinterpret_cast<::google::protobuf::RpcController*>(&cntl),
+                                     &req, &retry_res, nullptr);
+        ASSERT_EQ(retry_res.status().code(), MetaServiceCode::OK);
+        ASSERT_EQ(retry_res.status().msg().find("instance has already been 
recycled"),
+                  std::string::npos);
+        val.clear();
+        ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), 
TxnErrorCode::TXN_OK);
+        ASSERT_EQ(txn->get(key, &val), TxnErrorCode::TXN_OK);
+        instance.ParseFromString(val);
+        ASSERT_EQ(instance.status(), InstanceInfoPB::DELETED);
+        ASSERT_EQ(instance.recycle_state(), 
INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED);
     }
 
     // case: normal refresh instance
diff --git a/cloud/test/recycle_versioned_keys_test.cpp 
b/cloud/test/recycle_versioned_keys_test.cpp
index 670cccfac4a..00b99573fad 100644
--- a/cloud/test/recycle_versioned_keys_test.cpp
+++ b/cloud/test/recycle_versioned_keys_test.cpp
@@ -1859,6 +1859,8 @@ TEST(RecycleVersionedKeysTest, RecycleDeletedInstance) {
     {
         // Recycle deleted instance
         ASSERT_EQ(recycler.recycle_deleted_instance(), 0);
+        ASSERT_EQ(recycler.recycle_deleted_instance(), 0);
+        ASSERT_EQ(recycler.recycle_deleted_instance(), 0);
     }
 
     {
@@ -1890,6 +1892,11 @@ TEST(RecycleVersionedKeysTest, RecycleDeletedInstance) {
         std::string log_key_end = versioned::log_key_prefix(instance_id + 
'\x00');
         ASSERT_EQ(count_range(txn_kv.get(), log_key, log_key_end), 0) << 
dump_range(txn_kv.get());
 
+        std::string instance_key_st = instance_key(instance_id);
+        std::string instance_key_ed = instance_key(instance_id + '\x00');
+        ASSERT_EQ(count_range(txn_kv.get(), instance_key_st, instance_key_ed), 
0)
+                << dump_range(txn_kv.get());
+
         for (int i = 1; i < rowsets.size(); ++i) {
             std::unique_ptr<ListIterator> list_iter;
             ASSERT_EQ(0, accessor->list_directory(
diff --git a/cloud/test/recycler_operation_log_test.cpp 
b/cloud/test/recycler_operation_log_test.cpp
index b3f9c385fde..becac15276d 100644
--- a/cloud/test/recycler_operation_log_test.cpp
+++ b/cloud/test/recycler_operation_log_test.cpp
@@ -2577,9 +2577,11 @@ TEST(RecycleOperationLogTest, RecycleDeletedInstance) {
         update_instance_info(txn_kv.get(), deleted_instance);
     }
 
+    ASSERT_EQ(recycler.recycle_deleted_instance(), 0);
+    ASSERT_EQ(recycler.recycle_deleted_instance(), 0);
     ASSERT_EQ(recycler.recycle_deleted_instance(), 0);
 
-    // Verify all keys are deleted, expecting the instance_update
+    // Verify all data keys are deleted, keeping the instance status and 
instance_update keys.
     ASSERT_EQ(count_range(txn_kv.get()), 1) << dump_range(txn_kv.get());
 }
 
diff --git a/cloud/test/recycler_test.cpp b/cloud/test/recycler_test.cpp
index 28ce6b0f716..c1abd0a5e13 100644
--- a/cloud/test/recycler_test.cpp
+++ b/cloud/test/recycler_test.cpp
@@ -42,6 +42,7 @@
 #include "common/bvars.h"
 #include "common/config.h"
 #include "common/defer.h"
+#include "common/http_helper.h"
 #include "common/logging.h"
 #include "common/simple_thread_pool.h"
 #include "common/util.h"
@@ -58,6 +59,7 @@
 #include "rate-limiter/rate_limiter.h"
 #include "recycler/checker.h"
 #include "recycler/recycler.cpp"
+#include "recycler/recycler_service.h"
 #include "recycler/storage_vault_accessor.h"
 #include "recycler/util.h"
 #include "recycler/white_black_list.h"
@@ -1154,6 +1156,13 @@ static int create_instance(const std::string& 
internal_stage_id,
     return 0;
 }
 
+static void put_instance_info(TxnKv* txn_kv, const InstanceInfoPB& 
instance_info) {
+    std::unique_ptr<Transaction> txn;
+    ASSERT_EQ(TxnErrorCode::TXN_OK, txn_kv->create_txn(&txn));
+    txn->put(instance_key({instance_info.instance_id()}), 
instance_info.SerializeAsString());
+    ASSERT_EQ(TxnErrorCode::TXN_OK, txn->commit());
+}
+
 static int create_copy_job(TxnKv* txn_kv, const std::string& stage_id, int64_t 
table_id,
                            StagePB::StageType stage_type, CopyJobPB::JobStatus 
job_status,
                            std::vector<ObjectFilePB> object_files, int64_t 
timeout_time,
@@ -3611,6 +3620,8 @@ TEST(RecyclerTest, recycle_deleted_instance) {
 
     InstanceInfoPB instance_info;
     create_instance(internal_stage_id, external_stage_id, instance_info);
+    instance_info.set_status(InstanceInfoPB::DELETED);
+    put_instance_info(txn_kv.get(), instance_info);
     InstanceRecycler recycler(txn_kv, instance_info, thread_group,
                               std::make_shared<TxnLazyCommitter>(txn_kv));
     ASSERT_EQ(recycler.init(), 0);
@@ -3682,6 +3693,8 @@ TEST(RecyclerTest, recycle_deleted_instance) {
     }
 
     ASSERT_EQ(0, recycler.recycle_deleted_instance());
+    
ASSERT_EQ(InstanceRecycleState::INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING,
+              recycler.instance_info().recycle_state());
 
     {
         // No thing to recycle
@@ -3701,6 +3714,19 @@ TEST(RecyclerTest, recycle_deleted_instance) {
     }
 
     ASSERT_EQ(0, recycler.recycle_deleted_instance());
+    
ASSERT_EQ(InstanceRecycleState::INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING,
+              recycler.instance_info().recycle_state());
+    ASSERT_EQ(0, recycler.recycle_deleted_instance());
+    ASSERT_EQ(InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED,
+              recycler.instance_info().recycle_state());
+    ASSERT_EQ(0, recycler.recycle_deleted_instance());
+
+    {
+        std::unique_ptr<Transaction> txn;
+        ASSERT_EQ(TxnErrorCode::TXN_OK, txn_kv->create_txn(&txn));
+        std::string value;
+        ASSERT_EQ(TxnErrorCode::TXN_KEY_NOT_FOUND, 
txn->get(instance_key({instance_id}), &value));
+    }
 
     // check if all the objects are deleted
     std::for_each(recycler.accessor_map_.begin(), recycler.accessor_map_.end(),
@@ -3764,6 +3790,176 @@ TEST(RecyclerTest, recycle_deleted_instance) {
     }
 }
 
+TEST(RecyclerTest, init_deleted_instance_with_terminal_recycle_state) {
+    auto txn_kv = 
std::dynamic_pointer_cast<TxnKv>(std::make_shared<MemTxnKv>());
+    ASSERT_NE(txn_kv.get(), nullptr);
+    ASSERT_EQ(txn_kv->init(), 0);
+
+    InstanceInfoPB instance_info;
+    instance_info.set_instance_id(instance_id);
+    instance_info.set_status(InstanceInfoPB::DELETED);
+    
instance_info.set_recycle_state(InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED);
+    instance_info.add_resource_ids("deleted_vault");
+    put_instance_info(txn_kv.get(), instance_info);
+
+    InstanceRecycler recycler(txn_kv, instance_info, thread_group,
+                              std::make_shared<TxnLazyCommitter>(txn_kv));
+    ASSERT_EQ(recycler.init(), 0);
+    ASSERT_EQ(recycler.do_recycle(), 0);
+}
+
+TEST(RecyclerTest, recycle_legacy_deleted_instance_without_recycle_state) {
+    auto txn_kv = 
std::dynamic_pointer_cast<TxnKv>(std::make_shared<MemTxnKv>());
+    ASSERT_NE(txn_kv.get(), nullptr);
+    ASSERT_EQ(txn_kv->init(), 0);
+
+    InstanceInfoPB instance_info;
+    instance_info.set_instance_id(instance_id);
+    instance_info.set_status(InstanceInfoPB::DELETED);
+    instance_info.add_obj_info()->set_id("legacy_instance_obj");
+    ASSERT_FALSE(instance_info.has_recycle_state());
+    
ASSERT_EQ(InstanceRecycleState::INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING,
+              instance_info.recycle_state());
+    put_instance_info(txn_kv.get(), instance_info);
+
+    InstanceRecycler recycler(txn_kv, instance_info, thread_group,
+                              std::make_shared<TxnLazyCommitter>(txn_kv));
+    ASSERT_EQ(recycler.init(), 0);
+    ASSERT_EQ(recycler.recycle_deleted_instance(), 0);
+
+    std::unique_ptr<Transaction> txn;
+    ASSERT_EQ(TxnErrorCode::TXN_OK, txn_kv->create_txn(&txn));
+    std::string value;
+    ASSERT_EQ(TxnErrorCode::TXN_OK, txn->get(instance_key({instance_id}), 
&value));
+    ASSERT_TRUE(instance_info.ParseFromString(value));
+    ASSERT_TRUE(instance_info.has_recycle_state());
+    
ASSERT_EQ(InstanceRecycleState::INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING,
+              instance_info.recycle_state());
+}
+
+TEST(RecyclerTest, reject_stale_instance_recycle_state_update) {
+    auto txn_kv = 
std::dynamic_pointer_cast<TxnKv>(std::make_shared<MemTxnKv>());
+    ASSERT_NE(txn_kv.get(), nullptr);
+    ASSERT_EQ(txn_kv->init(), 0);
+
+    InstanceInfoPB stale_instance;
+    stale_instance.set_instance_id(instance_id);
+    stale_instance.set_status(InstanceInfoPB::DELETED);
+    stale_instance.set_recycle_state(
+            InstanceRecycleState::INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING);
+
+    InstanceInfoPB latest_instance = stale_instance;
+    latest_instance.set_recycle_state(
+            InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED);
+    put_instance_info(txn_kv.get(), latest_instance);
+
+    InstanceRecycler stale_recycler(txn_kv, stale_instance, thread_group,
+                                    
std::make_shared<TxnLazyCommitter>(txn_kv));
+    ASSERT_NE(stale_recycler.update_instance_recycle_state(
+                      
InstanceRecycleState::INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING,
+                      
InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED),
+              0);
+    ASSERT_NE(stale_recycler.update_instance_recycle_state(
+                      
InstanceRecycleState::INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING,
+                      
InstanceRecycleState::INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING),
+              0);
+
+    std::unique_ptr<Transaction> txn;
+    ASSERT_EQ(TxnErrorCode::TXN_OK, txn_kv->create_txn(&txn));
+    std::string value;
+    ASSERT_EQ(TxnErrorCode::TXN_OK, txn->get(instance_key({instance_id}), 
&value));
+    ASSERT_TRUE(latest_instance.ParseFromString(value));
+    ASSERT_EQ(InstanceRecycleState::INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED,
+              latest_instance.recycle_state());
+}
+
+TEST(RecyclerServiceTest, skip_instance_data_cleanup) {
+    auto txn_kv = 
std::dynamic_pointer_cast<TxnKv>(std::make_shared<MemTxnKv>());
+    ASSERT_NE(txn_kv.get(), nullptr);
+    ASSERT_EQ(txn_kv->init(), 0);
+
+    const std::string test_instance_id = "skip_instance_data_cleanup_instance";
+    RecyclerServiceImpl service(txn_kv, nullptr, nullptr, nullptr);
+    auto get_instance_info = [&](InstanceInfoPB* instance_info) {
+        std::unique_ptr<Transaction> txn;
+        ASSERT_EQ(TxnErrorCode::TXN_OK, txn_kv->create_txn(&txn));
+        std::string value;
+        ASSERT_EQ(TxnErrorCode::TXN_OK, 
txn->get(instance_key({test_instance_id}), &value));
+        ASSERT_TRUE(instance_info->ParseFromString(value));
+    };
+
+    InstanceInfoPB instance_info;
+    instance_info.set_instance_id(test_instance_id);
+    instance_info.set_status(InstanceInfoPB::DELETED);
+    
instance_info.set_recycle_state(INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING);
+    ASSERT_FALSE(instance_info.has_multi_version_status());
+    put_instance_info(txn_kv.get(), instance_info);
+
+    auto [code, msg] = 
service.skip_instance_data_cleanup(instance_info.instance_id());
+    ASSERT_EQ(code, MetaServiceCode::OK) << msg;
+
+    InstanceInfoPB persisted_instance;
+    get_instance_info(&persisted_instance);
+    ASSERT_FALSE(persisted_instance.has_multi_version_status());
+    ASSERT_EQ(persisted_instance.recycle_state(), 
INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING);
+    ASSERT_GT(persisted_instance.recycle_state_update_time_ms(), 0);
+
+    instance_info.set_multi_version_status(MULTI_VERSION_DISABLED);
+    
instance_info.set_recycle_state(INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING);
+    instance_info.clear_recycle_state_update_time_ms();
+    put_instance_info(txn_kv.get(), instance_info);
+
+    auto call_http = [&](std::string_view query) {
+        brpc::Controller ctrl;
+        ctrl.http_request().uri() = fmt::format("/?{}", query);
+        return process_skip_instance_data_cleanup(&service, &ctrl);
+    };
+    auto response = call_http("");
+    ASSERT_EQ(response.status_code, 400) << response.body;
+    ASSERT_EQ(response.msg, "instance_id is empty");
+
+    response = call_http(fmt::format("instance_id={}", test_instance_id));
+    ASSERT_EQ(response.status_code, 200) << response.body;
+    ASSERT_EQ(response.msg, "OK");
+
+    get_instance_info(&persisted_instance);
+    ASSERT_EQ(persisted_instance.multi_version_status(), 
MULTI_VERSION_DISABLED);
+    ASSERT_EQ(persisted_instance.recycle_state(), 
INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING);
+    ASSERT_GT(persisted_instance.recycle_state_update_time_ms(), 0);
+
+    constexpr std::array unsupported_multi_version_statuses {
+            MULTI_VERSION_WRITE_ONLY,
+            MULTI_VERSION_READ_WRITE,
+            MULTI_VERSION_ENABLED,
+    };
+    for (auto multi_version_status : unsupported_multi_version_statuses) {
+        instance_info.set_multi_version_status(multi_version_status);
+        
instance_info.set_recycle_state(INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING);
+        instance_info.set_recycle_state_update_time_ms(123);
+        put_instance_info(txn_kv.get(), instance_info);
+
+        std::tie(code, msg) = 
service.skip_instance_data_cleanup(instance_info.instance_id());
+        ASSERT_EQ(code, MetaServiceCode::INVALID_ARGUMENT)
+                << MultiVersionStatus_Name(multi_version_status);
+        ASSERT_EQ(msg,
+                  fmt::format("cannot skip instance data cleanup for a 
multi-version instance, "
+                              "instance_id={}, multi_version_status={}",
+                              test_instance_id, 
MultiVersionStatus_Name(multi_version_status)));
+
+        get_instance_info(&persisted_instance);
+        ASSERT_EQ(persisted_instance.multi_version_status(), 
multi_version_status);
+        ASSERT_EQ(persisted_instance.recycle_state(), 
INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING);
+        ASSERT_EQ(persisted_instance.recycle_state_update_time_ms(), 123);
+    }
+
+    instance_info.set_status(InstanceInfoPB::NORMAL);
+    instance_info.clear_multi_version_status();
+    put_instance_info(txn_kv.get(), instance_info);
+    std::tie(code, msg) = 
service.skip_instance_data_cleanup(instance_info.instance_id());
+    ASSERT_EQ(code, MetaServiceCode::INVALID_ARGUMENT);
+    ASSERT_NE(msg.find("instance is not deleted"), std::string::npos) << msg;
+}
+
 // Regression test: if commit_rowset is called but the txn is never committed,
 // the rowset stays in meta_rowset_tmp_key with data_rowset_ref_count_key=1.
 // recycle_deleted_instance() must call recycle_tmp_rowsets() first so that
@@ -3783,20 +3979,14 @@ TEST(RecyclerTest, 
recycle_deleted_instance_with_orphan_tmp_rowset) {
     // Create instance with multi-version read/write and snapshot support
     InstanceInfoPB instance_info;
     instance_info.set_instance_id(instance_id);
+    instance_info.set_status(InstanceInfoPB::DELETED);
     
instance_info.set_multi_version_status(MultiVersionStatus::MULTI_VERSION_READ_WRITE);
     
instance_info.set_snapshot_switch_status(SnapshotSwitchStatus::SNAPSHOT_SWITCH_ON);
     auto* obj_info = instance_info.add_obj_info();
     obj_info->set_id("orphan_tmp_rowset_test");
 
     // Write instance info to FDB (required by 
OperationLogRecycleChecker::init())
-    {
-        std::unique_ptr<Transaction> txn;
-        ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
-        std::string key = instance_key({instance_id});
-        std::string val = instance_info.SerializeAsString();
-        txn->put(key, val);
-        ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
-    }
+    put_instance_info(txn_kv.get(), instance_info);
 
     InstanceRecycler recycler(txn_kv, instance_info, thread_group,
                               std::make_shared<TxnLazyCommitter>(txn_kv));
diff --git a/gensrc/proto/cloud.proto b/gensrc/proto/cloud.proto
index 8dda9b0cf7c..94e4b923d00 100644
--- a/gensrc/proto/cloud.proto
+++ b/gensrc/proto/cloud.proto
@@ -100,12 +100,19 @@ enum KeySetType {
     MULTI_VERSION_META_ROWSET = 21;
 }
 
+enum InstanceRecycleState {
+    INSTANCE_RECYCLE_STATE_DATA_CLEANUP_PENDING = 0;
+    INSTANCE_RECYCLE_STATE_METADATA_CLEANUP_PENDING = 1;
+    INSTANCE_RECYCLE_STATE_CLEANUP_COMPLETED = 2;
+}
+
 message InstanceInfoPB {
     enum Status {
         NORMAL = 0;
         DELETED = 1;
         OVERDUE = 2;
     }
+
     optional string user_id = 1;
     optional string instance_id = 2;
     optional string name = 3;
@@ -152,6 +159,10 @@ message InstanceInfoPB {
     // It is not always same as source_instance_id, because the source 
instance may inherit from another
     // instance during rollback, and the predecessor instance is the real 
instance which to execute the rollback.
     optional string predecessor_instance_id = 124;
+
+    optional InstanceRecycleState recycle_state = 125;
+    // Last time when recycle_state was updated by a recycle step.
+    optional int64 recycle_state_update_time_ms = 126;
 }
 
 message StagePB {


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

Reply via email to