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 af7b9936265 branch-4.1: [fix](cloud) Prevent old clients from decoding 
unknown MetaServiceCode as OK (#64148) (#66353)
af7b9936265 is described below

commit af7b9936265e0a73b4137968f54a1571140385e2
Author: Yixuan Wang <[email protected]>
AuthorDate: Mon Aug 3 10:38:15 2026 +0800

    branch-4.1: [fix](cloud) Prevent old clients from decoding unknown 
MetaServiceCode as OK (#64148) (#66353)
    
    pick: https://github.com/apache/doris/pull/64148
---
 be/src/cloud/cloud_meta_mgr.cpp                    |  25 ++-
 be/src/cloud/cloud_meta_mgr.h                      |   4 +
 be/test/cloud/cloud_meta_mgr_test.cpp              |  23 ++
 cloud/src/meta-service/meta_service.h              |   7 +
 cloud/src/meta-service/meta_service_helper.h       |  40 +++-
 cloud/test/meta_service_helper_test.cpp            | 250 +++++++++++++++++++++
 cloud/test/meta_service_http_test.cpp              |   9 +-
 cloud/test/meta_service_test.cpp                   |   8 +
 .../apache/doris/cloud/rpc/MetaServiceClient.java  |  70 +++++-
 .../doris/cloud/rpc/MetaServiceProxyTest.java      | 170 +++++++++++++-
 gensrc/proto/cloud.proto                           |  23 +-
 11 files changed, 598 insertions(+), 31 deletions(-)

diff --git a/be/src/cloud/cloud_meta_mgr.cpp b/be/src/cloud/cloud_meta_mgr.cpp
index f731186679d..7781297942b 100644
--- a/be/src/cloud/cloud_meta_mgr.cpp
+++ b/be/src/cloud/cloud_meta_mgr.cpp
@@ -153,9 +153,20 @@ Status 
bthread_fork_join(std::vector<std::function<Status()>>&& tasks, int concu
     return Status::OK();
 }
 
+MetaServiceCode get_response_code(const MetaServiceResponseStatus& status) {
+    if (status.has_actual_code() && 
MetaServiceCode_IsValid(status.actual_code())) {
+        return static_cast<MetaServiceCode>(status.actual_code());
+    }
+    return status.code();
+}
+
 namespace {
 constexpr int kBrpcRetryTimes = 3;
 
+void restore_actual_code(MetaServiceResponseStatus* status) {
+    status->set_code(get_response_code(*status));
+}
+
 bvar::LatencyRecorder _get_rowset_latency("doris_cloud_meta_mgr_get_rowset");
 bvar::LatencyRecorder 
g_cloud_commit_txn_resp_redirect_latency("cloud_table_stats_report_latency");
 bvar::Adder<uint64_t> 
g_cloud_meta_mgr_rpc_timeout_count("cloud_meta_mgr_rpc_timeout_count");
@@ -417,6 +428,16 @@ using MetaServiceMethod = void 
(MetaService_Stub::*)(::google::protobuf::RpcCont
                                                      const Request*, Response*,
                                                      
::google::protobuf::Closure*);
 
+template <typename Request, typename Response>
+void call_ms(MetaService_Stub* stub, MetaServiceMethod<Request, Response> 
method,
+             brpc::Controller* cntl, const Request& req, Response* res) {
+    (stub->*method)(cntl, &req, res, nullptr);
+    if (!cntl->Failed()) {
+        // Meta Service may downgrade code for wire compatibility; restore the 
exact value.
+        restore_actual_code(res->mutable_status());
+    }
+}
+
 // Rate limiting context for retry_rpc
 struct RpcRateLimitCtx {
     HostLevelMSRpcRateLimiters* host_limiters {nullptr};
@@ -502,7 +523,7 @@ Status retry_rpc(MetaServiceRPC rpc, const Request& req, 
Response* res,
         cntl.set_max_retry(kBrpcRetryTimes);
         res->Clear();
         int error_code = 0;
-        (stub.get()->*method)(&cntl, &req, res, nullptr);
+        call_ms(stub.get(), method, &cntl, req, res);
 
         // Record QPS statistics for all RPCs sent to MS (success or failure)
         record_rpc_qps(rpc, rate_limit_ctx);
@@ -728,7 +749,7 @@ Status 
CloudMetaMgr::sync_tablet_rowsets_unlocked(CloudTablet* tablet,
         }
 
         auto start = std::chrono::steady_clock::now();
-        stub->get_rowset(&cntl, &req, &resp, nullptr);
+        call_ms(stub.get(), &MetaService_Stub::get_rowset, &cntl, req, &resp);
         auto end = std::chrono::steady_clock::now();
         int64_t latency = cntl.latency_us();
         _get_rowset_latency << latency;
diff --git a/be/src/cloud/cloud_meta_mgr.h b/be/src/cloud/cloud_meta_mgr.h
index b66418f0d3d..2a4419cff0b 100644
--- a/be/src/cloud/cloud_meta_mgr.h
+++ b/be/src/cloud/cloud_meta_mgr.h
@@ -65,6 +65,10 @@ Status bthread_fork_join(const 
std::vector<std::function<Status()>>& tasks, int
 Status bthread_fork_join(std::vector<std::function<Status()>>&& tasks, int 
concurrency,
                          std::future<Status>* fut);
 
+// Returns the exact actual_code when recognized, otherwise the 
legacy-compatible code.
+// Exposed for unit tests.
+MetaServiceCode get_response_code(const MetaServiceResponseStatus& status);
+
 class CloudMetaMgr {
 public:
     CloudMetaMgr() = default;
diff --git a/be/test/cloud/cloud_meta_mgr_test.cpp 
b/be/test/cloud/cloud_meta_mgr_test.cpp
index c6e99cb26b1..ff87378348e 100644
--- a/be/test/cloud/cloud_meta_mgr_test.cpp
+++ b/be/test/cloud/cloud_meta_mgr_test.cpp
@@ -21,6 +21,8 @@
 #include <gtest/gtest.h>
 
 #include <chrono>
+#include <cstdint>
+#include <limits>
 #include <memory>
 #include <random>
 #include <set>
@@ -44,6 +46,27 @@ class CloudMetaMgrTest : public testing::Test {
     void TearDown() override {}
 };
 
+TEST_F(CloudMetaMgrTest, response_status_uses_actual_code_when_valid) {
+    MetaServiceResponseStatus status;
+    status.set_code(MetaServiceCode::KV_TXN_CONFLICT);
+    status.set_actual_code(static_cast<int32_t>(MetaServiceCode::MS_TOO_BUSY));
+    EXPECT_EQ(get_response_code(status), MetaServiceCode::MS_TOO_BUSY);
+
+    status.set_code(MetaServiceCode::KV_TXN_CONFLICT);
+    
status.set_actual_code(static_cast<int32_t>(MetaServiceCode::KV_TXN_CONFLICT));
+    EXPECT_EQ(get_response_code(status), MetaServiceCode::KV_TXN_CONFLICT);
+
+    status.clear_actual_code();
+    EXPECT_EQ(get_response_code(status), MetaServiceCode::KV_TXN_CONFLICT);
+}
+
+TEST_F(CloudMetaMgrTest, response_status_falls_back_for_invalid_actual_code) {
+    MetaServiceResponseStatus status;
+    status.set_code(MetaServiceCode::KV_TXN_CONFLICT);
+    status.set_actual_code(std::numeric_limits<int32_t>::max());
+    EXPECT_EQ(get_response_code(status), MetaServiceCode::KV_TXN_CONFLICT);
+}
+
 static AbortTxnRequest get_abort_txn_request(CloudMetaMgr* meta_mgr, const 
StreamLoadContext& ctx) {
     auto* sp = SyncPoint::get_instance();
     sp->clear_all_call_backs();
diff --git a/cloud/src/meta-service/meta_service.h 
b/cloud/src/meta-service/meta_service.h
index a0594a945d6..70fb410e117 100644
--- a/cloud/src/meta-service/meta_service.h
+++ b/cloud/src/meta-service/meta_service.h
@@ -28,9 +28,11 @@
 #include <type_traits>
 
 #include "common/config.h"
+#include "common/defer.h"
 #include "common/stats.h"
 #include "cpp/sync_point.h"
 #include "meta-service/delete_bitmap_lock_white_list.h"
+#include "meta-service/meta_service_helper.h"
 #include "meta-service/txn_lazy_committer.h"
 #include "meta-store/txn_kv.h"
 #include "rate-limiter/rate_limiter.h"
@@ -1035,6 +1037,11 @@ private:
         using namespace std::chrono;
         brpc::ClosureGuard done_guard(done);
 
+        DORIS_CLOUD_DEFER {
+            auto* status = resp->mutable_status();
+            set_response_code(status, status->code(), status->msg());
+        };
+
         // life span of this defer MUST be longer than `done`
         std::unique_ptr<int, std::function<void(int*)>> defer_injection(
                 (int*)(0x01), [&, this](int*) { idempotent_injection(method, 
req, resp); });
diff --git a/cloud/src/meta-service/meta_service_helper.h 
b/cloud/src/meta-service/meta_service_helper.h
index 3ee903e3dfd..9b1fbed4648 100644
--- a/cloud/src/meta-service/meta_service_helper.h
+++ b/cloud/src/meta-service/meta_service_helper.h
@@ -25,6 +25,7 @@
 #include <memory>
 #include <string>
 #include <string_view>
+#include <utility>
 
 #include "common/bvars.h"
 #include "common/config.h"
@@ -41,6 +42,24 @@
 #include "resource-manager/resource_manager.h"
 
 namespace doris::cloud {
+inline MetaServiceCode get_legacy_code(MetaServiceCode code) {
+    switch (code) {
+    // MS_TOO_BUSY is a overload signal. Map it to KV_TXN_CONFLICT so the BE's 
existing
+    // conflict-retry path can retry the request.
+    case MetaServiceCode::MS_TOO_BUSY:
+        return MetaServiceCode::KV_TXN_CONFLICT;
+    default:
+        return code;
+    }
+}
+
+inline void set_response_code(MetaServiceResponseStatus* status, 
MetaServiceCode code,
+                              std::string msg) {
+    status->set_actual_code(static_cast<int32_t>(code));
+    status->set_code(get_legacy_code(code));
+    status->set_msg(std::move(msg));
+}
+
 inline std::string md5(const std::string& str) {
     unsigned char digest[MD5_DIGEST_LENGTH];
     MD5_CTX context;
@@ -315,17 +334,16 @@ inline MetaServiceCode cast_as(TxnErrorCode code) {
     [[maybe_unused]] MsStressDecision ms_stress_decision;                      
               \
     if (config::enable_ms_rate_limit || 
config::enable_ms_rate_limit_injection) {             \
         ms_stress_decision = get_ms_stress_decision();                         
               \
-    }                                                                          
               \
-    if ((config::enable_ms_rate_limit || 
config::enable_ms_rate_limit_injection) &&           \
-        RpcRateLimitWhitelist::instance().should_rate_limit(#func_name) &&     
               \
-        ms_stress_decision.under_great_stress()) {                             
               \
-        drop_request = true;                                                   
               \
-        code = MetaServiceCode::MS_TOO_BUSY;                                   
               \
-        msg = ms_stress_decision.debug_string();                               
               \
-        response->mutable_status()->set_code(code);                            
               \
-        response->mutable_status()->set_msg(msg);                              
               \
-        finish_rpc(#func_name, ctrl, request, response);                       
               \
-        return;                                                                
               \
+        if (RpcRateLimitWhitelist::instance().should_rate_limit(#func_name) && 
               \
+            ms_stress_decision.under_great_stress()) {                         
               \
+            drop_request = true;                                               
               \
+            msg = ms_stress_decision.debug_string();                           
               \
+            code = MetaServiceCode::MS_TOO_BUSY;                               
               \
+            response->mutable_status()->set_code(code);                        
               \
+            response->mutable_status()->set_msg(msg);                          
               \
+            finish_rpc(#func_name, ctrl, request, response);                   
               \
+            return;                                                            
               \
+        }                                                                      
               \
     }                                                                          
               \
     DORIS_CLOUD_DEFER {                                                        
               \
         response->mutable_status()->set_code(code);                            
               \
diff --git a/cloud/test/meta_service_helper_test.cpp 
b/cloud/test/meta_service_helper_test.cpp
index 73b3d37de5b..7b50792f88b 100644
--- a/cloud/test/meta_service_helper_test.cpp
+++ b/cloud/test/meta_service_helper_test.cpp
@@ -15,10 +15,17 @@
 // specific language governing permissions and limitations
 // under the License.
 
+#include "meta-service/meta_service_helper.h"
+
+#include <google/protobuf/descriptor.pb.h>
+#include <google/protobuf/dynamic_message.h>
 #include <gtest/gtest.h>
 
 #include <limits>
+#include <memory>
 #include <optional>
+#include <set>
+#include <string>
 #include <string_view>
 
 #include "common/config.h"
@@ -43,8 +50,67 @@ struct MsRateLimitInjectionConfigGuard {
     bool original_enable {config::enable_ms_rate_limit_injection};
     int32_t original_probability {config::ms_rate_limit_injection_probability};
 };
+
+google::protobuf::FileDescriptorProto legacy_status_file_descriptor() {
+    // Frozen subset of the pre-actual_code schema used by released clients.
+    google::protobuf::FileDescriptorProto file;
+    file.set_name("legacy_meta_service_status.proto");
+    file.set_package("doris.cloud.legacy");
+    file.set_syntax("proto2");
+
+    auto* code = file.add_enum_type();
+    code->set_name("MetaServiceCode");
+    auto* ok = code->add_value();
+    ok->set_name("OK");
+    ok->set_number(0);
+    auto* conflict = code->add_value();
+    conflict->set_name("KV_TXN_CONFLICT");
+    conflict->set_number(1005);
+
+    auto* status = file.add_message_type();
+    status->set_name("MetaServiceResponseStatus");
+    auto* code_field = status->add_field();
+    code_field->set_name("code");
+    code_field->set_number(1);
+    
code_field->set_label(google::protobuf::FieldDescriptorProto::LABEL_OPTIONAL);
+    code_field->set_type(google::protobuf::FieldDescriptorProto::TYPE_ENUM);
+    code_field->set_type_name(".doris.cloud.legacy.MetaServiceCode");
+    auto* msg_field = status->add_field();
+    msg_field->set_name("msg");
+    msg_field->set_number(2);
+    
msg_field->set_label(google::protobuf::FieldDescriptorProto::LABEL_OPTIONAL);
+    msg_field->set_type(google::protobuf::FieldDescriptorProto::TYPE_STRING);
+    return file;
+}
 } // namespace
 
+class MetaServiceWireCompatibilityTest : public testing::Test {
+protected:
+    void SetUp() override {
+        const auto* file = 
legacy_pool_.BuildFile(legacy_status_file_descriptor());
+        ASSERT_NE(file, nullptr);
+        legacy_status_descriptor_ = 
file->FindMessageTypeByName("MetaServiceResponseStatus");
+        ASSERT_NE(legacy_status_descriptor_, nullptr);
+        legacy_code_field_ = 
legacy_status_descriptor_->FindFieldByName("code");
+        ASSERT_NE(legacy_code_field_, nullptr);
+        legacy_msg_field_ = legacy_status_descriptor_->FindFieldByName("msg");
+        ASSERT_NE(legacy_msg_field_, nullptr);
+        legacy_status_prototype_ = 
legacy_factory_.GetPrototype(legacy_status_descriptor_);
+        ASSERT_NE(legacy_status_prototype_, nullptr);
+    }
+
+    std::unique_ptr<google::protobuf::Message> new_legacy_status() const {
+        return 
std::unique_ptr<google::protobuf::Message>(legacy_status_prototype_->New());
+    }
+
+    google::protobuf::DescriptorPool legacy_pool_;
+    google::protobuf::DynamicMessageFactory legacy_factory_ {&legacy_pool_};
+    const google::protobuf::Descriptor* legacy_status_descriptor_ = nullptr;
+    const google::protobuf::FieldDescriptor* legacy_code_field_ = nullptr;
+    const google::protobuf::FieldDescriptor* legacy_msg_field_ = nullptr;
+    const google::protobuf::Message* legacy_status_prototype_ = nullptr;
+};
+
 TEST(MetaServiceHelperTest, FdbClusterPressureNeedsLatencyAndNonWorkload) {
     MsStressMetrics metrics;
     metrics.fdb_commit_latency_ns = 51L * 1000 * 1000;
@@ -148,4 +214,188 @@ TEST(MetaServiceHelperTest, 
UsagePercentCalculationUsesEffectiveLimit) {
     ASSERT_EQ(internal::calculate_cpu_usage_percent(15e8, 1e9, 2.0), 75);
     ASSERT_EQ(internal::calculate_cpu_usage_percent(1, 0, 2.0), -1);
 }
+
+TEST_F(MetaServiceWireCompatibilityTest, 
LegacyClientReadsFallbackAndIgnoresActualCode) {
+    MetaServiceResponseStatus current_status;
+    set_response_code(&current_status, MetaServiceCode::MS_TOO_BUSY, "busy");
+
+    std::string wire;
+    ASSERT_TRUE(current_status.SerializeToString(&wire));
+    auto legacy_status = new_legacy_status();
+    ASSERT_TRUE(legacy_status->ParseFromString(wire));
+
+    const auto* reflection = legacy_status->GetReflection();
+    ASSERT_TRUE(reflection->HasField(*legacy_status, legacy_code_field_));
+    EXPECT_EQ(reflection->GetEnumValue(*legacy_status, legacy_code_field_),
+              MetaServiceCode::KV_TXN_CONFLICT);
+    EXPECT_EQ(reflection->GetString(*legacy_status, legacy_msg_field_), 
"busy");
+    EXPECT_EQ(legacy_status_descriptor_->FindFieldByName("actual_code"), 
nullptr);
+
+    const auto& unknown_fields = reflection->GetUnknownFields(*legacy_status);
+    ASSERT_EQ(unknown_fields.field_count(), 1);
+    EXPECT_EQ(unknown_fields.field(0).number(), 3);
+    EXPECT_EQ(unknown_fields.field(0).type(), 
google::protobuf::UnknownField::TYPE_VARINT);
+    EXPECT_EQ(unknown_fields.field(0).varint(), MetaServiceCode::MS_TOO_BUSY);
+
+    ASSERT_TRUE(legacy_status->SerializeToString(&wire));
+    MetaServiceResponseStatus round_trip_status;
+    ASSERT_TRUE(round_trip_status.ParseFromString(wire));
+    EXPECT_EQ(round_trip_status.code(), MetaServiceCode::KV_TXN_CONFLICT);
+    ASSERT_TRUE(round_trip_status.has_actual_code());
+    EXPECT_EQ(round_trip_status.actual_code(), MetaServiceCode::MS_TOO_BUSY);
+}
+
+TEST_F(MetaServiceWireCompatibilityTest, 
LegacyClientReadsUnknownEnumAsDefaultOk) {
+    MetaServiceResponseStatus incompatible_status;
+    incompatible_status.set_code(MetaServiceCode::MS_TOO_BUSY);
+
+    std::string wire;
+    ASSERT_TRUE(incompatible_status.SerializeToString(&wire));
+    auto legacy_status = new_legacy_status();
+    ASSERT_TRUE(legacy_status->ParseFromString(wire));
+
+    const auto* reflection = legacy_status->GetReflection();
+    EXPECT_FALSE(reflection->HasField(*legacy_status, legacy_code_field_));
+    EXPECT_EQ(reflection->GetEnumValue(*legacy_status, legacy_code_field_), 
MetaServiceCode::OK);
+    const auto& unknown_fields = reflection->GetUnknownFields(*legacy_status);
+    ASSERT_EQ(unknown_fields.field_count(), 1);
+    EXPECT_EQ(unknown_fields.field(0).number(), 1);
+    EXPECT_EQ(unknown_fields.field(0).type(), 
google::protobuf::UnknownField::TYPE_VARINT);
+    EXPECT_EQ(unknown_fields.field(0).varint(), MetaServiceCode::MS_TOO_BUSY);
+}
+
+TEST_F(MetaServiceWireCompatibilityTest, NewClientFallsBackForLegacyResponse) {
+    auto legacy_status = new_legacy_status();
+    const auto* reflection = legacy_status->GetReflection();
+    const auto* conflict =
+            
legacy_code_field_->enum_type()->FindValueByNumber(MetaServiceCode::KV_TXN_CONFLICT);
+    ASSERT_NE(conflict, nullptr);
+    reflection->SetEnum(legacy_status.get(), legacy_code_field_, conflict);
+    reflection->SetString(legacy_status.get(), legacy_msg_field_, "conflict");
+
+    std::string wire;
+    ASSERT_TRUE(legacy_status->SerializeToString(&wire));
+    MetaServiceResponseStatus current_status;
+    ASSERT_TRUE(current_status.ParseFromString(wire));
+    EXPECT_EQ(current_status.code(), MetaServiceCode::KV_TXN_CONFLICT);
+    EXPECT_FALSE(current_status.has_actual_code());
+}
+
+TEST(MetaServiceHelperTest, ResponseStatusUsesExactAndLegacyCodes) {
+    MetaServiceResponseStatus status;
+
+    set_response_code(&status, MetaServiceCode::MS_TOO_BUSY, "busy");
+    EXPECT_EQ(status.code(), MetaServiceCode::KV_TXN_CONFLICT);
+    EXPECT_EQ(status.actual_code(), MetaServiceCode::MS_TOO_BUSY);
+    EXPECT_EQ(status.msg(), "busy");
+
+    set_response_code(&status, MetaServiceCode::KV_TXN_CONFLICT, "conflict");
+    EXPECT_EQ(status.code(), MetaServiceCode::KV_TXN_CONFLICT);
+    EXPECT_EQ(status.actual_code(), MetaServiceCode::KV_TXN_CONFLICT);
+    EXPECT_EQ(status.msg(), "conflict");
+}
+
+TEST(MetaServiceHelperTest, ResponseStatusCoversEveryMetaServiceCode) {
+    std::set<MetaServiceCode> covered_codes;
+    auto expect_response_status = [&](MetaServiceCode code, MetaServiceCode 
expected_legacy_code) {
+        EXPECT_TRUE(covered_codes.insert(code).second)
+                << "Duplicate MetaServiceCode: " << MetaServiceCode_Name(code);
+
+        MetaServiceResponseStatus status;
+        set_response_code(&status, code, "");
+        EXPECT_EQ(status.code(), expected_legacy_code)
+                << "MetaServiceCode: " << MetaServiceCode_Name(code);
+        EXPECT_EQ(status.actual_code(), static_cast<int32_t>(code))
+                << "MetaServiceCode: " << MetaServiceCode_Name(code);
+    };
+
+    expect_response_status(MetaServiceCode::OK, MetaServiceCode::OK);
+    expect_response_status(MetaServiceCode::INVALID_ARGUMENT, 
MetaServiceCode::INVALID_ARGUMENT);
+    expect_response_status(MetaServiceCode::KV_TXN_CREATE_ERR, 
MetaServiceCode::KV_TXN_CREATE_ERR);
+    expect_response_status(MetaServiceCode::KV_TXN_GET_ERR, 
MetaServiceCode::KV_TXN_GET_ERR);
+    expect_response_status(MetaServiceCode::KV_TXN_COMMIT_ERR, 
MetaServiceCode::KV_TXN_COMMIT_ERR);
+    expect_response_status(MetaServiceCode::KV_TXN_CONFLICT, 
MetaServiceCode::KV_TXN_CONFLICT);
+    expect_response_status(MetaServiceCode::PROTOBUF_PARSE_ERR,
+                           MetaServiceCode::PROTOBUF_PARSE_ERR);
+    expect_response_status(MetaServiceCode::PROTOBUF_SERIALIZE_ERR,
+                           MetaServiceCode::PROTOBUF_SERIALIZE_ERR);
+    expect_response_status(MetaServiceCode::KV_TXN_STORE_GET_RETRYABLE,
+                           MetaServiceCode::KV_TXN_STORE_GET_RETRYABLE);
+    expect_response_status(MetaServiceCode::KV_TXN_STORE_COMMIT_RETRYABLE,
+                           MetaServiceCode::KV_TXN_STORE_COMMIT_RETRYABLE);
+    expect_response_status(MetaServiceCode::KV_TXN_STORE_CREATE_RETRYABLE,
+                           MetaServiceCode::KV_TXN_STORE_CREATE_RETRYABLE);
+    expect_response_status(MetaServiceCode::KV_TXN_TOO_OLD, 
MetaServiceCode::KV_TXN_TOO_OLD);
+    expect_response_status(MetaServiceCode::KV_TXN_MAYBE_COMMITTED,
+                           MetaServiceCode::KV_TXN_MAYBE_COMMITTED);
+    expect_response_status(MetaServiceCode::TXN_GEN_ID_ERR, 
MetaServiceCode::TXN_GEN_ID_ERR);
+    expect_response_status(MetaServiceCode::TXN_DUPLICATED_REQ,
+                           MetaServiceCode::TXN_DUPLICATED_REQ);
+    expect_response_status(MetaServiceCode::TXN_LABEL_ALREADY_USED,
+                           MetaServiceCode::TXN_LABEL_ALREADY_USED);
+    expect_response_status(MetaServiceCode::TXN_INVALID_STATUS,
+                           MetaServiceCode::TXN_INVALID_STATUS);
+    expect_response_status(MetaServiceCode::TXN_LABEL_NOT_FOUND,
+                           MetaServiceCode::TXN_LABEL_NOT_FOUND);
+    expect_response_status(MetaServiceCode::TXN_ID_NOT_FOUND, 
MetaServiceCode::TXN_ID_NOT_FOUND);
+    expect_response_status(MetaServiceCode::TXN_ALREADY_ABORTED,
+                           MetaServiceCode::TXN_ALREADY_ABORTED);
+    expect_response_status(MetaServiceCode::TXN_ALREADY_VISIBLE,
+                           MetaServiceCode::TXN_ALREADY_VISIBLE);
+    expect_response_status(MetaServiceCode::TXN_ALREADY_PRECOMMITED,
+                           MetaServiceCode::TXN_ALREADY_PRECOMMITED);
+    expect_response_status(MetaServiceCode::VERSION_NOT_FOUND, 
MetaServiceCode::VERSION_NOT_FOUND);
+    expect_response_status(MetaServiceCode::TABLET_NOT_FOUND, 
MetaServiceCode::TABLET_NOT_FOUND);
+    expect_response_status(MetaServiceCode::STALE_TABLET_CACHE,
+                           MetaServiceCode::STALE_TABLET_CACHE);
+    expect_response_status(MetaServiceCode::STALE_PREPARE_ROWSET,
+                           MetaServiceCode::STALE_PREPARE_ROWSET);
+    expect_response_status(MetaServiceCode::TXN_ALREADY_COMMITED,
+                           MetaServiceCode::TXN_ALREADY_COMMITED);
+    expect_response_status(MetaServiceCode::CLUSTER_NOT_FOUND, 
MetaServiceCode::CLUSTER_NOT_FOUND);
+    expect_response_status(MetaServiceCode::ALREADY_EXISTED, 
MetaServiceCode::ALREADY_EXISTED);
+    expect_response_status(MetaServiceCode::CLUSTER_ENDPOINT_MISSING,
+                           MetaServiceCode::CLUSTER_ENDPOINT_MISSING);
+    expect_response_status(MetaServiceCode::STORAGE_VAULT_NOT_FOUND,
+                           MetaServiceCode::STORAGE_VAULT_NOT_FOUND);
+    expect_response_status(MetaServiceCode::STAGE_NOT_FOUND, 
MetaServiceCode::STAGE_NOT_FOUND);
+    expect_response_status(MetaServiceCode::STAGE_GET_ERR, 
MetaServiceCode::STAGE_GET_ERR);
+    expect_response_status(MetaServiceCode::STATE_ALREADY_EXISTED_FOR_USER,
+                           MetaServiceCode::STATE_ALREADY_EXISTED_FOR_USER);
+    expect_response_status(MetaServiceCode::COPY_JOB_NOT_FOUND,
+                           MetaServiceCode::COPY_JOB_NOT_FOUND);
+    expect_response_status(MetaServiceCode::JOB_EXPIRED, 
MetaServiceCode::JOB_EXPIRED);
+    expect_response_status(MetaServiceCode::JOB_TABLET_BUSY, 
MetaServiceCode::JOB_TABLET_BUSY);
+    expect_response_status(MetaServiceCode::JOB_ALREADY_SUCCESS,
+                           MetaServiceCode::JOB_ALREADY_SUCCESS);
+    expect_response_status(MetaServiceCode::ROUTINE_LOAD_DATA_INCONSISTENT,
+                           MetaServiceCode::ROUTINE_LOAD_DATA_INCONSISTENT);
+    expect_response_status(MetaServiceCode::ROUTINE_LOAD_PROGRESS_NOT_FOUND,
+                           MetaServiceCode::ROUTINE_LOAD_PROGRESS_NOT_FOUND);
+    expect_response_status(MetaServiceCode::JOB_CHECK_ALTER_VERSION,
+                           MetaServiceCode::JOB_CHECK_ALTER_VERSION);
+    expect_response_status(MetaServiceCode::STREAMING_JOB_PROGRESS_NOT_FOUND,
+                           MetaServiceCode::STREAMING_JOB_PROGRESS_NOT_FOUND);
+    expect_response_status(MetaServiceCode::MAX_QPS_LIMIT, 
MetaServiceCode::MAX_QPS_LIMIT);
+    expect_response_status(MetaServiceCode::MS_TOO_BUSY, 
MetaServiceCode::KV_TXN_CONFLICT);
+    expect_response_status(MetaServiceCode::ERR_ENCRYPT, 
MetaServiceCode::ERR_ENCRYPT);
+    expect_response_status(MetaServiceCode::ERR_DECPYPT, 
MetaServiceCode::ERR_DECPYPT);
+    expect_response_status(MetaServiceCode::LOCK_EXPIRED, 
MetaServiceCode::LOCK_EXPIRED);
+    expect_response_status(MetaServiceCode::LOCK_CONFLICT, 
MetaServiceCode::LOCK_CONFLICT);
+    expect_response_status(MetaServiceCode::ROWSETS_EXPIRED, 
MetaServiceCode::ROWSETS_EXPIRED);
+    expect_response_status(MetaServiceCode::VERSION_NOT_MATCH, 
MetaServiceCode::VERSION_NOT_MATCH);
+    expect_response_status(MetaServiceCode::UPDATE_OVERRIDE_EXISTING_KV,
+                           MetaServiceCode::UPDATE_OVERRIDE_EXISTING_KV);
+    expect_response_status(MetaServiceCode::ROWSET_META_NOT_FOUND,
+                           MetaServiceCode::ROWSET_META_NOT_FOUND);
+    
expect_response_status(MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES,
+                           
MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES);
+    expect_response_status(MetaServiceCode::SCHEMA_DICT_NOT_FOUND,
+                           MetaServiceCode::SCHEMA_DICT_NOT_FOUND);
+    expect_response_status(MetaServiceCode::UNDEFINED_ERR, 
MetaServiceCode::UNDEFINED_ERR);
+
+    EXPECT_EQ(covered_codes.size(),
+              
static_cast<size_t>(MetaServiceCode_descriptor()->value_count()));
+}
+
 } // namespace doris::cloud
diff --git a/cloud/test/meta_service_http_test.cpp 
b/cloud/test/meta_service_http_test.cpp
index 7452fed53ec..c043a99f5ef 100644
--- a/cloud/test/meta_service_http_test.cpp
+++ b/cloud/test/meta_service_http_test.cpp
@@ -1413,6 +1413,11 @@ TEST(MetaServiceHttpTest, GetStageTest) {
 TEST(MetaServiceHttpTest, GetTabletStatsTest) {
     HttpContext ctx(true);
     auto& meta_service = ctx.meta_service_;
+    auto expected_http_body = [](GetTabletStatsResponse response) {
+        // The HTTP handler bypasses MetaServiceProxy, so its text response 
has no actual_code.
+        response.mutable_status()->clear_actual_code();
+        return response.DebugString() + "\n";
+    };
 
     constexpr auto db_id = 1000, table_id = 10001, index_id = 10002, 
partition_id = 10003,
                    tablet_id = 10004;
@@ -1437,7 +1442,7 @@ TEST(MetaServiceHttpTest, GetTabletStatsTest) {
         idx->set_tablet_id(tablet_id);
         auto [status_code, content] = 
ctx.forward<std::string>("get_tablet_stats", req);
         ASSERT_EQ(status_code, 200);
-        ASSERT_EQ(content, res.DebugString() + "\n");
+        ASSERT_EQ(content, expected_http_body(res));
     }
 
     // Insert rowset
@@ -1504,7 +1509,7 @@ TEST(MetaServiceHttpTest, GetTabletStatsTest) {
         idx->set_tablet_id(tablet_id);
         auto [status_code, content] = 
ctx.forward<std::string>("get_tablet_stats", req);
         ASSERT_EQ(status_code, 200);
-        ASSERT_EQ(content, res.DebugString() + "\n");
+        ASSERT_EQ(content, expected_http_body(res));
     }
 }
 
diff --git a/cloud/test/meta_service_test.cpp b/cloud/test/meta_service_test.cpp
index be3425c3736..05ef1f65357 100644
--- a/cloud/test/meta_service_test.cpp
+++ b/cloud/test/meta_service_test.cpp
@@ -8640,6 +8640,8 @@ TEST(MetaServiceTxnStoreRetryableTest, 
MaybeCommittedCodeWithoutRetryReturnsComm
 
     ASSERT_EQ(resp.status().code(), MetaServiceCode::KV_TXN_COMMIT_ERR)
             << " status is " << resp.status().msg() << ", code=" << 
resp.status().code();
+    ASSERT_TRUE(resp.status().has_actual_code());
+    EXPECT_EQ(resp.status().actual_code(), MetaServiceCode::KV_TXN_COMMIT_ERR);
     EXPECT_EQ(index, 1);
 
     SyncPoint::get_instance()->disable_processing();
@@ -8680,6 +8682,8 @@ TEST(MetaServiceTxnStoreRetryableTest, 
ReadMaybeCommittedCodeWithoutRetryReturns
 
     ASSERT_EQ(resp.status().code(), MetaServiceCode::KV_TXN_COMMIT_ERR)
             << " status is " << resp.status().msg() << ", code=" << 
resp.status().code();
+    ASSERT_TRUE(resp.status().has_actual_code());
+    EXPECT_EQ(resp.status().actual_code(), MetaServiceCode::KV_TXN_COMMIT_ERR);
     EXPECT_EQ(resp.version(), 2);
     EXPECT_EQ(index, 1);
 }
@@ -8719,6 +8723,8 @@ TEST(MetaServiceTxnStoreRetryableTest, 
RetryMaybeCommittedCodeReturnsCommitErr)
 
     ASSERT_EQ(resp.status().code(), MetaServiceCode::KV_TXN_COMMIT_ERR)
             << " status is " << resp.status().msg() << ", code=" << 
resp.status().code();
+    ASSERT_TRUE(resp.status().has_actual_code());
+    EXPECT_EQ(resp.status().actual_code(), MetaServiceCode::KV_TXN_COMMIT_ERR);
     EXPECT_GE(index, static_cast<size_t>(config::txn_store_retry_times + 1));
 
     SyncPoint::get_instance()->disable_processing();
@@ -8763,6 +8769,8 @@ TEST(MetaServiceTxnStoreRetryableTest, 
RetryReadMaybeCommittedCodeReturnsCommitE
 
     ASSERT_EQ(resp.status().code(), MetaServiceCode::KV_TXN_COMMIT_ERR)
             << " status is " << resp.status().msg() << ", code=" << 
resp.status().code();
+    ASSERT_TRUE(resp.status().has_actual_code());
+    EXPECT_EQ(resp.status().actual_code(), MetaServiceCode::KV_TXN_COMMIT_ERR);
     EXPECT_EQ(resp.version(), 2);
     EXPECT_GE(index, static_cast<size_t>(config::txn_store_retry_times + 1));
 }
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceClient.java 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceClient.java
index 2f85ef17a8e..044d69ac7a7 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceClient.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceClient.java
@@ -24,8 +24,19 @@ import org.apache.doris.common.Config;
 import com.google.common.base.Preconditions;
 import com.google.gson.Gson;
 import com.google.gson.stream.JsonReader;
+import com.google.protobuf.Descriptors;
+import com.google.protobuf.Message;
+import io.grpc.CallOptions;
+import io.grpc.Channel;
+import io.grpc.ClientCall;
+import io.grpc.ClientInterceptor;
+import io.grpc.ClientInterceptors;
 import io.grpc.ConnectivityState;
+import io.grpc.ForwardingClientCall;
+import io.grpc.ForwardingClientCallListener;
 import io.grpc.ManagedChannel;
+import io.grpc.Metadata;
+import io.grpc.MethodDescriptor;
 import io.grpc.NameResolverRegistry;
 import io.grpc.netty.shaded.io.grpc.netty.NettyChannelBuilder;
 import io.grpc.netty.shaded.io.netty.channel.ChannelOption;
@@ -77,11 +88,66 @@ public class MetaServiceClient {
                 .enableRetry()
                 .usePlaintext()
                 .withOption(ChannelOption.CONNECT_TIMEOUT_MILLIS, 
Config.meta_service_brpc_connect_timeout_ms).build();
-        stub = MetaServiceGrpc.newFutureStub(channel);
-        blockingStub = MetaServiceGrpc.newBlockingStub(channel);
+        Channel intercepted = ClientInterceptors.intercept(channel, new 
MetaServiceResponseStatusInterceptor());
+        stub = MetaServiceGrpc.newFutureStub(intercepted);
+        blockingStub = MetaServiceGrpc.newBlockingStub(intercepted);
         expiredAt = connectionAgeExpiredAt();
     }
 
+    private static final class MetaServiceResponseStatusInterceptor implements 
ClientInterceptor {
+        @Override
+        public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
+                MethodDescriptor<ReqT, RespT> method, CallOptions callOptions, 
Channel next) {
+            ClientCall<ReqT, RespT> call = next.newCall(method, callOptions);
+            return new ForwardingClientCall.SimpleForwardingClientCall<ReqT, 
RespT>(call) {
+                @Override
+                public void start(Listener<RespT> listener, Metadata headers) {
+                    Listener<RespT> normalizingListener =
+                            new 
ForwardingClientCallListener.SimpleForwardingClientCallListener<RespT>(listener)
 {
+                                @Override
+                                public void onMessage(RespT response) {
+                                    
super.onMessage(restoreActualCode(response));
+                                }
+                            };
+                    super.start(normalizingListener, headers);
+                }
+            };
+        }
+    }
+
+    @SuppressWarnings("unchecked")
+    // Restore the exact status code from actual_code when this FE recognizes 
it.
+    // Otherwise, keep
+    // the legacy-compatible value in code so responses from a newer Meta 
Service
+    // remain readable.
+    private static <Response> Response restoreActualCode(Response response) {
+        if (!(response instanceof Message)) {
+            return response;
+        }
+        Message message = (Message) response;
+        Descriptors.FieldDescriptor statusField = 
message.getDescriptorForType().findFieldByName("status");
+        if (statusField == null || !message.hasField(statusField)) {
+            return response;
+        }
+        Object statusObject = message.getField(statusField);
+        if (!(statusObject instanceof Cloud.MetaServiceResponseStatus)) {
+            return response;
+        }
+        Cloud.MetaServiceResponseStatus status = 
(Cloud.MetaServiceResponseStatus) statusObject;
+
+        if (!status.hasActualCode()) {
+            return response;
+        }
+        Cloud.MetaServiceCode code = 
Cloud.MetaServiceCode.forNumber(status.getActualCode());
+        if (code == null || code == status.getCode()) {
+            return response;
+        }
+        Cloud.MetaServiceResponseStatus restoredStatus = 
status.toBuilder().setCode(code).build();
+        Message.Builder builder = message.toBuilder();
+        builder.setField(statusField, restoredStatus);
+        return (Response) builder.build();
+    }
+
     private long connectionAgeExpiredAt() {
         long connectionAgeBase = 
Config.meta_service_connection_age_base_minutes;
         if (connectionAgeBase > 0) {
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/rpc/MetaServiceProxyTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/rpc/MetaServiceProxyTest.java
index 85c93ebf84d..6aa430f1e82 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/rpc/MetaServiceProxyTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/rpc/MetaServiceProxyTest.java
@@ -23,6 +23,9 @@ import org.apache.doris.common.jmockit.Deencapsulation;
 import org.apache.doris.rpc.RpcException;
 
 import com.google.common.util.concurrent.SettableFuture;
+import com.google.protobuf.DescriptorProtos;
+import com.google.protobuf.Descriptors;
+import com.google.protobuf.DynamicMessage;
 import org.junit.After;
 import org.junit.Assert;
 import org.junit.Before;
@@ -31,6 +34,7 @@ import org.mockito.Mockito;
 
 import java.util.Map;
 import java.util.Queue;
+import java.util.concurrent.Future;
 import java.util.concurrent.atomic.AtomicInteger;
 
 public class MetaServiceProxyTest {
@@ -196,11 +200,15 @@ public class MetaServiceProxyTest {
         serviceMap.put(Config.meta_service_endpoint, client);
 
         MetaServiceProxy.MetaServiceClientWrapper wrapper = 
Deencapsulation.getField(proxy, "w");
-        Cloud.GetVersionResponse tooBusyResponse = 
Cloud.GetVersionResponse.newBuilder()
-                .setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
-                        .setCode(Cloud.MetaServiceCode.MS_TOO_BUSY)
-                        .setMsg("server is overloaded"))
-                .build();
+        Cloud.GetVersionResponse tooBusyResponse = Deencapsulation.invoke(
+                MetaServiceClient.class,
+                "restoreActualCode",
+                Cloud.GetVersionResponse.newBuilder()
+                        .setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
+                                .setCode(Cloud.MetaServiceCode.KV_TXN_CONFLICT)
+                                
.setActualCode(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber())
+                                .setMsg("server is overloaded"))
+                        .build());
         Cloud.GetVersionResponse okResponse = 
Cloud.GetVersionResponse.newBuilder()
                 .setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
                         .setCode(Cloud.MetaServiceCode.OK))
@@ -215,6 +223,109 @@ public class MetaServiceProxyTest {
         Mockito.verify(client, Mockito.never()).shutdown(Mockito.anyBoolean());
     }
 
+    @Test
+    public void testGetInstancePrefersKnownActualCode() throws RpcException {
+        Cloud.MetaServiceResponseStatus status = callGetInstanceWithStatus(
+                Cloud.MetaServiceResponseStatus.newBuilder()
+                        .setCode(Cloud.MetaServiceCode.KV_TXN_CONFLICT)
+                        
.setActualCode(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber())
+                        .build());
+
+        Assert.assertEquals(Cloud.MetaServiceCode.MS_TOO_BUSY, 
status.getCode());
+        Assert.assertEquals(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber(), 
status.getActualCode());
+    }
+
+    @Test
+    public void testGetInstanceKeepsLegacyCodeForUnknownActualCode() throws 
RpcException {
+        Cloud.MetaServiceResponseStatus status = callGetInstanceWithStatus(
+                Cloud.MetaServiceResponseStatus.newBuilder()
+                        .setCode(Cloud.MetaServiceCode.KV_TXN_CONFLICT)
+                        .setActualCode(Integer.MAX_VALUE)
+                        .build());
+
+        Assert.assertEquals(Cloud.MetaServiceCode.KV_TXN_CONFLICT, 
status.getCode());
+        Assert.assertEquals(Integer.MAX_VALUE, status.getActualCode());
+    }
+
+    @Test
+    public void testGetInstanceKeepsLegacyCodeWithoutActualCode() throws 
RpcException {
+        Cloud.MetaServiceResponseStatus status = callGetInstanceWithStatus(
+                Cloud.MetaServiceResponseStatus.newBuilder()
+                        .setCode(Cloud.MetaServiceCode.KV_TXN_CONFLICT)
+                        .build());
+
+        Assert.assertEquals(Cloud.MetaServiceCode.KV_TXN_CONFLICT, 
status.getCode());
+        Assert.assertFalse(status.hasActualCode());
+    }
+
+    @Test
+    public void testGetVisibleVersionAsyncPrefersKnownActualCode() throws 
Exception {
+        MetaServiceProxy proxy = new MetaServiceProxy();
+        MetaServiceClient client = mockNormalClient();
+        putClient(proxy, client);
+        SettableFuture<Cloud.GetVersionResponse> rpcFuture = 
SettableFuture.create();
+        
Mockito.when(client.getVisibleVersionAsync(Mockito.any())).thenReturn(rpcFuture);
+
+        Future<Cloud.GetVersionResponse> normalizedFuture = 
proxy.getVisibleVersionAsync(
+                Cloud.GetVersionRequest.newBuilder().build());
+        rpcFuture.set(Cloud.GetVersionResponse.newBuilder()
+                .setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
+                        .setCode(Cloud.MetaServiceCode.KV_TXN_CONFLICT)
+                        
.setActualCode(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber()))
+                .build());
+
+        Cloud.GetVersionResponse response = normalizedFuture.get();
+        response = Deencapsulation.invoke(MetaServiceClient.class, 
"restoreActualCode", response);
+        Cloud.MetaServiceResponseStatus status = response.getStatus();
+        Assert.assertEquals(Cloud.MetaServiceCode.MS_TOO_BUSY, 
status.getCode());
+        Assert.assertEquals(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber(), 
status.getActualCode());
+    }
+
+    @Test
+    public void testLegacySchemaWireCompatibility() throws Exception {
+        Descriptors.Descriptor legacyStatusDescriptor = 
legacyStatusDescriptor();
+        Descriptors.FieldDescriptor legacyCodeField = 
legacyStatusDescriptor.findFieldByName("code");
+        Cloud.MetaServiceResponseStatus currentStatus = 
Cloud.MetaServiceResponseStatus.newBuilder()
+                .setCode(Cloud.MetaServiceCode.KV_TXN_CONFLICT)
+                .setActualCode(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber())
+                .setMsg("busy")
+                .build();
+
+        DynamicMessage legacyStatus = DynamicMessage.parseFrom(
+                legacyStatusDescriptor, currentStatus.toByteArray());
+        Descriptors.EnumValueDescriptor legacyCode =
+                (Descriptors.EnumValueDescriptor) 
legacyStatus.getField(legacyCodeField);
+        Assert.assertEquals(Cloud.MetaServiceCode.KV_TXN_CONFLICT.getNumber(), 
legacyCode.getNumber());
+        
Assert.assertNull(legacyStatusDescriptor.findFieldByName("actual_code"));
+        Assert.assertTrue(legacyStatus.getUnknownFields().hasField(3));
+        
Assert.assertEquals(Long.valueOf(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber()),
+                
legacyStatus.getUnknownFields().getField(3).getVarintList().get(0));
+
+        Cloud.MetaServiceResponseStatus roundTripStatus =
+                
Cloud.MetaServiceResponseStatus.parseFrom(legacyStatus.toByteArray());
+        Assert.assertEquals(Cloud.MetaServiceCode.KV_TXN_CONFLICT, 
roundTripStatus.getCode());
+        Assert.assertEquals(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber(), 
roundTripStatus.getActualCode());
+    }
+
+    @Test
+    public void testLegacySchemaReadsUnknownEnumAsDefaultOk() throws Exception 
{
+        Descriptors.Descriptor legacyStatusDescriptor = 
legacyStatusDescriptor();
+        Descriptors.FieldDescriptor legacyCodeField = 
legacyStatusDescriptor.findFieldByName("code");
+        Cloud.MetaServiceResponseStatus incompatibleStatus = 
Cloud.MetaServiceResponseStatus.newBuilder()
+                .setCode(Cloud.MetaServiceCode.MS_TOO_BUSY)
+                .build();
+
+        DynamicMessage legacyStatus = DynamicMessage.parseFrom(
+                legacyStatusDescriptor, incompatibleStatus.toByteArray());
+        Descriptors.EnumValueDescriptor legacyCode =
+                (Descriptors.EnumValueDescriptor) 
legacyStatus.getField(legacyCodeField);
+        Assert.assertFalse(legacyStatus.hasField(legacyCodeField));
+        Assert.assertEquals(Cloud.MetaServiceCode.OK.getNumber(), 
legacyCode.getNumber());
+        Assert.assertTrue(legacyStatus.getUnknownFields().hasField(1));
+        
Assert.assertEquals(Long.valueOf(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber()),
+                
legacyStatus.getUnknownFields().getField(1).getVarintList().get(0));
+    }
+
     @Test
     public void testExecuteRequestFailureAfterTooBusyRetries() throws 
RpcException {
         Config.meta_service_rpc_retry_cnt = 2;
@@ -383,6 +494,55 @@ public class MetaServiceProxyTest {
         }
     }
 
+    private Cloud.MetaServiceResponseStatus callGetInstanceWithStatus(
+            Cloud.MetaServiceResponseStatus responseStatus) throws 
RpcException {
+        MetaServiceProxy proxy = new MetaServiceProxy();
+        MetaServiceClient client = mockNormalClient();
+        putClient(proxy, client);
+        
Mockito.when(client.getInstance(Mockito.any())).thenReturn(Cloud.GetInstanceResponse.newBuilder()
+                .setStatus(responseStatus)
+                .build());
+        Cloud.GetInstanceResponse response = 
proxy.getInstance(Cloud.GetInstanceRequest.newBuilder().build());
+        response = Deencapsulation.invoke(MetaServiceClient.class, 
"restoreActualCode", response);
+        return response.getStatus();
+    }
+
+    private Descriptors.Descriptor legacyStatusDescriptor() throws 
Descriptors.DescriptorValidationException {
+        // Frozen subset of the actual_code schema used by released clients.
+        DescriptorProtos.EnumDescriptorProto legacyCode = 
DescriptorProtos.EnumDescriptorProto.newBuilder()
+                .setName("MetaServiceCode")
+                
.addValue(DescriptorProtos.EnumValueDescriptorProto.newBuilder()
+                        .setName("OK")
+                        .setNumber(0))
+                
.addValue(DescriptorProtos.EnumValueDescriptorProto.newBuilder()
+                        .setName("KV_TXN_CONFLICT")
+                        
.setNumber(Cloud.MetaServiceCode.KV_TXN_CONFLICT.getNumber()))
+                .build();
+        DescriptorProtos.DescriptorProto legacyStatus = 
DescriptorProtos.DescriptorProto.newBuilder()
+                .setName("MetaServiceResponseStatus")
+                .addField(DescriptorProtos.FieldDescriptorProto.newBuilder()
+                        .setName("code")
+                        .setNumber(1)
+                        
.setLabel(DescriptorProtos.FieldDescriptorProto.Label.LABEL_OPTIONAL)
+                        
.setType(DescriptorProtos.FieldDescriptorProto.Type.TYPE_ENUM)
+                        .setTypeName(".doris.cloud.legacy.MetaServiceCode"))
+                .addField(DescriptorProtos.FieldDescriptorProto.newBuilder()
+                        .setName("msg")
+                        .setNumber(2)
+                        
.setLabel(DescriptorProtos.FieldDescriptorProto.Label.LABEL_OPTIONAL)
+                        
.setType(DescriptorProtos.FieldDescriptorProto.Type.TYPE_STRING))
+                .build();
+        DescriptorProtos.FileDescriptorProto legacyFile = 
DescriptorProtos.FileDescriptorProto.newBuilder()
+                .setName("legacy_meta_service_status.proto")
+                .setPackage("doris.cloud.legacy")
+                .setSyntax("proto2")
+                .addEnumType(legacyCode)
+                .addMessageType(legacyStatus)
+                .build();
+        return Descriptors.FileDescriptor.buildFrom(
+                legacyFile, new 
Descriptors.FileDescriptor[0]).findMessageTypeByName("MetaServiceResponseStatus");
+    }
+
     private Cloud.GetVersionRequest buildBatchPartitionVersionRequest(int 
partitionNum) {
         Cloud.GetVersionRequest.Builder builder = 
Cloud.GetVersionRequest.newBuilder()
                 .setBatchMode(true);
diff --git a/gensrc/proto/cloud.proto b/gensrc/proto/cloud.proto
index dcf1e3b28aa..51c522b6272 100644
--- a/gensrc/proto/cloud.proto
+++ b/gensrc/proto/cloud.proto
@@ -1462,8 +1462,14 @@ message RestoreJobResponse {
 }
 
 message MetaServiceResponseStatus {
+    // Legacy-compatible status code. Keep this value recognizable by all 
released clients.
     optional MetaServiceCode code = 1;
     optional string msg = 2;
+    // Exact client-visible status code encoded as int32, so proto2 clients do 
not drop unknown
+    // enum values. Internal retry signals must be converted before the 
response is sent.
+    // New clients should use this field when the local enum descriptor 
recognizes the value,
+    // otherwise fall back to `code`.
+    optional int32 actual_code = 3;
 }
 
 message MetaServiceHttpRequest {
@@ -1800,6 +1806,13 @@ message RecycleInstanceResponse {
 }
 
 enum MetaServiceCode {
+    // Compatibility rule: proto2 optional enum fields drop unknown enum 
values into
+    // UnknownFieldSet, so old clients read an unset 
`MetaServiceResponseStatus.code` as OK.
+
+    // MetaService must write the exact client-visible code to
+    // `MetaServiceResponseStatus.actual_code` and write only a legacy 
fallback code to
+    // `MetaServiceResponseStatus.code`. Any newly added error code that may 
be returned to
+    // clients must be mapped in get_legacy_code().
     OK = 0;
 
     //Meta service internal error
@@ -1817,10 +1830,7 @@ enum MetaServiceCode {
     KV_TXN_STORE_COMMIT_RETRYABLE = 1009;
     KV_TXN_STORE_CREATE_RETRYABLE = 1010;
     KV_TXN_TOO_OLD = 1011;
-    // WARNING: KV_TXN_MAYBE_COMMITTED is NOT returned to clients. It is kept 
as an
-    // internal retry signal inside MetaServiceProxy::call_impl(), then 
downgraded to
-    // KV_TXN_COMMIT_ERR before the response is sent back. Older BE/FE 
versions do not
-    // recognize this enum value, and proto2 would otherwise fall back to OK 
(= 0).
+    // WARNING: KV_TXN_MAYBE_COMMITTED must be downgraded through the legacy 
status channel.
     KV_TXN_MAYBE_COMMITTED = 1012;
 
     //Doris error
@@ -1884,11 +1894,6 @@ enum MetaServiceCode {
 
     SCHEMA_DICT_NOT_FOUND = 11001;
 
-    // WARNING: Before adding a new MetaServiceCode, consider backward 
compatibility.
-    // This is a proto2 optional enum field. If an older client receives an 
unrecognized
-    // enum value, the field is treated as unset and defaults to OK (= 0), 
silently
-    // turning errors into success. Any new error code that may be sent to 
older clients
-    // MUST be downgraded to a legacy code before the response is sent to the 
client.
     UNDEFINED_ERR = 1000000;
 }
 


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

Reply via email to