From 6d9fadfacc90287ce64ec9fb84779d62c4fd9543 Mon Sep 17 00:00:00 2001 From: Yixuan Wang Date: Fri, 5 Jun 2026 16:51:45 +0800 Subject: [PATCH 1/3] 1 --- be/src/cloud/cloud_meta_mgr.cpp | 49 ++++--- .../src/meta-service/injection_point_http.cpp | 6 +- cloud/src/meta-service/meta_service.h | 24 ++-- cloud/src/meta-service/meta_service_helper.h | 52 ++++++-- cloud/test/meta_service_helper_test.cpp | 123 ++++++++++++++++++ cloud/test/meta_service_test.cpp | 8 ++ common/cpp/cloud_proto_util.h | 33 +++++ .../doris/cloud/rpc/MetaServiceProxy.java | 44 ++++++- gensrc/proto/cloud.proto | 22 ++-- 9 files changed, 299 insertions(+), 62 deletions(-) create mode 100644 common/cpp/cloud_proto_util.h diff --git a/be/src/cloud/cloud_meta_mgr.cpp b/be/src/cloud/cloud_meta_mgr.cpp index c4b7f1af68e0c0..f02d2f63343a83 100644 --- a/be/src/cloud/cloud_meta_mgr.cpp +++ b/be/src/cloud/cloud_meta_mgr.cpp @@ -54,6 +54,7 @@ #include "common/config.h" #include "common/logging.h" #include "common/status.h" +#include "cpp/cloud_proto_util.h" #include "cpp/sync_point.h" #include "io/fs/obj_storage_client.h" #include "load/stream_load/stream_load_context.h" @@ -508,21 +509,22 @@ Status retry_rpc(MetaServiceRPC rpc, const Request& req, Response* res, // Record QPS statistics for all RPCs sent to MS (success or failure) record_rpc_qps(rpc, rate_limit_ctx); + MetaServiceCode status_code = get_response_code(res->status()); if (cntl.Failed()) [[unlikely]] { error_msg = cntl.ErrorText(); error_code = cntl.ErrorCode(); proxy->set_unhealthy(); - } else if (res->status().code() == MetaServiceCode::OK) { + } else if (status_code == MetaServiceCode::OK) { return Status::OK(); - } else if (res->status().code() == MetaServiceCode::INVALID_ARGUMENT) { + } else if (status_code == MetaServiceCode::INVALID_ARGUMENT) { return Status::Error("failed to {}: {}", op_name, res->status().msg()); - } else if (res->status().code() == MetaServiceCode::MS_TOO_BUSY) { + } else if (status_code == MetaServiceCode::MS_TOO_BUSY) { // MS_BUSY should also be retried if (rate_limit_ctx.backpressure_handler) { rate_limit_ctx.backpressure_handler->on_ms_busy(); } - } else if (res->status().code() != MetaServiceCode::KV_TXN_CONFLICT) { + } else if (status_code != MetaServiceCode::KV_TXN_CONFLICT) { return Status::Error("failed to {}: {}", op_name, res->status().msg()); } else { @@ -538,7 +540,7 @@ Status retry_rpc(MetaServiceRPC rpc, const Request& req, Response* res, (retry_times > config::meta_service_rpc_timeout_retry_times && error_code == brpc::ERPCTIMEDOUT) || (retry_times > config::meta_service_conflict_error_retry_times && - res->status().code() == MetaServiceCode::KV_TXN_CONFLICT)) { + status_code == MetaServiceCode::KV_TXN_CONFLICT)) { break; } @@ -567,7 +569,7 @@ Status CloudMetaMgr::get_tablet_meta(int64_t tablet_id, TabletMetaSharedPtr* tab .backpressure_handler = ms_backpressure_handler_, }); if (!st.ok()) { - if (resp.status().code() == MetaServiceCode::TABLET_NOT_FOUND) { + if (get_response_code(resp.status()) == MetaServiceCode::TABLET_NOT_FOUND) { return Status::NotFound("failed to get tablet meta: {}", resp.status().msg()); } return st; @@ -750,13 +752,14 @@ Status CloudMetaMgr::sync_tablet_rowsets_unlocked(CloudTablet* tablet, } return Status::RpcError("failed to get rowset meta: {}", cntl.ErrorText()); } - if (resp.status().code() == MetaServiceCode::TABLET_NOT_FOUND) { + MetaServiceCode status_code = get_response_code(resp.status()); + if (status_code == MetaServiceCode::TABLET_NOT_FOUND) { LOG(WARNING) << "failed to get rowset meta, err=" << resp.status().msg() << " " << tablet_info; return Status::NotFound("failed to get rowset meta: {}, {}", resp.status().msg(), tablet_info); } - if (resp.status().code() == MetaServiceCode::MS_TOO_BUSY) { + if (status_code == MetaServiceCode::MS_TOO_BUSY) { // MS_BUSY should also be retried if (ms_backpressure_handler_) { ms_backpressure_handler_->on_ms_busy(); @@ -775,7 +778,7 @@ Status CloudMetaMgr::sync_tablet_rowsets_unlocked(CloudTablet* tablet, } return Status::RpcError("failed to get rowset meta: {}", resp.status().msg()); } - if (resp.status().code() != MetaServiceCode::OK) { + if (status_code != MetaServiceCode::OK) { LOG(WARNING) << " failed to get rowset meta, err=" << resp.status().msg() << " " << tablet_info; return Status::InternalError("failed to get rowset meta: {}, {}", resp.status().msg(), @@ -998,7 +1001,8 @@ Status CloudMetaMgr::_get_delete_bitmap_from_ms(GetDeleteBitmapRequest& req, return st; } - if (res.status().code() == MetaServiceCode::TABLET_NOT_FOUND) { + MetaServiceCode status_code = get_response_code(res.status()); + if (status_code == MetaServiceCode::TABLET_NOT_FOUND) { return Status::NotFound("failed to get delete bitmap: {}", res.status().msg()); } // The delete bitmap of stale rowsets will be removed when commit compaction job, @@ -1021,11 +1025,11 @@ Status CloudMetaMgr::_get_delete_bitmap_from_ms(GetDeleteBitmapRequest& req, // | return get delete bitmap | | // |<---------------------------| | // | | | - if (res.status().code() == MetaServiceCode::ROWSETS_EXPIRED) { + if (status_code == MetaServiceCode::ROWSETS_EXPIRED) { return Status::Error("failed to get delete bitmap: {}", res.status().msg()); } - if (res.status().code() != MetaServiceCode::OK) { + if (status_code != MetaServiceCode::OK) { return Status::Error("failed to get delete bitmap: {}", res.status().msg()); } @@ -1464,7 +1468,7 @@ Status CloudMetaMgr::prepare_rowset(const RowsetMeta& rs_meta, const std::string .backpressure_handler = ms_backpressure_handler_, .table_id = table_id, }); - if (!st.ok() && resp.status().code() == MetaServiceCode::ALREADY_EXISTED) { + if (!st.ok() && get_response_code(resp.status()) == MetaServiceCode::ALREADY_EXISTED) { if (existed_rs_meta != nullptr && resp.has_existed_rowset_meta()) { RowsetMetaPB doris_rs_meta_tmp = cloud_rowset_meta_to_doris(std::move(*resp.mutable_existed_rowset_meta())); @@ -1500,7 +1504,7 @@ Status CloudMetaMgr::commit_rowset(RowsetMeta& rs_meta, const std::string& job_i .backpressure_handler = ms_backpressure_handler_, .table_id = table_id, }); - if (!st.ok() && resp.status().code() == MetaServiceCode::ALREADY_EXISTED) { + if (!st.ok() && get_response_code(resp.status()) == MetaServiceCode::ALREADY_EXISTED) { if (existed_rs_meta != nullptr && resp.has_existed_rowset_meta()) { RowsetMetaPB doris_rs_meta = cloud_rowset_meta_to_doris(std::move(*resp.mutable_existed_rowset_meta())); @@ -1562,7 +1566,7 @@ Status CloudMetaMgr::update_tmp_rowset(const RowsetMeta& rs_meta, int64_t table_ .backpressure_handler = ms_backpressure_handler_, .table_id = table_id, }); - if (!st.ok() && resp.status().code() == MetaServiceCode::ROWSET_META_NOT_FOUND) { + if (!st.ok() && get_response_code(resp.status()) == MetaServiceCode::ROWSET_META_NOT_FOUND) { return Status::InternalError("failed to update committed rowset: {}", resp.status().msg()); } return st; @@ -1844,7 +1848,8 @@ Status CloudMetaMgr::commit_tablet_job(const TabletJobInfoPB& job, FinishTabletJ .host_limiters = host_level_ms_rpc_rate_limiters_, .backpressure_handler = ms_backpressure_handler_, }); - if (res->status().code() == MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES) { + if (get_response_code(res->status()) == + MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES) { return Status::Error( "txn conflict when commit tablet job {}", job.ShortDebugString()); } @@ -2082,13 +2087,14 @@ Status CloudMetaMgr::update_delete_bitmap(const CloudTablet& tablet, int64_t loc .backpressure_handler = ms_backpressure_handler_, .table_id = table_id, }); + MetaServiceCode status_code = get_response_code(res.status()); if (config::enable_update_delete_bitmap_kv_check_core && - res.status().code() == MetaServiceCode::UPDATE_OVERRIDE_EXISTING_KV) { + status_code == MetaServiceCode::UPDATE_OVERRIDE_EXISTING_KV) { auto& msg = res.status().msg(); LOG_WARNING(msg); CHECK(false) << msg; } - if (res.status().code() == MetaServiceCode::LOCK_EXPIRED) { + if (status_code == MetaServiceCode::LOCK_EXPIRED) { return Status::Error( "lock expired when update delete bitmap, tablet_id: {}, lock_id: {}, initiator: " "{}, error_msg: {}", @@ -2186,7 +2192,7 @@ Status CloudMetaMgr::get_delete_bitmap_update_lock(const CloudTablet& tablet, in }); DBUG_EXECUTE_IF("CloudMetaMgr::test_get_delete_bitmap_update_lock_conflict", { test_conflict = true; }); - if (!test_conflict && res.status().code() != MetaServiceCode::LOCK_CONFLICT) { + if (!test_conflict && get_response_code(res.status()) != MetaServiceCode::LOCK_CONFLICT) { break; } @@ -2212,12 +2218,13 @@ Status CloudMetaMgr::get_delete_bitmap_update_lock(const CloudTablet& tablet, in std::this_thread::sleep_for(std::chrono::seconds(sleep_time)); } }); - if (res.status().code() == MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES) { + MetaServiceCode status_code = get_response_code(res.status()); + if (status_code == MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES) { return Status::Error( "txn conflict when get delete bitmap update lock, table_id {}, lock_id {}, " "initiator {}", tablet.table_id(), lock_id, initiator); - } else if (res.status().code() == MetaServiceCode::LOCK_CONFLICT) { + } else if (status_code == MetaServiceCode::LOCK_CONFLICT) { return Status::Error( "lock conflict when get delete bitmap update lock, table_id {}, lock_id {}, " "initiator {}", diff --git a/cloud/src/meta-service/injection_point_http.cpp b/cloud/src/meta-service/injection_point_http.cpp index 3c2ae4fee33a65..7705ed7c5831ee 100644 --- a/cloud/src/meta-service/injection_point_http.cpp +++ b/cloud/src/meta-service/injection_point_http.cpp @@ -131,8 +131,8 @@ static void register_suites() { std::bernoulli_distribution inject_fault {p}; if (inject_fault(gen)) { auto* status = try_any_cast(args[1]); - status->set_code(MetaServiceCode::MS_TOO_BUSY); - status->set_msg("injected ms too busy"); + set_response_status(status, MetaServiceCode::MS_TOO_BUSY, + "injected ms too busy"); LOG_WARNING("inject ms too busy on {} with probability {}", *req_name, p); *try_any_cast(args.back()) = true; } @@ -395,4 +395,4 @@ HttpResponse process_injection_point(MetaServiceImpl* service, brpc::Controller* return http_json_reply(MetaServiceCode::INVALID_ARGUMENT, "unknown op:" + op); } -} // namespace doris::cloud \ No newline at end of file +} // namespace doris::cloud diff --git a/cloud/src/meta-service/meta_service.h b/cloud/src/meta-service/meta_service.h index a0594a945d6f6f..84efd2d5335f37 100644 --- a/cloud/src/meta-service/meta_service.h +++ b/cloud/src/meta-service/meta_service.h @@ -28,9 +28,11 @@ #include #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,19 +1037,19 @@ class MetaServiceProxy final : public MetaService { using namespace std::chrono; brpc::ClosureGuard done_guard(done); + DORIS_CLOUD_DEFER { + auto* status = resp->mutable_status(); + set_response_status(status, get_response_code(*status), status->msg()); + }; + // life span of this defer MUST be longer than `done` std::unique_ptr> defer_injection( (int*)(0x01), [&, this](int*) { idempotent_injection(method, req, resp); }); if (!config::enable_txn_store_retry) { (impl_.get()->*method)(ctrl, req, resp, brpc::DoNothing()); - if (resp->status().code() == MetaServiceCode::KV_TXN_MAYBE_COMMITTED) { - // Keep maybe-committed as an internal retry signal only. Older proto2 - // clients may treat unknown enum values as unset and fall back to OK. - resp->mutable_status()->set_code(MetaServiceCode::KV_TXN_COMMIT_ERR); - } + MetaServiceCode code = get_legacy_code(resp->status().code()); if (DCHECK_IS_ON()) { - MetaServiceCode code = resp->status().code(); DCHECK_NE(code, MetaServiceCode::KV_TXN_STORE_GET_RETRYABLE) << "KV_TXN_STORE_GET_RETRYABLE should not be sent back to client"; DCHECK_NE(code, MetaServiceCode::KV_TXN_STORE_COMMIT_RETRYABLE) @@ -1092,16 +1094,6 @@ class MetaServiceProxy final : public MetaService { if (retry_times >= config::txn_store_retry_times || // Retrying KV_TXN_TOO_OLD is very expensive, so we only retry once. (retry_times > 1 && code == MetaServiceCode::KV_TXN_TOO_OLD)) { - // For KV_TXN_CONFLICT, we should return KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES, - // because BE will retries the KV_TXN_CONFLICT error. - resp->mutable_status()->set_code( - code == MetaServiceCode::KV_TXN_STORE_COMMIT_RETRYABLE ? KV_TXN_COMMIT_ERR - : code == MetaServiceCode::KV_TXN_STORE_GET_RETRYABLE ? KV_TXN_GET_ERR - : code == MetaServiceCode::KV_TXN_STORE_CREATE_RETRYABLE ? KV_TXN_CREATE_ERR - : code == MetaServiceCode::KV_TXN_MAYBE_COMMITTED ? KV_TXN_COMMIT_ERR - : code == MetaServiceCode::KV_TXN_CONFLICT - ? KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES - : MetaServiceCode::KV_TXN_TOO_OLD); return; } diff --git a/cloud/src/meta-service/meta_service_helper.h b/cloud/src/meta-service/meta_service_helper.h index 3ee903e3dfda75..2c8b465c104115 100644 --- a/cloud/src/meta-service/meta_service_helper.h +++ b/cloud/src/meta-service/meta_service_helper.h @@ -25,6 +25,7 @@ #include #include #include +#include #include "common/bvars.h" #include "common/config.h" @@ -33,6 +34,7 @@ #include "common/stats.h" #include "common/stopwatch.h" #include "common/util.h" +#include "cpp/cloud_proto_util.h" #include "cpp/sync_point.h" #include "meta-service/meta_service_rate_limit_helper.h" #include "meta-store/keys.h" @@ -41,6 +43,35 @@ #include "resource-manager/resource_manager.h" namespace doris::cloud { +// MetaServiceResponseStatus has two status-code channels: +// - aux_code: the real code encoded as int32, readable by new clients if their enum knows it. +// - code: a legacy fallback enum value, readable by old proto2 clients without falling back to OK. +inline MetaServiceCode get_legacy_code(MetaServiceCode code) { + switch (code) { + case MetaServiceCode::KV_TXN_STORE_GET_RETRYABLE: + return MetaServiceCode::KV_TXN_GET_ERR; + case MetaServiceCode::KV_TXN_STORE_COMMIT_RETRYABLE: + return MetaServiceCode::KV_TXN_COMMIT_ERR; + case MetaServiceCode::KV_TXN_STORE_CREATE_RETRYABLE: + return MetaServiceCode::KV_TXN_CREATE_ERR; + case MetaServiceCode::KV_TXN_MAYBE_COMMITTED: + return MetaServiceCode::KV_TXN_COMMIT_ERR; + case MetaServiceCode::MS_TOO_BUSY: + return MetaServiceCode::KV_TXN_CONFLICT; + case MetaServiceCode::KV_TXN_CONFLICT: + return MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES; + default: + return code; + } +} + +inline void set_response_status(MetaServiceResponseStatus* status, MetaServiceCode code, + std::string msg) { + status->set_aux_code(static_cast(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 +346,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 73b3d37de5b560..cf00667b5fb77f 100644 --- a/cloud/test/meta_service_helper_test.cpp +++ b/cloud/test/meta_service_helper_test.cpp @@ -15,10 +15,13 @@ // specific language governing permissions and limitations // under the License. +#include "meta-service/meta_service_helper.h" + #include #include #include +#include #include #include "common/config.h" @@ -45,6 +48,126 @@ struct MsRateLimitInjectionConfigGuard { }; } // namespace +TEST(MetaServiceHelperTest, ResponseStatusUsesExactAndLegacyCodes) { + MetaServiceResponseStatus status; + + set_response_status(&status, MetaServiceCode::MS_TOO_BUSY, "busy"); + EXPECT_EQ(status.code(), MetaServiceCode::KV_TXN_CONFLICT); + EXPECT_EQ(status.aux_code(), MetaServiceCode::MS_TOO_BUSY); + EXPECT_EQ(get_response_code(status), MetaServiceCode::MS_TOO_BUSY); + + set_response_status(&status, MetaServiceCode::KV_TXN_MAYBE_COMMITTED, "maybe committed"); + EXPECT_EQ(status.code(), MetaServiceCode::KV_TXN_COMMIT_ERR); + EXPECT_EQ(status.aux_code(), MetaServiceCode::KV_TXN_MAYBE_COMMITTED); + EXPECT_EQ(get_response_code(status), MetaServiceCode::KV_TXN_MAYBE_COMMITTED); +} + +TEST(MetaServiceHelperTest, ResponseStatusCoversEveryMetaServiceCode) { + std::set 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_status(&status, code, ""); + EXPECT_EQ(status.code(), expected_legacy_code) + << "MetaServiceCode: " << MetaServiceCode_Name(code); + EXPECT_EQ(status.aux_code(), static_cast(code)) + << "MetaServiceCode: " << MetaServiceCode_Name(code); + EXPECT_EQ(get_response_code(status), 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_RETRY_EXCEEDED_MAX_TIMES); + 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_GET_ERR); + expect_response_status(MetaServiceCode::KV_TXN_STORE_COMMIT_RETRYABLE, + MetaServiceCode::KV_TXN_COMMIT_ERR); + expect_response_status(MetaServiceCode::KV_TXN_STORE_CREATE_RETRYABLE, + MetaServiceCode::KV_TXN_CREATE_ERR); + expect_response_status(MetaServiceCode::KV_TXN_TOO_OLD, MetaServiceCode::KV_TXN_TOO_OLD); + expect_response_status(MetaServiceCode::KV_TXN_MAYBE_COMMITTED, + MetaServiceCode::KV_TXN_COMMIT_ERR); + 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(MetaServiceCode_descriptor()->value_count())); +} + TEST(MetaServiceHelperTest, FdbClusterPressureNeedsLatencyAndNonWorkload) { MsStressMetrics metrics; metrics.fdb_commit_latency_ns = 51L * 1000 * 1000; diff --git a/cloud/test/meta_service_test.cpp b/cloud/test/meta_service_test.cpp index 6f18493ee2ed3a..c944013e7bc033 100644 --- a/cloud/test/meta_service_test.cpp +++ b/cloud/test/meta_service_test.cpp @@ -8518,6 +8518,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_aux_code()); + EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_MAYBE_COMMITTED); EXPECT_EQ(index, 1); SyncPoint::get_instance()->disable_processing(); @@ -8558,6 +8560,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_aux_code()); + EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_MAYBE_COMMITTED); EXPECT_EQ(resp.version(), 2); EXPECT_EQ(index, 1); } @@ -8597,6 +8601,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_aux_code()); + EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_MAYBE_COMMITTED); EXPECT_GE(index, static_cast(config::txn_store_retry_times + 1)); SyncPoint::get_instance()->disable_processing(); @@ -8641,6 +8647,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_aux_code()); + EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_MAYBE_COMMITTED); EXPECT_EQ(resp.version(), 2); EXPECT_GE(index, static_cast(config::txn_store_retry_times + 1)); } diff --git a/common/cpp/cloud_proto_util.h b/common/cpp/cloud_proto_util.h new file mode 100644 index 00000000000000..a82c32d77fd2a2 --- /dev/null +++ b/common/cpp/cloud_proto_util.h @@ -0,0 +1,33 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#pragma once + +#include + +namespace doris::cloud { + +// Reads the exact status code when this binary knows the enum value. If aux_code contains a +// future enum value unknown to this binary, fall back to the legacy-compatible code field. +inline MetaServiceCode get_response_code(const MetaServiceResponseStatus& status) { + if (status.has_aux_code() && MetaServiceCode_IsValid(status.aux_code())) { + return static_cast(status.aux_code()); + } + return status.code(); +} + +} // namespace doris::cloud diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java index 69db7115c83446..e582b2749c75fe 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java @@ -24,6 +24,8 @@ import org.apache.doris.rpc.RpcException; import com.google.common.collect.Maps; +import com.google.protobuf.Descriptors; +import com.google.protobuf.Message; import io.grpc.StatusRuntimeException; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -111,7 +113,7 @@ public Cloud.GetInstanceResponse getInstance(Cloud.GetInstanceRequest request) CloudMetrics.META_SERVICE_RPC_LATENCY.getOrAdd(methodName) .update(System.currentTimeMillis() - startTime); } - return response; + return preferExactStatusCode(response); } catch (Exception e) { if (MetricRepo.isInit && Config.isCloudMode()) { CloudMetrics.META_SERVICE_RPC_ALL_FAILED.increase(1L); @@ -212,7 +214,7 @@ public Response executeRequest(String methodName, Function Response executeRequest(String methodName, Function Response preferExactStatusCode(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.hasAuxCode()) { + return response; + } + Cloud.MetaServiceCode code = Cloud.MetaServiceCode.forNumber(status.getAuxCode()); + if (code == null || code == status.getCode()) { + return response; + } + Cloud.MetaServiceResponseStatus normalizedStatus = status.toBuilder().setCode(code).build(); + Message.Builder builder = message.toBuilder(); + builder.setField(statusField, normalizedStatus); + return (Response) builder.build(); + } + private final MetaServiceClientWrapper w = new MetaServiceClientWrapper(this); /** @@ -307,6 +343,9 @@ public Future getVisibleVersionAsync(Cloud.GetVersionR if (future instanceof com.google.common.util.concurrent.ListenableFuture) { com.google.common.util.concurrent.ListenableFuture listenableFuture = (com.google.common.util.concurrent.ListenableFuture) future; + listenableFuture = com.google.common.util.concurrent.Futures.transform( + listenableFuture, MetaServiceProxy::preferExactStatusCode, + com.google.common.util.concurrent.MoreExecutors.directExecutor()); MetaServiceClient finalClient = client; com.google.common.util.concurrent.Futures.addCallback(listenableFuture, new com.google.common.util.concurrent.FutureCallback() { @@ -331,6 +370,7 @@ public void onFailure(Throwable t) { } } }, com.google.common.util.concurrent.MoreExecutors.directExecutor()); + return listenableFuture; } return future; } catch (Exception e) { diff --git a/gensrc/proto/cloud.proto b/gensrc/proto/cloud.proto index bff85eee942756..fe8ccc9f3d4b95 100644 --- a/gensrc/proto/cloud.proto +++ b/gensrc/proto/cloud.proto @@ -1455,8 +1455,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 status code encoded as int32, so proto2 clients do not drop unknown enum values. + // New clients should use this field when the local enum descriptor recognizes the value, + // otherwise fall back to `code`. + + optional int32 aux_code = 3; } message MetaServiceHttpRequest { @@ -1793,6 +1799,12 @@ 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 code to `MetaServiceResponseStatus.aux_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 @@ -1810,10 +1822,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 @@ -1877,11 +1886,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; } From 3eccc8c51e21a0133e6912106f643910cca831a4 Mon Sep 17 00:00:00 2001 From: Yixuan Wang Date: Fri, 31 Jul 2026 12:06:00 +0800 Subject: [PATCH 2/3] [fix](be) Normalize meta-service response codes at RPC boundaries ### What problem does this PR solve? Issue Number: None Related PR: None Problem Summary: New meta-service versions expose exact status codes through aux_code while preserving a legacy-compatible code. Updating individual BE readers can miss compaction, schema change, or future call sites. Normalize the response code once after each BE-to-meta-service RPC so all existing status().code() readers observe the exact known code. The direct get_rowset RPC is normalized at its response boundary as well. ### Release note Fix Cloud BE handling of exact meta-service response codes during rolling upgrades. ### Check List (For Author) - Test: Not run (targeted helper assertions added; build was not requested) - Behavior changed: Yes (all BE readers now observe exact known meta-service response codes) - Does this need documentation: No --- be/src/cloud/cloud_meta_mgr.cpp | 50 ++++++++++++------------- cloud/test/meta_service_helper_test.cpp | 2 + common/cpp/cloud_proto_util.h | 4 ++ 3 files changed, 29 insertions(+), 27 deletions(-) diff --git a/be/src/cloud/cloud_meta_mgr.cpp b/be/src/cloud/cloud_meta_mgr.cpp index 260831e0f4a888..0f59addf97cf8f 100644 --- a/be/src/cloud/cloud_meta_mgr.cpp +++ b/be/src/cloud/cloud_meta_mgr.cpp @@ -505,26 +505,26 @@ Status retry_rpc(MetaServiceRPC rpc, const Request& req, Response* res, res->Clear(); int error_code = 0; (stub.get()->*method)(&cntl, &req, res, nullptr); + normalize_response_status(res->mutable_status()); // Record QPS statistics for all RPCs sent to MS (success or failure) record_rpc_qps(rpc, rate_limit_ctx); - MetaServiceCode status_code = get_response_code(res->status()); if (cntl.Failed()) [[unlikely]] { error_msg = cntl.ErrorText(); error_code = cntl.ErrorCode(); proxy->set_unhealthy(); - } else if (status_code == MetaServiceCode::OK) { + } else if (res->status().code() == MetaServiceCode::OK) { return Status::OK(); - } else if (status_code == MetaServiceCode::INVALID_ARGUMENT) { + } else if (res->status().code() == MetaServiceCode::INVALID_ARGUMENT) { return Status::Error("failed to {}: {}", op_name, res->status().msg()); - } else if (status_code == MetaServiceCode::MS_TOO_BUSY) { + } else if (res->status().code() == MetaServiceCode::MS_TOO_BUSY) { // MS_BUSY should also be retried if (rate_limit_ctx.backpressure_handler) { rate_limit_ctx.backpressure_handler->on_ms_busy(); } - } else if (status_code != MetaServiceCode::KV_TXN_CONFLICT) { + } else if (res->status().code() != MetaServiceCode::KV_TXN_CONFLICT) { return Status::Error("failed to {}: {}", op_name, res->status().msg()); } else { @@ -540,7 +540,7 @@ Status retry_rpc(MetaServiceRPC rpc, const Request& req, Response* res, (retry_times > config::meta_service_rpc_timeout_retry_times && error_code == brpc::ERPCTIMEDOUT) || (retry_times > config::meta_service_conflict_error_retry_times && - status_code == MetaServiceCode::KV_TXN_CONFLICT)) { + res->status().code() == MetaServiceCode::KV_TXN_CONFLICT)) { break; } @@ -569,7 +569,7 @@ Status CloudMetaMgr::get_tablet_meta(int64_t tablet_id, TabletMetaSharedPtr* tab .backpressure_handler = ms_backpressure_handler_, }); if (!st.ok()) { - if (get_response_code(resp.status()) == MetaServiceCode::TABLET_NOT_FOUND) { + if (resp.status().code() == MetaServiceCode::TABLET_NOT_FOUND) { return Status::NotFound("failed to get tablet meta: {}", resp.status().msg()); } return st; @@ -732,6 +732,7 @@ Status CloudMetaMgr::sync_tablet_rowsets_unlocked(CloudTablet* tablet, auto start = std::chrono::steady_clock::now(); stub->get_rowset(&cntl, &req, &resp, nullptr); + normalize_response_status(resp.mutable_status()); auto end = std::chrono::steady_clock::now(); int64_t latency = cntl.latency_us(); _get_rowset_latency << latency; @@ -752,14 +753,13 @@ Status CloudMetaMgr::sync_tablet_rowsets_unlocked(CloudTablet* tablet, } return Status::RpcError("failed to get rowset meta: {}", cntl.ErrorText()); } - MetaServiceCode status_code = get_response_code(resp.status()); - if (status_code == MetaServiceCode::TABLET_NOT_FOUND) { + if (resp.status().code() == MetaServiceCode::TABLET_NOT_FOUND) { LOG(WARNING) << "failed to get rowset meta, err=" << resp.status().msg() << " " << tablet_info; return Status::NotFound("failed to get rowset meta: {}, {}", resp.status().msg(), tablet_info); } - if (status_code == MetaServiceCode::MS_TOO_BUSY) { + if (resp.status().code() == MetaServiceCode::MS_TOO_BUSY) { // MS_BUSY should also be retried if (ms_backpressure_handler_) { ms_backpressure_handler_->on_ms_busy(); @@ -778,7 +778,7 @@ Status CloudMetaMgr::sync_tablet_rowsets_unlocked(CloudTablet* tablet, } return Status::RpcError("failed to get rowset meta: {}", resp.status().msg()); } - if (status_code != MetaServiceCode::OK) { + if (resp.status().code() != MetaServiceCode::OK) { LOG(WARNING) << " failed to get rowset meta, err=" << resp.status().msg() << " " << tablet_info; return Status::InternalError("failed to get rowset meta: {}, {}", resp.status().msg(), @@ -1001,8 +1001,7 @@ Status CloudMetaMgr::_get_delete_bitmap_from_ms(GetDeleteBitmapRequest& req, return st; } - MetaServiceCode status_code = get_response_code(res.status()); - if (status_code == MetaServiceCode::TABLET_NOT_FOUND) { + if (res.status().code() == MetaServiceCode::TABLET_NOT_FOUND) { return Status::NotFound("failed to get delete bitmap: {}", res.status().msg()); } // The delete bitmap of stale rowsets will be removed when commit compaction job, @@ -1025,11 +1024,11 @@ Status CloudMetaMgr::_get_delete_bitmap_from_ms(GetDeleteBitmapRequest& req, // | return get delete bitmap | | // |<---------------------------| | // | | | - if (status_code == MetaServiceCode::ROWSETS_EXPIRED) { + if (res.status().code() == MetaServiceCode::ROWSETS_EXPIRED) { return Status::Error("failed to get delete bitmap: {}", res.status().msg()); } - if (status_code != MetaServiceCode::OK) { + if (res.status().code() != MetaServiceCode::OK) { return Status::Error("failed to get delete bitmap: {}", res.status().msg()); } @@ -1468,7 +1467,7 @@ Status CloudMetaMgr::prepare_rowset(const RowsetMeta& rs_meta, const std::string .backpressure_handler = ms_backpressure_handler_, .table_id = table_id, }); - if (!st.ok() && get_response_code(resp.status()) == MetaServiceCode::ALREADY_EXISTED) { + if (!st.ok() && resp.status().code() == MetaServiceCode::ALREADY_EXISTED) { if (existed_rs_meta != nullptr && resp.has_existed_rowset_meta()) { RowsetMetaPB doris_rs_meta_tmp = cloud_rowset_meta_to_doris(std::move(*resp.mutable_existed_rowset_meta())); @@ -1504,7 +1503,7 @@ Status CloudMetaMgr::commit_rowset(RowsetMeta& rs_meta, const std::string& job_i .backpressure_handler = ms_backpressure_handler_, .table_id = table_id, }); - if (!st.ok() && get_response_code(resp.status()) == MetaServiceCode::ALREADY_EXISTED) { + if (!st.ok() && resp.status().code() == MetaServiceCode::ALREADY_EXISTED) { if (existed_rs_meta != nullptr && resp.has_existed_rowset_meta()) { RowsetMetaPB doris_rs_meta = cloud_rowset_meta_to_doris(std::move(*resp.mutable_existed_rowset_meta())); @@ -1566,7 +1565,7 @@ Status CloudMetaMgr::update_tmp_rowset(const RowsetMeta& rs_meta, int64_t table_ .backpressure_handler = ms_backpressure_handler_, .table_id = table_id, }); - if (!st.ok() && get_response_code(resp.status()) == MetaServiceCode::ROWSET_META_NOT_FOUND) { + if (!st.ok() && resp.status().code() == MetaServiceCode::ROWSET_META_NOT_FOUND) { return Status::InternalError("failed to update committed rowset: {}", resp.status().msg()); } return st; @@ -1849,8 +1848,7 @@ Status CloudMetaMgr::commit_tablet_job(const TabletJobInfoPB& job, FinishTabletJ .host_limiters = host_level_ms_rpc_rate_limiters_, .backpressure_handler = ms_backpressure_handler_, }); - if (get_response_code(res->status()) == - MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES) { + if (res->status().code() == MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES) { return Status::Error( "txn conflict when commit tablet job {}", job.ShortDebugString()); } @@ -2088,14 +2086,13 @@ Status CloudMetaMgr::update_delete_bitmap(const CloudTablet& tablet, int64_t loc .backpressure_handler = ms_backpressure_handler_, .table_id = table_id, }); - MetaServiceCode status_code = get_response_code(res.status()); if (config::enable_update_delete_bitmap_kv_check_core && - status_code == MetaServiceCode::UPDATE_OVERRIDE_EXISTING_KV) { + res.status().code() == MetaServiceCode::UPDATE_OVERRIDE_EXISTING_KV) { auto& msg = res.status().msg(); LOG_WARNING(msg); CHECK(false) << msg; } - if (status_code == MetaServiceCode::LOCK_EXPIRED) { + if (res.status().code() == MetaServiceCode::LOCK_EXPIRED) { return Status::Error( "lock expired when update delete bitmap, tablet_id: {}, lock_id: {}, initiator: " "{}, error_msg: {}", @@ -2193,7 +2190,7 @@ Status CloudMetaMgr::get_delete_bitmap_update_lock(const CloudTablet& tablet, in }); DBUG_EXECUTE_IF("CloudMetaMgr::test_get_delete_bitmap_update_lock_conflict", { test_conflict = true; }); - if (!test_conflict && get_response_code(res.status()) != MetaServiceCode::LOCK_CONFLICT) { + if (!test_conflict && res.status().code() != MetaServiceCode::LOCK_CONFLICT) { break; } @@ -2219,13 +2216,12 @@ Status CloudMetaMgr::get_delete_bitmap_update_lock(const CloudTablet& tablet, in std::this_thread::sleep_for(std::chrono::seconds(sleep_time)); } }); - MetaServiceCode status_code = get_response_code(res.status()); - if (status_code == MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES) { + if (res.status().code() == MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES) { return Status::Error( "txn conflict when get delete bitmap update lock, table_id {}, lock_id {}, " "initiator {}", tablet.table_id(), lock_id, initiator); - } else if (status_code == MetaServiceCode::LOCK_CONFLICT) { + } else if (res.status().code() == MetaServiceCode::LOCK_CONFLICT) { return Status::Error( "lock conflict when get delete bitmap update lock, table_id {}, lock_id {}, " "initiator {}", diff --git a/cloud/test/meta_service_helper_test.cpp b/cloud/test/meta_service_helper_test.cpp index cfaa888189a6d8..b2e5c159ce4023 100644 --- a/cloud/test/meta_service_helper_test.cpp +++ b/cloud/test/meta_service_helper_test.cpp @@ -55,6 +55,8 @@ TEST(MetaServiceHelperTest, ResponseStatusUsesExactAndLegacyCodes) { EXPECT_EQ(status.code(), MetaServiceCode::KV_TXN_CONFLICT); EXPECT_EQ(status.aux_code(), MetaServiceCode::MS_TOO_BUSY); EXPECT_EQ(get_response_code(status), MetaServiceCode::MS_TOO_BUSY); + normalize_response_status(&status); + EXPECT_EQ(status.code(), MetaServiceCode::MS_TOO_BUSY); set_response_status(&status, MetaServiceCode::KV_TXN_CONFLICT, "conflict"); EXPECT_EQ(status.code(), MetaServiceCode::KV_TXN_CONFLICT); diff --git a/common/cpp/cloud_proto_util.h b/common/cpp/cloud_proto_util.h index a82c32d77fd2a2..675ad4d2b3c75d 100644 --- a/common/cpp/cloud_proto_util.h +++ b/common/cpp/cloud_proto_util.h @@ -30,4 +30,8 @@ inline MetaServiceCode get_response_code(const MetaServiceResponseStatus& status return status.code(); } +inline void normalize_response_status(MetaServiceResponseStatus* status) { + status->set_code(get_response_code(*status)); +} + } // namespace doris::cloud From b8a9fe77c8c6a40a37007ed24884afd16008e19e Mon Sep 17 00:00:00 2001 From: Yixuan Wang Date: Fri, 31 Jul 2026 12:33:30 +0800 Subject: [PATCH 3/3] 1 --- be/src/cloud/cloud_meta_mgr.cpp | 28 +- be/src/cloud/cloud_meta_mgr.h | 4 + be/test/cloud/cloud_meta_mgr_test.cpp | 23 ++ .../src/meta-service/injection_point_http.cpp | 6 +- cloud/src/meta-service/meta_service.h | 7 +- cloud/src/meta-service/meta_service_helper.h | 7 +- cloud/test/meta_service_helper_test.cpp | 356 ++++++++++++------ cloud/test/meta_service_http_test.cpp | 9 +- cloud/test/meta_service_test.cpp | 59 +-- common/cpp/cloud_proto_util.h | 37 -- .../doris/cloud/rpc/MetaServiceClient.java | 70 +++- .../doris/cloud/rpc/MetaServiceProxy.java | 44 +-- .../doris/cloud/rpc/MetaServiceProxyTest.java | 170 ++++++++- gensrc/proto/cloud.proto | 5 +- 14 files changed, 554 insertions(+), 271 deletions(-) delete mode 100644 common/cpp/cloud_proto_util.h diff --git a/be/src/cloud/cloud_meta_mgr.cpp b/be/src/cloud/cloud_meta_mgr.cpp index 0f59addf97cf8f..53837654443b06 100644 --- a/be/src/cloud/cloud_meta_mgr.cpp +++ b/be/src/cloud/cloud_meta_mgr.cpp @@ -54,7 +54,6 @@ #include "common/config.h" #include "common/logging.h" #include "common/status.h" -#include "cpp/cloud_proto_util.h" #include "cpp/sync_point.h" #include "io/fs/obj_storage_client.h" #include "load/stream_load/stream_load_context.h" @@ -154,9 +153,20 @@ Status bthread_fork_join(std::vector>&& 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(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 g_cloud_meta_mgr_rpc_timeout_count("cloud_meta_mgr_rpc_timeout_count"); @@ -419,6 +429,16 @@ using MetaServiceMethod = void (MetaService_Stub::*)(::google::protobuf::RpcCont const Request*, Response*, ::google::protobuf::Closure*); +template +void call_ms(MetaService_Stub* stub, MetaServiceMethod 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}; @@ -504,8 +524,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); - normalize_response_status(res->mutable_status()); + 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); @@ -731,8 +750,7 @@ Status CloudMetaMgr::sync_tablet_rowsets_unlocked(CloudTablet* tablet, } auto start = std::chrono::steady_clock::now(); - stub->get_rowset(&cntl, &req, &resp, nullptr); - normalize_response_status(resp.mutable_status()); + 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 657265594bcd62..9084df3d6ea42a 100644 --- a/be/src/cloud/cloud_meta_mgr.h +++ b/be/src/cloud/cloud_meta_mgr.h @@ -64,6 +64,10 @@ Status bthread_fork_join(const std::vector>& tasks, int Status bthread_fork_join(std::vector>&& tasks, int concurrency, std::future* 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 c6e99cb26b1da6..ff87378348ef10 100644 --- a/be/test/cloud/cloud_meta_mgr_test.cpp +++ b/be/test/cloud/cloud_meta_mgr_test.cpp @@ -21,6 +21,8 @@ #include #include +#include +#include #include #include #include @@ -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(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(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::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/injection_point_http.cpp b/cloud/src/meta-service/injection_point_http.cpp index 7705ed7c5831ee..3c2ae4fee33a65 100644 --- a/cloud/src/meta-service/injection_point_http.cpp +++ b/cloud/src/meta-service/injection_point_http.cpp @@ -131,8 +131,8 @@ static void register_suites() { std::bernoulli_distribution inject_fault {p}; if (inject_fault(gen)) { auto* status = try_any_cast(args[1]); - set_response_status(status, MetaServiceCode::MS_TOO_BUSY, - "injected ms too busy"); + status->set_code(MetaServiceCode::MS_TOO_BUSY); + status->set_msg("injected ms too busy"); LOG_WARNING("inject ms too busy on {} with probability {}", *req_name, p); *try_any_cast(args.back()) = true; } @@ -395,4 +395,4 @@ HttpResponse process_injection_point(MetaServiceImpl* service, brpc::Controller* return http_json_reply(MetaServiceCode::INVALID_ARGUMENT, "unknown op:" + op); } -} // namespace doris::cloud +} // namespace doris::cloud \ No newline at end of file diff --git a/cloud/src/meta-service/meta_service.h b/cloud/src/meta-service/meta_service.h index b4b437ea770292..70fb410e1176c8 100644 --- a/cloud/src/meta-service/meta_service.h +++ b/cloud/src/meta-service/meta_service.h @@ -1039,7 +1039,7 @@ class MetaServiceProxy final : public MetaService { DORIS_CLOUD_DEFER { auto* status = resp->mutable_status(); - set_response_status(status, get_response_code(*status), status->msg()); + set_response_code(status, status->code(), status->msg()); }; // life span of this defer MUST be longer than `done` @@ -1049,6 +1049,8 @@ class MetaServiceProxy final : public MetaService { if (!config::enable_txn_store_retry) { (impl_.get()->*method)(ctrl, req, resp, brpc::DoNothing()); if (resp->status().code() == MetaServiceCode::KV_TXN_MAYBE_COMMITTED) { + // Keep maybe-committed as an internal retry signal only. Older proto2 + // clients may treat unknown enum values as unset and fall back to OK. resp->mutable_status()->set_code(MetaServiceCode::KV_TXN_COMMIT_ERR); } if (DCHECK_IS_ON()) { @@ -1097,7 +1099,8 @@ class MetaServiceProxy final : public MetaService { if (retry_times >= config::txn_store_retry_times || // Retrying KV_TXN_TOO_OLD is very expensive, so we only retry once. (retry_times > 1 && code == MetaServiceCode::KV_TXN_TOO_OLD)) { - // Convert internal retry signals only after MetaService stops retrying. + // For KV_TXN_CONFLICT, we should return KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES, + // because BE will retries the KV_TXN_CONFLICT error. resp->mutable_status()->set_code( code == MetaServiceCode::KV_TXN_STORE_COMMIT_RETRYABLE ? KV_TXN_COMMIT_ERR : code == MetaServiceCode::KV_TXN_STORE_GET_RETRYABLE ? KV_TXN_GET_ERR diff --git a/cloud/src/meta-service/meta_service_helper.h b/cloud/src/meta-service/meta_service_helper.h index f9cc8b121e3ac9..9b1fbed4648436 100644 --- a/cloud/src/meta-service/meta_service_helper.h +++ b/cloud/src/meta-service/meta_service_helper.h @@ -34,7 +34,6 @@ #include "common/stats.h" #include "common/stopwatch.h" #include "common/util.h" -#include "cpp/cloud_proto_util.h" #include "cpp/sync_point.h" #include "meta-service/meta_service_rate_limit_helper.h" #include "meta-store/keys.h" @@ -54,9 +53,9 @@ inline MetaServiceCode get_legacy_code(MetaServiceCode code) { } } -inline void set_response_status(MetaServiceResponseStatus* status, MetaServiceCode code, - std::string msg) { - status->set_aux_code(static_cast(code)); +inline void set_response_code(MetaServiceResponseStatus* status, MetaServiceCode code, + std::string msg) { + status->set_actual_code(static_cast(code)); status->set_code(get_legacy_code(code)); status->set_msg(std::move(msg)); } diff --git a/cloud/test/meta_service_helper_test.cpp b/cloud/test/meta_service_helper_test.cpp index b2e5c159ce4023..7b50792f88b631 100644 --- a/cloud/test/meta_service_helper_test.cpp +++ b/cloud/test/meta_service_helper_test.cpp @@ -17,11 +17,15 @@ #include "meta-service/meta_service_helper.h" +#include +#include #include #include +#include #include #include +#include #include #include "common/config.h" @@ -46,128 +50,66 @@ struct MsRateLimitInjectionConfigGuard { bool original_enable {config::enable_ms_rate_limit_injection}; int32_t original_probability {config::ms_rate_limit_injection_probability}; }; -} // namespace -TEST(MetaServiceHelperTest, ResponseStatusUsesExactAndLegacyCodes) { - MetaServiceResponseStatus status; +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"); - set_response_status(&status, MetaServiceCode::MS_TOO_BUSY, "busy"); - EXPECT_EQ(status.code(), MetaServiceCode::KV_TXN_CONFLICT); - EXPECT_EQ(status.aux_code(), MetaServiceCode::MS_TOO_BUSY); - EXPECT_EQ(get_response_code(status), MetaServiceCode::MS_TOO_BUSY); - normalize_response_status(&status); - EXPECT_EQ(status.code(), MetaServiceCode::MS_TOO_BUSY); + 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); - set_response_status(&status, MetaServiceCode::KV_TXN_CONFLICT, "conflict"); - EXPECT_EQ(status.code(), MetaServiceCode::KV_TXN_CONFLICT); - EXPECT_EQ(status.aux_code(), MetaServiceCode::KV_TXN_CONFLICT); - EXPECT_EQ(get_response_code(status), MetaServiceCode::KV_TXN_CONFLICT); + 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 -TEST(MetaServiceHelperTest, ResponseStatusCoversEveryMetaServiceCode) { - std::set 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_status(&status, code, ""); - EXPECT_EQ(status.code(), expected_legacy_code) - << "MetaServiceCode: " << MetaServiceCode_Name(code); - EXPECT_EQ(status.aux_code(), static_cast(code)) - << "MetaServiceCode: " << MetaServiceCode_Name(code); - EXPECT_EQ(get_response_code(status), code) - << "MetaServiceCode: " << MetaServiceCode_Name(code); - }; +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); + } - 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); + std::unique_ptr new_legacy_status() const { + return std::unique_ptr(legacy_status_prototype_->New()); + } - EXPECT_EQ(covered_codes.size(), - static_cast(MetaServiceCode_descriptor()->value_count())); -} + 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; @@ -272,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(¤t_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 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(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(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 7452fed53ec417..c043a99f5efa3e 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("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("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 b7c73e76c41045..e58fae0d3a65da 100644 --- a/cloud/test/meta_service_test.cpp +++ b/cloud/test/meta_service_test.cpp @@ -8462,55 +8462,12 @@ TEST(MetaServiceTxnStoreRetryableTest, DoNotReturnRetryableCode) { ASSERT_EQ(resp.status().code(), MetaServiceCode::KV_TXN_GET_ERR) << " status is " << resp.status().msg() << ", code=" << resp.status().code(); - ASSERT_TRUE(resp.status().has_aux_code()); - EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_GET_ERR); SyncPoint::get_instance()->disable_processing(); SyncPoint::get_instance()->clear_all_call_backs(); config::txn_store_retry_times = retry_times; } -TEST(MetaServiceTxnStoreRetryableTest, ConflictRetryExhaustionReturnsFinalCode) { - size_t index = 0; - auto* sync_point = SyncPoint::get_instance(); - sync_point->set_call_back("get_version_code", [&](auto&& args) { - ++index; - *doris::try_any_cast(args[0]) = MetaServiceCode::KV_TXN_CONFLICT; - }); - sync_point->enable_processing(); - int32_t retry_times = config::txn_store_retry_times; - bool enable_retry_txn_conflict = config::enable_retry_txn_conflict; - DORIS_CLOUD_DEFER { - config::txn_store_retry_times = retry_times; - config::enable_retry_txn_conflict = enable_retry_txn_conflict; - sync_point->disable_processing(); - sync_point->clear_all_call_backs(); - }; - config::txn_store_retry_times = 2; - config::enable_retry_txn_conflict = true; - - auto service = get_meta_service(); - create_tablet(service.get(), 1, 1, 1, 1); - insert_rowset(service.get(), 1, "conflict_retry_exhaustion", 1, 1, 1); - - brpc::Controller ctrl; - GetVersionRequest req; - req.set_cloud_unique_id("test_cloud_unique_id"); - req.set_db_id(1); - req.set_table_id(1); - req.set_partition_id(1); - - GetVersionResponse resp; - service->get_version(&ctrl, &req, &resp, nullptr); - - EXPECT_EQ(resp.status().code(), MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES); - ASSERT_TRUE(resp.status().has_aux_code()); - EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES); - EXPECT_EQ(get_response_code(resp.status()), - MetaServiceCode::KV_TXN_CONFLICT_RETRY_EXCEEDED_MAX_TIMES); - EXPECT_GE(index, static_cast(config::txn_store_retry_times + 1)); -} - TEST(MetaServiceTxnStoreRetryableTest, CastAsPreservesMaybeCommittedForProxyRetry) { bool enable_retry = config::enable_txn_store_retry; DORIS_CLOUD_DEFER { @@ -8571,8 +8528,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_aux_code()); - EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_COMMIT_ERR); + 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(); @@ -8613,8 +8570,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_aux_code()); - EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_COMMIT_ERR); + 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); } @@ -8654,8 +8611,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_aux_code()); - EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_COMMIT_ERR); + ASSERT_TRUE(resp.status().has_actual_code()); + EXPECT_EQ(resp.status().actual_code(), MetaServiceCode::KV_TXN_COMMIT_ERR); EXPECT_GE(index, static_cast(config::txn_store_retry_times + 1)); SyncPoint::get_instance()->disable_processing(); @@ -8700,8 +8657,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_aux_code()); - EXPECT_EQ(resp.status().aux_code(), MetaServiceCode::KV_TXN_COMMIT_ERR); + 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(config::txn_store_retry_times + 1)); } diff --git a/common/cpp/cloud_proto_util.h b/common/cpp/cloud_proto_util.h deleted file mode 100644 index 675ad4d2b3c75d..00000000000000 --- a/common/cpp/cloud_proto_util.h +++ /dev/null @@ -1,37 +0,0 @@ -// Licensed to the Apache Software Foundation (ASF) under one -// or more contributor license agreements. See the NOTICE file -// distributed with this work for additional information -// regarding copyright ownership. The ASF licenses this file -// to you under the Apache License, Version 2.0 (the -// "License"); you may not use this file except in compliance -// with the License. You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, -// software distributed under the License is distributed on an -// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY -// KIND, either express or implied. See the License for the -// specific language governing permissions and limitations -// under the License. - -#pragma once - -#include - -namespace doris::cloud { - -// Reads the exact status code when this binary knows the enum value. If aux_code contains a -// future enum value unknown to this binary, fall back to the legacy-compatible code field. -inline MetaServiceCode get_response_code(const MetaServiceResponseStatus& status) { - if (status.has_aux_code() && MetaServiceCode_IsValid(status.aux_code())) { - return static_cast(status.aux_code()); - } - return status.code(); -} - -inline void normalize_response_status(MetaServiceResponseStatus* status) { - status->set_code(get_response_code(*status)); -} - -} // namespace doris::cloud 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 7e02b530a060f1..7fbb70900f0f07 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 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 org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -72,11 +83,66 @@ public MetaServiceClient(String address) { Preconditions.checkNotNull(serviceConfig, "serviceConfig is null"); channelConfigVersion = CHANNEL_PROVIDER.currentConfigVersion(); channel = CHANNEL_PROVIDER.createChannel(target); - 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 ClientCall interceptCall( + MethodDescriptor method, CallOptions callOptions, Channel next) { + ClientCall call = next.newCall(method, callOptions); + return new ForwardingClientCall.SimpleForwardingClientCall<>(call) { + @Override + public void start(Listener listener, Metadata headers) { + Listener normalizingListener = + new ForwardingClientCallListener.SimpleForwardingClientCallListener<>(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 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/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java index c0e1a416a48e11..6b81f084717afe 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java @@ -25,8 +25,6 @@ import org.apache.doris.rpc.RpcException; import com.google.common.collect.Maps; -import com.google.protobuf.Descriptors; -import com.google.protobuf.Message; import io.grpc.StatusRuntimeException; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -116,7 +114,7 @@ public Cloud.GetInstanceResponse getInstance(Cloud.GetInstanceRequest request) CloudMetrics.META_SERVICE_RPC_LATENCY.getOrAdd(methodName) .update(System.currentTimeMillis() - startTime); } - return preferExactStatusCode(response); + return response; } catch (MetaServiceRateLimitException e) { recordRpcRateLimited(methodName); throw e; @@ -264,7 +262,7 @@ public Response executeRequest(String methodName, Function Response executeRequest(String methodName, Function Response preferExactStatusCode(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.hasAuxCode()) { - return response; - } - Cloud.MetaServiceCode code = Cloud.MetaServiceCode.forNumber(status.getAuxCode()); - if (code == null || code == status.getCode()) { - return response; - } - Cloud.MetaServiceResponseStatus normalizedStatus = status.toBuilder().setCode(code).build(); - Message.Builder builder = message.toBuilder(); - builder.setField(statusField, normalizedStatus); - return (Response) builder.build(); - } - private final MetaServiceClientWrapper w = new MetaServiceClientWrapper(this); /** @@ -403,9 +367,6 @@ public Future getVisibleVersionAsync(Cloud.GetVersionR if (future instanceof com.google.common.util.concurrent.ListenableFuture) { com.google.common.util.concurrent.ListenableFuture listenableFuture = (com.google.common.util.concurrent.ListenableFuture) future; - listenableFuture = com.google.common.util.concurrent.Futures.transform( - listenableFuture, MetaServiceProxy::preferExactStatusCode, - com.google.common.util.concurrent.MoreExecutors.directExecutor()); MetaServiceClient finalClient = client; com.google.common.util.concurrent.Futures.addCallback(listenableFuture, new com.google.common.util.concurrent.FutureCallback() { @@ -425,7 +386,6 @@ public void onFailure(Throwable t) { } } }, com.google.common.util.concurrent.MoreExecutors.directExecutor()); - return listenableFuture; } return future; } catch (MetaServiceRateLimitException e) { 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 7638bfa774d581..31e0887f02e24c 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.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 java.util.Map; import java.util.Queue; +import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicInteger; public class MetaServiceProxyTest { @@ -190,11 +194,15 @@ public void testExecuteRequestRetryOnTooBusy() throws RpcException { 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)) @@ -209,6 +217,109 @@ public void testExecuteRequestRetryOnTooBusy() throws RpcException { 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 rpcFuture = SettableFuture.create(); + Mockito.when(client.getVisibleVersionAsync(Mockito.any())).thenReturn(rpcFuture); + + Future 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; @@ -377,6 +488,55 @@ private void consumeRateLimitPermits(MetaServiceProxy proxy, String methodName) } } + 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 90919c9971676e..b3a1de924f4f99 100644 --- a/gensrc/proto/cloud.proto +++ b/gensrc/proto/cloud.proto @@ -1464,8 +1464,7 @@ message MetaServiceResponseStatus { // 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 aux_code = 3; + optional int32 actual_code = 3; } message MetaServiceHttpRequest { @@ -1806,7 +1805,7 @@ enum MetaServiceCode { // UnknownFieldSet, so old clients read an unset `MetaServiceResponseStatus.code` as OK. // MetaService must write the exact client-visible code to - // `MetaServiceResponseStatus.aux_code` and write only a legacy fallback 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;