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 140863942b6 branch-4.1: [fix](cloud) Fail closed on unknown
MetaService response codes #66364 (#66716)
140863942b6 is described below
commit 140863942b65447b456a32e1f78eef52fcd719de
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Thu Aug 13 18:34:41 2026 +0800
branch-4.1: [fix](cloud) Fail closed on unknown MetaService response codes
#66364 (#66716)
Cherry-picked from #66364
Co-authored-by: Yixuan Wang <[email protected]>
---
be/src/cloud/cloud_meta_mgr.cpp | 50 ++++-
be/src/cloud/cloud_meta_mgr.h | 3 +-
be/test/cloud/cloud_meta_mgr_test.cpp | 32 ++-
cloud/src/meta-service/meta_service_helper.h | 26 ++-
cloud/test/meta_service_helper_test.cpp | 220 +++++++++++----------
cloud/test/txn_lazy_commit_test.cpp | 2 +-
.../apache/doris/cloud/rpc/MetaServiceClient.java | 55 +++++-
.../doris/cloud/rpc/MetaServiceProxyTest.java | 48 +++++
gensrc/proto/cloud.proto | 7 +-
9 files changed, 320 insertions(+), 123 deletions(-)
diff --git a/be/src/cloud/cloud_meta_mgr.cpp b/be/src/cloud/cloud_meta_mgr.cpp
index c912576d0b3..8c4069e0e22 100644
--- a/be/src/cloud/cloud_meta_mgr.cpp
+++ b/be/src/cloud/cloud_meta_mgr.cpp
@@ -154,11 +154,55 @@ Status
bthread_fork_join(std::vector<std::function<Status()>>&& tasks, int concu
return Status::OK();
}
+// Resolve the status code returned by Meta Service (MS) for BE/FE clients of
different version.
+// Assuming MS is always the latest version, it sends both the meta-service
error code and a code that
+// older clients can decode:
+//
+// latest MS
+// +---------------------------------------+
+// | actual_code = meta-service error code |
+// | code = compatible code |
+// +----------------+----------------------+
+// |
+// +-------------+-------------+
+// | |
+// old BE/FE without old BE/FE with
+// the actual_code field the actual_code field
+// | |
+// ignores actual_code local enum recognizes
+// and reads code actual_code value?
+// yes no
+// | |
+// use actual code use code only when it
+// is explicit and non-OK
+// |
+// otherwise return
+// UNDEFINED_ERR
+//
+// After MS adds an error code, an older actual_code-aware client may not have
that enum value;
+// MetaServiceCode_IsValid detects this case. The non-OK fallback check is
essential:
+// if MS ignore or incorrectly converts the compatible code to OK, an unknown
error
+// must remain an error instead of becoming a false success.
MetaServiceCode get_response_code(const MetaServiceResponseStatus& status) {
- if (status.has_actual_code() &&
MetaServiceCode_IsValid(status.actual_code())) {
- return static_cast<MetaServiceCode>(status.actual_code());
+ if (status.has_actual_code()) {
+ // Check whether this client build contains the code in its
MetaServiceCode enum.
+ if (MetaServiceCode_IsValid(status.actual_code())) {
+ return static_cast<MetaServiceCode>(status.actual_code());
+ }
+ // An older client may use the compatible code, but unsupported cases
must return an explicit error.
+ // Return the non-OK compatible code prepared by MS for older clients.
+ if (status.has_code() && status.code() != MetaServiceCode::OK) {
+ return status.code();
+ }
+ // Never return OK when the compatible code is absent or invalid.
+ return MetaServiceCode::UNDEFINED_ERR;
+ }
+ // A legacy response has only code, so return its explicit value,
including a real OK.
+ if (status.has_code()) {
+ return status.code();
}
- return status.code();
+ // A response missing both fields is invalid and must be rejected.
+ return MetaServiceCode::UNDEFINED_ERR;
}
namespace {
diff --git a/be/src/cloud/cloud_meta_mgr.h b/be/src/cloud/cloud_meta_mgr.h
index 2a4419cff0b..50b8353b7e6 100644
--- a/be/src/cloud/cloud_meta_mgr.h
+++ b/be/src/cloud/cloud_meta_mgr.h
@@ -65,7 +65,8 @@ 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.
+// Returns the exact actual_code when recognized. An unknown actual_code uses
an explicit non-OK
+// legacy fallback and otherwise fails closed. Responses from a legacy Meta
Service use code.
// Exposed for unit tests.
MetaServiceCode get_response_code(const MetaServiceResponseStatus& status);
diff --git a/be/test/cloud/cloud_meta_mgr_test.cpp
b/be/test/cloud/cloud_meta_mgr_test.cpp
index ff87378348e..9b10f80a218 100644
--- a/be/test/cloud/cloud_meta_mgr_test.cpp
+++ b/be/test/cloud/cloud_meta_mgr_test.cpp
@@ -56,17 +56,47 @@ TEST_F(CloudMetaMgrTest,
response_status_uses_actual_code_when_valid) {
status.set_actual_code(static_cast<int32_t>(MetaServiceCode::KV_TXN_CONFLICT));
EXPECT_EQ(get_response_code(status), MetaServiceCode::KV_TXN_CONFLICT);
+ status.set_code(MetaServiceCode::KV_TXN_CONFLICT);
+ status.set_actual_code(static_cast<int32_t>(MetaServiceCode::OK));
+ EXPECT_EQ(get_response_code(status), MetaServiceCode::OK);
+
+ status.clear_code();
+ 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.clear_actual_code();
EXPECT_EQ(get_response_code(status), MetaServiceCode::KV_TXN_CONFLICT);
}
-TEST_F(CloudMetaMgrTest, response_status_falls_back_for_invalid_actual_code) {
+TEST_F(CloudMetaMgrTest,
response_status_falls_back_to_non_ok_code_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);
}
+TEST_F(CloudMetaMgrTest,
response_status_returns_undefined_for_invalid_actual_code_with_ok) {
+ MetaServiceResponseStatus status;
+ status.set_code(MetaServiceCode::OK);
+ status.set_actual_code(std::numeric_limits<int32_t>::max());
+ EXPECT_EQ(get_response_code(status), MetaServiceCode::UNDEFINED_ERR);
+}
+
+TEST_F(CloudMetaMgrTest,
response_status_returns_undefined_for_invalid_actual_code_without_code) {
+ MetaServiceResponseStatus status;
+ status.set_actual_code(std::numeric_limits<int32_t>::max());
+ EXPECT_EQ(get_response_code(status), MetaServiceCode::UNDEFINED_ERR);
+}
+
+TEST_F(CloudMetaMgrTest, response_status_returns_undefined_without_any_code) {
+ MetaServiceResponseStatus status;
+ EXPECT_EQ(get_response_code(status), MetaServiceCode::UNDEFINED_ERR);
+
+ status.set_code(MetaServiceCode::OK);
+ EXPECT_EQ(get_response_code(status), MetaServiceCode::OK);
+}
+
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_helper.h
b/cloud/src/meta-service/meta_service_helper.h
index 9b1fbed4648..ab8605c20a1 100644
--- a/cloud/src/meta-service/meta_service_helper.h
+++ b/cloud/src/meta-service/meta_service_helper.h
@@ -18,6 +18,7 @@
#pragma once
#include <brpc/controller.h>
+#include <fmt/core.h>
#include <gen_cpp/cloud.pb.h>
#include <openssl/md5.h>
@@ -42,22 +43,35 @@
#include "resource-manager/resource_manager.h"
namespace doris::cloud {
-inline MetaServiceCode get_legacy_code(MetaServiceCode code) {
+// Converts a response code and message to values that older clients can read.
+// set_response_code() stores the original code in actual_code
+// Call this function only from set_response_code() or from unit tests; do not
call it from other production code.
+// When adding an error code that may be returned to clients, must add its
conversion here.
+inline std::pair<MetaServiceCode, std::string>
resolve_response_code_and_msg(MetaServiceCode code,
+
std::string msg) {
switch (code) {
- // MS_TOO_BUSY is a overload signal. Map it to KV_TXN_CONFLICT so the BE's
existing
+ // MS_TOO_BUSY is an 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;
+ msg += std::string((msg.empty() ? "" : ", ")) +
+ "[MS_TOO_BUSY will be converted to code=KV_TXN_CONFLICT for old
version clients]";
+ return {MetaServiceCode::KV_TXN_CONFLICT, std::move(msg)};
+ case MetaServiceCode::TXN_ALREADY_COMMITED:
+ msg += std::string((msg.empty() ? "" : ", ")) +
+ "[TXN_ALREADY_COMMITED will be converted to code=UNDEFINED_ERR
for old version "
+ "clients]";
+ return {MetaServiceCode::UNDEFINED_ERR, std::move(msg)};
default:
- return code;
+ return {code, std::move(msg)};
}
}
inline void set_response_code(MetaServiceResponseStatus* status,
MetaServiceCode code,
std::string msg) {
+ auto [resolved_code, resolved_msg] = resolve_response_code_and_msg(code,
std::move(msg));
status->set_actual_code(static_cast<int32_t>(code));
- status->set_code(get_legacy_code(code));
- status->set_msg(std::move(msg));
+ status->set_code(resolved_code);
+ status->set_msg(std::move(resolved_msg));
}
inline std::string md5(const std::string& str) {
diff --git a/cloud/test/meta_service_helper_test.cpp
b/cloud/test/meta_service_helper_test.cpp
index 7b50792f88b..73ccf12f42d 100644
--- a/cloud/test/meta_service_helper_test.cpp
+++ b/cloud/test/meta_service_helper_test.cpp
@@ -41,6 +41,105 @@ int64_t calculate_cpu_usage_percent(double delta_cpu_ns,
double delta_wall_ns, d
} // namespace internal
namespace {
+// IMPORTANT: Never-Never-Never add new codes to this snapshot. New codes must
be mapped to a
+// LegacyFallbackCode and verified with
expect_legacy_fallback_response_status().
+const std::set<MetaServiceCode> identity_snapshot = {
+ MetaServiceCode::OK,
+ MetaServiceCode::INVALID_ARGUMENT,
+ MetaServiceCode::KV_TXN_CREATE_ERR,
+ MetaServiceCode::KV_TXN_GET_ERR,
+ MetaServiceCode::KV_TXN_COMMIT_ERR,
+ MetaServiceCode::KV_TXN_CONFLICT,
+ MetaServiceCode::PROTOBUF_PARSE_ERR,
+ MetaServiceCode::PROTOBUF_SERIALIZE_ERR,
+ MetaServiceCode::KV_TXN_STORE_GET_RETRYABLE,
+ MetaServiceCode::KV_TXN_STORE_COMMIT_RETRYABLE,
+ MetaServiceCode::KV_TXN_STORE_CREATE_RETRYABLE,
+ MetaServiceCode::KV_TXN_TOO_OLD,
+ MetaServiceCode::KV_TXN_MAYBE_COMMITTED,
+ MetaServiceCode::TXN_GEN_ID_ERR,
+ MetaServiceCode::TXN_DUPLICATED_REQ,
+ MetaServiceCode::TXN_LABEL_ALREADY_USED,
+ MetaServiceCode::TXN_INVALID_STATUS,
+ MetaServiceCode::TXN_LABEL_NOT_FOUND,
+ MetaServiceCode::TXN_ID_NOT_FOUND,
+ MetaServiceCode::TXN_ALREADY_ABORTED,
+ MetaServiceCode::TXN_ALREADY_VISIBLE,
+ MetaServiceCode::TXN_ALREADY_PRECOMMITED,
+ MetaServiceCode::VERSION_NOT_FOUND,
+ MetaServiceCode::TABLET_NOT_FOUND,
+ MetaServiceCode::STALE_TABLET_CACHE,
+ MetaServiceCode::STALE_PREPARE_ROWSET,
+ MetaServiceCode::CLUSTER_NOT_FOUND,
+ MetaServiceCode::ALREADY_EXISTED,
+ MetaServiceCode::CLUSTER_ENDPOINT_MISSING,
+ MetaServiceCode::STORAGE_VAULT_NOT_FOUND,
+ MetaServiceCode::STAGE_NOT_FOUND,
+ MetaServiceCode::STAGE_GET_ERR,
+ MetaServiceCode::STATE_ALREADY_EXISTED_FOR_USER,
+ MetaServiceCode::COPY_JOB_NOT_FOUND,
+ MetaServiceCode::JOB_EXPIRED,
+ MetaServiceCode::JOB_TABLET_BUSY,
+ MetaServiceCode::JOB_ALREADY_SUCCESS,
+ MetaServiceCode::ROUTINE_LOAD_DATA_INCONSISTENT,
+ MetaServiceCode::ROUTINE_LOAD_PROGRESS_NOT_FOUND,
+ MetaServiceCode::JOB_CHECK_ALTER_VERSION,
+ MetaServiceCode::STREAMING_JOB_PROGRESS_NOT_FOUND,
+ MetaServiceCode::MAX_QPS_LIMIT,
+ MetaServiceCode::ERR_ENCRYPT,
+ MetaServiceCode::ERR_DECPYPT,
+ MetaServiceCode::LOCK_EXPIRED,
+ MetaServiceCode::LOCK_CONFLICT,
+ MetaServiceCode::ROWSETS_EXPIRED,
+ MetaServiceCode::VERSION_NOT_MATCH,
+ MetaServiceCode::UPDATE_OVERRIDE_EXISTING_KV,
+ MetaServiceCode::ROWSET_META_NOT_FOUND,
+ MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES,
+ MetaServiceCode::SCHEMA_DICT_NOT_FOUND,
+ MetaServiceCode::UNDEFINED_ERR,
+};
+
+// IMPORTANT: Never-Never-Never modify or extend this enum. New error codes
must be mapped to one of the
+// existing legacy fallback codes below.
+enum class LegacyFallbackCode : int32_t {
+ UNDEFINED_ERR = static_cast<int32_t>(MetaServiceCode::UNDEFINED_ERR),
+ KV_TXN_CONFLICT = static_cast<int32_t>(MetaServiceCode::KV_TXN_CONFLICT),
+};
+
+void verify_response_status_impl(std::set<MetaServiceCode>& covered_codes,
MetaServiceCode code,
+ int32_t expected_legacy_code) {
+ EXPECT_TRUE(covered_codes.insert(code).second)
+ << "Duplicate MetaServiceCode: " << MetaServiceCode_Name(code);
+
+ MetaServiceResponseStatus status;
+ set_response_code(&status, code, "");
+ EXPECT_EQ(static_cast<int32_t>(status.code()), expected_legacy_code)
+ << "MetaServiceCode: " << MetaServiceCode_Name(code);
+ EXPECT_EQ(status.actual_code(), static_cast<int32_t>(code))
+ << "MetaServiceCode: " << MetaServiceCode_Name(code);
+}
+
+void verify_response_status(std::set<MetaServiceCode>& covered_codes,
MetaServiceCode code,
+ int32_t expected_legacy_code) {
+ if (!identity_snapshot.contains(code)) {
+ EXPECT_TRUE(false)
+ << "MetaServiceCode " << MetaServiceCode_Name(code)
+ << " is not in identity_snapshot. New error codes must be
mapped to a "
+ "LegacyFallbackCode in resolve_response_code_and_msg() and
verified with "
+ "expect_legacy_fallback_response_status().";
+ }
+ verify_response_status_impl(covered_codes, code, expected_legacy_code);
+}
+
+// New error codes may only be converted to a value allowed by
LegacyFallbackCode.
+// Resolve the conversion in resolve_response_code_and_msg();
+// For example, MS_TOO_BUSY maps to KV_TXN_CONFLICT so that the BE can retry
it.
+void expect_legacy_fallback_response_status(std::set<MetaServiceCode>&
covered_codes,
+ MetaServiceCode code,
+ LegacyFallbackCode
expected_legacy_code) {
+ verify_response_status_impl(covered_codes, code,
static_cast<int32_t>(expected_legacy_code));
+}
+
struct MsRateLimitInjectionConfigGuard {
~MsRateLimitInjectionConfigGuard() {
config::enable_ms_rate_limit_injection = original_enable;
@@ -228,7 +327,9 @@ TEST_F(MetaServiceWireCompatibilityTest,
LegacyClientReadsFallbackAndIgnoresActu
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(reflection->GetString(*legacy_status, legacy_msg_field_),
+ "busy, [MS_TOO_BUSY will be converted to code=KV_TXN_CONFLICT
for old version "
+ "clients]");
EXPECT_EQ(legacy_status_descriptor_->FindFieldByName("actual_code"),
nullptr);
const auto& unknown_fields = reflection->GetUnknownFields(*legacy_status);
@@ -287,115 +388,36 @@ TEST(MetaServiceHelperTest,
ResponseStatusUsesExactAndLegacyCodes) {
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");
+ EXPECT_EQ(status.msg(),
+ "busy, [MS_TOO_BUSY will be converted to code=KV_TXN_CONFLICT
for old version "
+ "clients]");
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");
+
+ set_response_code(&status, MetaServiceCode::MS_TOO_BUSY, "");
+ EXPECT_EQ(status.msg(),
+ "[MS_TOO_BUSY will be converted to code=KV_TXN_CONFLICT for old
version clients]");
}
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);
+ for (auto code : identity_snapshot) {
+ verify_response_status(covered_codes, code,
static_cast<int32_t>(code));
+ }
+
+ expect_legacy_fallback_response_status(covered_codes,
MetaServiceCode::MS_TOO_BUSY,
+
LegacyFallbackCode::KV_TXN_CONFLICT);
+ expect_legacy_fallback_response_status(covered_codes,
MetaServiceCode::TXN_ALREADY_COMMITED,
+ LegacyFallbackCode::UNDEFINED_ERR);
EXPECT_EQ(covered_codes.size(),
-
static_cast<size_t>(MetaServiceCode_descriptor()->value_count()));
+ static_cast<size_t>(MetaServiceCode_descriptor()->value_count()))
+ << "A new MetaServiceCode was added. Map it to a
LegacyFallbackCode in "
+ "resolve_response_code_and_msg() and verify it with "
+ "expect_legacy_fallback_response_status().";
}
} // namespace doris::cloud
diff --git a/cloud/test/txn_lazy_commit_test.cpp
b/cloud/test/txn_lazy_commit_test.cpp
index 60a33067143..41374e67a75 100644
--- a/cloud/test/txn_lazy_commit_test.cpp
+++ b/cloud/test/txn_lazy_commit_test.cpp
@@ -3448,7 +3448,7 @@ TEST(TxnLazyCommitTest,
CommitTxnEventuallyWithAbortAfterCommitTest) {
req.set_cloud_unique_id("test_cloud_unique_id");
meta_service->abort_txn(reinterpret_cast<::google::protobuf::RpcController*>(&cntl),
&req,
&res, nullptr);
- ASSERT_EQ(res.status().code(), MetaServiceCode::TXN_ALREADY_COMMITED);
+ ASSERT_EQ(res.status().actual_code(),
MetaServiceCode::TXN_ALREADY_COMMITED);
});
// mock rowset and tablet
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 7a9f3b68e31..2d2a12fd409 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
@@ -110,18 +110,44 @@ public class MetaServiceClient {
}
}
+ // Resolve the status code returned by Meta Service (MS) for FE clients of
different versions.
+ // Assuming MS is always the latest version, it sends both the
meta-service error code and a code that
+ // older clients can decode:
+ //
+ // latest MS
+ // +---------------------------------------+
+ // | actual_code = meta-service error code |
+ // | code = compatible code |
+ // +----------------+----------------------+
+ // |
+ // +-------------+-------------+
+ // | |
+ // old FE without the old FE with the
+ // actual_code field actual_code field
+ // | |
+ // ignores actual_code local enum recognizes
+ // and reads code actual_code value?
+ // / \
+ // yes no
+ // | |
+ // use actual code use code only
when it
+ // is explicit and
non-OK
+ // |
+ // otherwise
return
+ // UNDEFINED_ERR
+ //
+ // After MS adds an error code, an older actual_code-aware client may not
have that enum value;
+ // Cloud.MetaServiceCode.forNumber returns null in this case. The non-OK
fallback check is essential:
+ // if MS ignores or incorrectly converts the compatible code to OK, an
unknown error must remain an
+ // error instead of becoming a false success.
@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)) {
+ if (statusField == null) {
return response;
}
Object statusObject = message.getField(statusField);
@@ -130,11 +156,22 @@ public class MetaServiceClient {
}
Cloud.MetaServiceResponseStatus status =
(Cloud.MetaServiceResponseStatus) statusObject;
- if (!status.hasActualCode()) {
- return response;
+ Cloud.MetaServiceCode code;
+ if (status.hasActualCode()) {
+ code = Cloud.MetaServiceCode.forNumber(status.getActualCode());
+ if (code == null) {
+ if (status.hasCode() && status.getCode() !=
Cloud.MetaServiceCode.OK) {
+ return response;
+ }
+ code = Cloud.MetaServiceCode.UNDEFINED_ERR;
+ }
+ } else {
+ if (status.hasCode()) {
+ return response;
+ }
+ code = Cloud.MetaServiceCode.UNDEFINED_ERR;
}
- Cloud.MetaServiceCode code =
Cloud.MetaServiceCode.forNumber(status.getActualCode());
- if (code == null || code == status.getCode()) {
+ if (status.hasCode() && code == status.getCode()) {
return response;
}
Cloud.MetaServiceResponseStatus restoredStatus =
status.toBuilder().setCode(code).build();
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 31e0887f02e..a9f25586bee 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
@@ -229,6 +229,16 @@ public class MetaServiceProxyTest {
Assert.assertEquals(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber(),
status.getActualCode());
}
+ @Test
+ public void testGetInstanceUsesKnownActualCodeWithoutFallback() throws
RpcException {
+ Cloud.MetaServiceResponseStatus status = callGetInstanceWithStatus(
+ Cloud.MetaServiceResponseStatus.newBuilder()
+
.setActualCode(Cloud.MetaServiceCode.MS_TOO_BUSY.getNumber())
+ .build());
+
+ Assert.assertEquals(Cloud.MetaServiceCode.MS_TOO_BUSY,
status.getCode());
+ }
+
@Test
public void testGetInstanceKeepsLegacyCodeForUnknownActualCode() throws
RpcException {
Cloud.MetaServiceResponseStatus status = callGetInstanceWithStatus(
@@ -241,6 +251,44 @@ public class MetaServiceProxyTest {
Assert.assertEquals(Integer.MAX_VALUE, status.getActualCode());
}
+ @Test
+ public void
testGetInstanceFailsClosedForUnknownActualCodeWithoutErrorFallback() throws
RpcException {
+ Cloud.MetaServiceResponseStatus status = callGetInstanceWithStatus(
+ Cloud.MetaServiceResponseStatus.newBuilder()
+ .setCode(Cloud.MetaServiceCode.OK)
+ .setActualCode(Integer.MAX_VALUE)
+ .build());
+
+ Assert.assertEquals(Cloud.MetaServiceCode.UNDEFINED_ERR,
status.getCode());
+ Assert.assertEquals(Integer.MAX_VALUE, status.getActualCode());
+
+ status =
callGetInstanceWithStatus(Cloud.MetaServiceResponseStatus.newBuilder()
+ .setActualCode(Integer.MAX_VALUE)
+ .build());
+
+ Assert.assertEquals(Cloud.MetaServiceCode.UNDEFINED_ERR,
status.getCode());
+ Assert.assertEquals(Integer.MAX_VALUE, status.getActualCode());
+ }
+
+ @Test
+ public void testGetInstanceFailsClosedWithoutAnyCode() throws RpcException
{
+ Cloud.MetaServiceResponseStatus status = callGetInstanceWithStatus(
+ Cloud.MetaServiceResponseStatus.getDefaultInstance());
+
+ Assert.assertEquals(Cloud.MetaServiceCode.UNDEFINED_ERR,
status.getCode());
+ }
+
+ @Test
+ public void testResponseFailsClosedWithoutStatus() {
+ Cloud.GetInstanceResponse response = Deencapsulation.invoke(
+ MetaServiceClient.class,
+ "restoreActualCode",
+ Cloud.GetInstanceResponse.getDefaultInstance());
+
+ Assert.assertTrue(response.hasStatus());
+ Assert.assertEquals(Cloud.MetaServiceCode.UNDEFINED_ERR,
response.getStatus().getCode());
+ }
+
@Test
public void testGetInstanceKeepsLegacyCodeWithoutActualCode() throws
RpcException {
Cloud.MetaServiceResponseStatus status = callGetInstanceWithStatus(
diff --git a/gensrc/proto/cloud.proto b/gensrc/proto/cloud.proto
index 0b8d2738ca4..e87bd439efd 100644
--- a/gensrc/proto/cloud.proto
+++ b/gensrc/proto/cloud.proto
@@ -1478,8 +1478,9 @@ message MetaServiceResponseStatus {
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`.
+ // New clients should use this field when the local enum descriptor
recognizes the value.
+ // Otherwise, use `code` only when it is explicitly present and non-OK,
and fail closed with
+ // UNDEFINED_ERR when no recognizable error code is available.
optional int32 actual_code = 3;
}
@@ -1823,7 +1824,7 @@ enum MetaServiceCode {
// 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().
+ // clients must be mapped in resolve_response_code_and_msg().
OK = 0;
//Meta service internal error
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]