已合并
refactor(routing): 收敛 worker 选择策略接口 #2117
Qin Yufei创建于 6 天前
refactor(routing): 收敛 worker 选择策略接口 #2117
已合并
共 14 个文件变更+71-135
| @@ -14,7 +14,7 @@ ds_cc_library( | |||
| 14 | deps = [ | 14 | deps = [ |
| 15 | "//include/datasystem/utils:utils_headers", | 15 | "//include/datasystem/utils:utils_headers", |
| 16 | "//src/datasystem/client:transport_exist_types", | 16 | "//src/datasystem/client:transport_exist_types", |
| 17 | - "//src/datasystem/client/routing:select_strategy", | 17 | + "//src/datasystem/client/routing:data_placement_policy", |
| 18 | "//src/datasystem/common/metrics:common_metrics", | 18 | "//src/datasystem/common/metrics:common_metrics", |
| 19 | "//src/datasystem/common/util:format", | 19 | "//src/datasystem/common/util:format", |
| 20 | "//src/datasystem/common/util:net_util", | 20 | "//src/datasystem/common/util:net_util", |
| @@ -335,12 +335,12 @@ ExistHandler::ExistHandler(IExistRouting *routing, IExistTransport *transport, s | |||
| 335 | { | 335 | { |
| 336 | } | 336 | } |
| 337 | 337 | ||
| 338 | -Status ExistHandler::SelectWorkers(const std::vector<std::string> &keys, client::SelectStrategy strategy, | 338 | +Status ExistHandler::SelectWorkers(const std::vector<std::string> &keys, client::DataPlacementPolicy policy, |
| 339 | std::unordered_map<HostPort, std::vector<std::string>> &groups, | 339 | std::unordered_map<HostPort, std::vector<std::string>> &groups, |
| 340 | const std::vector<HostPort> &exclude) | 340 | const std::vector<HostPort> &exclude) |
| 341 | { | 341 | { |
| 342 | RETURN_RUNTIME_ERROR_IF_NULL(routing_); | 342 | RETURN_RUNTIME_ERROR_IF_NULL(routing_); |
| 343 | - return routing_->SelectWorkers(keys, strategy, groups, exclude); | 343 | + return routing_->SelectWorkers(keys, policy, groups, exclude); |
| 344 | } | 344 | } |
| 345 | 345 | ||
| 346 | void ExistHandler::UpdateRoutingState(const HostPort &addr, StatusCode status) | 346 | void ExistHandler::UpdateRoutingState(const HostPort &addr, StatusCode status) |
| @@ -373,7 +373,7 @@ Status ExistHandler::RunSelectedWorkers(const ExistHandlerRequest &request, | |||
| 373 | std::vector<bool> &exists) | 373 | std::vector<bool> &exists) |
| 374 | { | 374 | { |
| 375 | std::unordered_map<HostPort, std::vector<std::string>> groups; | 375 | std::unordered_map<HostPort, std::vector<std::string>> groups; |
| 376 | - RETURN_IF_NOT_OK(SelectWorkers(request.keys, client::SelectStrategy::HASH_RING_AFFINITY, groups, | 376 | + RETURN_IF_NOT_OK(SelectWorkers(request.keys, client::DataPlacementPolicy::PREFERRED_META_OWNER, groups, |
| 377 | excludedWorkers_)); | 377 | excludedWorkers_)); |
| 378 | std::vector<std::pair<HostPort, std::vector<std::string>>> orderedGroups(groups.begin(), groups.end()); | 378 | std::vector<std::pair<HostPort, std::vector<std::string>>> orderedGroups(groups.begin(), groups.end()); |
| 379 | std::sort(orderedGroups.begin(), orderedGroups.end(), [](const auto &lhs, const auto &rhs) { | 379 | std::sort(orderedGroups.begin(), orderedGroups.end(), [](const auto &lhs, const auto &rhs) { |
| @@ -24,7 +24,7 @@ | |||
| 24 | 24 | ||
| 25 | 25 | ||
| 26 | 26 | ||
| 27 | -#include "datasystem/client/routing/select_strategy.h" | 27 | +#include "datasystem/client/routing/data_placement_policy.h" |
| 28 | 28 | ||
| 29 | 29 | ||
| 30 | 30 | ||
| @@ -37,7 +37,7 @@ namespace object_cache { | |||
| 37 | class IExistRouting { | 37 | class IExistRouting { |
| 38 | public: | 38 | public: |
| 39 | virtual ~IExistRouting() = default; | 39 | virtual ~IExistRouting() = default; |
| 40 | - virtual Status SelectWorkers(const std::vector<std::string> &keys, client::SelectStrategy strategy, | 40 | + virtual Status SelectWorkers(const std::vector<std::string> &keys, client::DataPlacementPolicy policy, |
| 41 | std::unordered_map<HostPort, std::vector<std::string>> &groups, | 41 | std::unordered_map<HostPort, std::vector<std::string>> &groups, |
| 42 | const std::vector<HostPort> &exclude = {}) = 0; | 42 | const std::vector<HostPort> &exclude = {}) = 0; |
| 43 | virtual void UpdateState(const HostPort &addr, StatusCode status) = 0; | 43 | virtual void UpdateState(const HostPort &addr, StatusCode status) = 0; |
| @@ -114,7 +114,7 @@ private: | |||
| 114 | std::vector<bool> &exists; | 114 | std::vector<bool> &exists; |
| 115 | }; | 115 | }; |
| 116 | 116 | ||
| 117 | - Status SelectWorkers(const std::vector<std::string> &keys, client::SelectStrategy strategy, | 117 | + Status SelectWorkers(const std::vector<std::string> &keys, client::DataPlacementPolicy policy, |
| 118 | std::unordered_map<HostPort, std::vector<std::string>> &groups, | 118 | std::unordered_map<HostPort, std::vector<std::string>> &groups, |
| 119 | const std::vector<HostPort> &exclude); | 119 | const std::vector<HostPort> &exclude); |
| 120 | 120 | ||
| @@ -591,12 +591,12 @@ public: | |||
| 591 | explicit RoutingExistAdapter(std::shared_ptr<client::Routing> routing) : routing_(std::move(routing)) {} | 591 | explicit RoutingExistAdapter(std::shared_ptr<client::Routing> routing) : routing_(std::move(routing)) {} |
| 592 | ~RoutingExistAdapter() override = default; | 592 | ~RoutingExistAdapter() override = default; |
| 593 | 593 | ||
| 594 | - Status SelectWorkers(const std::vector<std::string> &keys, client::SelectStrategy strategy, | 594 | + Status SelectWorkers(const std::vector<std::string> &keys, client::DataPlacementPolicy policy, |
| 595 | std::unordered_map<HostPort, std::vector<std::string>> &groups, | 595 | std::unordered_map<HostPort, std::vector<std::string>> &groups, |
| 596 | const std::vector<HostPort> &exclude) override | 596 | const std::vector<HostPort> &exclude) override |
| 597 | { | 597 | { |
| 598 | RETURN_RUNTIME_ERROR_IF_NULL(routing_); | 598 | RETURN_RUNTIME_ERROR_IF_NULL(routing_); |
| 599 | - return routing_->SelectWorkers(keys, strategy, groups, exclude); | 599 | + return routing_->SelectWorkers(keys, policy, groups, exclude); |
| 600 | } | 600 | } |
| 601 | 601 | ||
| 602 | void UpdateState(const HostPort &addr, StatusCode status) override | 602 | void UpdateState(const HostPort &addr, StatusCode status) override |
| @@ -3,8 +3,8 @@ load("//bazel:build_defs.bzl", "ds_cc_library") | |||
| 3 | package(default_visibility = ["//visibility:public"]) | 3 | package(default_visibility = ["//visibility:public"]) |
| 4 | 4 | ||
| 5 | ds_cc_library( | 5 | ds_cc_library( |
| 6 | - name = "select_strategy", | 6 | + name = "data_placement_policy", |
| 7 | - hdrs = ["select_strategy.h"], | 7 | + hdrs = ["data_placement_policy.h"], |
| 8 | ) | 8 | ) |
| 9 | 9 | ||
| 10 | ds_cc_library( | 10 | ds_cc_library( |
| @@ -29,7 +29,7 @@ ds_cc_library( | |||
| 29 | "worker_router.h", | 29 | "worker_router.h", |
| 30 | ], | 30 | ], |
| 31 | deps = [ | 31 | deps = [ |
| 32 | - ":select_strategy", | 32 | + ":data_placement_policy", |
| 33 | "//include/datasystem/utils:utils_headers", | 33 | "//include/datasystem/utils:utils_headers", |
| 34 | "//src/datasystem/common/ak_sk:signature", | 34 | "//src/datasystem/common/ak_sk:signature", |
| 35 | "//src/datasystem/common/inject:common_inject", | 35 | "//src/datasystem/common/inject:common_inject", |
Rsrc/datasystem/client/routing/select_strategy.h→src/datasystem/client/routing/data_placement_policy.h+9-7
| @@ -14,19 +14,21 @@ | |||
| 14 | * limitations under the License. | 14 | * limitations under the License. |
| 15 | */ | 15 | */ |
| 16 | 16 | ||
| 17 | -#ifndef DATASYSTEM_CLIENT_ROUTING_SELECT_STRATEGY_H | 17 | +#ifndef DATASYSTEM_CLIENT_ROUTING_DATA_PLACEMENT_POLICY_H |
| 18 | -#define DATASYSTEM_CLIENT_ROUTING_SELECT_STRATEGY_H | 18 | +#define DATASYSTEM_CLIENT_ROUTING_DATA_PLACEMENT_POLICY_H |
| 19 | + | ||
| 20 | + | ||
| 19 | 21 | ||
| 20 | namespace datasystem { | 22 | namespace datasystem { |
| 21 | namespace client { | 23 | namespace client { |
| 22 | 24 | ||
| 23 | -// Legacy compatibility enum. New routing code must use DataPlacementPolicy so REQUIRED_SAME_NODE remains expressible. | 25 | +enum class DataPlacementPolicy : uint8_t { |
| 24 | -enum class SelectStrategy { | 26 | + PREFERRED_SAME_NODE, // Prefer a same-node worker, then fall back to the metadata owner. |
| 25 | - HASH_RING_AFFINITY, // Select by key hash on ring (metadata owner) | 27 | + REQUIRED_SAME_NODE, // Select only from same-node workers. |
| 26 | - SAME_NODE_PREFERRED, // Prefer same-node worker, fallback to ring | 28 | + PREFERRED_META_OWNER, // Prefer the metadata owner selected by the hash ring. |
| 27 | }; | 29 | }; |
| 28 | 30 | ||
| 29 | } // namespace client | 31 | } // namespace client |
| 30 | } // namespace datasystem | 32 | } // namespace datasystem |
| 31 | 33 | ||
| 32 | -#endif // DATASYSTEM_CLIENT_ROUTING_SELECT_STRATEGY_H | 34 | +#endif // DATASYSTEM_CLIENT_ROUTING_DATA_PLACEMENT_POLICY_H |
| @@ -131,13 +131,6 @@ Status Routing::FetchHashRing(const HostPort &workerAddr, uint64_t currentVersio | |||
| 131 | return Status::OK(); | 131 | return Status::OK(); |
| 132 | } | 132 | } |
| 133 | 133 | ||
| 134 | -Status Routing::SelectWorker(const std::string &key, SelectStrategy strategy, HostPort &worker, | ||
| 135 | - const std::vector<HostPort> &exclude) | ||
| 136 | -{ | ||
| 137 | - CHECK_FAIL_RETURN_STATUS(initialized_.load(), K_NOT_READY, "Routing is not initialized"); | ||
| 138 | - return router_->SelectWorker(key, strategy, worker, exclude); | ||
| 139 | -} | ||
| 140 | - | ||
| 141 | Status Routing::SelectWorker(const std::string &key, DataPlacementPolicy policy, HostPort &worker, | 134 | Status Routing::SelectWorker(const std::string &key, DataPlacementPolicy policy, HostPort &worker, |
| 142 | const std::vector<HostPort> &exclude) | 135 | const std::vector<HostPort> &exclude) |
| 143 | { | 136 | { |
| @@ -145,22 +138,6 @@ Status Routing::SelectWorker(const std::string &key, DataPlacementPolicy policy, | |||
| 145 | return router_->SelectWorker(key, policy, worker, exclude); | 138 | return router_->SelectWorker(key, policy, worker, exclude); |
| 146 | } | 139 | } |
| 147 | 140 | ||
| 148 | -Status Routing::SelectWorkers(const std::vector<std::string> &keys, SelectStrategy strategy, | ||
| 149 | - std::unordered_map<HostPort, std::vector<std::string>> &groups) | ||
| 150 | -{ | ||
| 151 | - CHECK_FAIL_RETURN_STATUS(initialized_.load(), K_NOT_READY, "Routing is not initialized"); | ||
| 152 | - return router_->SelectWorkers(keys, strategy, groups); | ||
| 153 | -} | ||
| 154 | - | ||
| 155 | -Status Routing::SelectWorkers(const std::vector<std::string> &keys, SelectStrategy strategy, | ||
| 156 | - std::unordered_map<HostPort, std::vector<std::string>> &groups, | ||
| 157 | - const std::vector<HostPort> &exclude) | ||
| 158 | -{ | ||
| 159 | - auto policy = strategy == SelectStrategy::SAME_NODE_PREFERRED ? DataPlacementPolicy::PREFERRED_SAME_NODE | ||
| 160 | - : DataPlacementPolicy::PREFERRED_META_OWNER; | ||
| 161 | - return SelectWorkers(keys, policy, groups, exclude); | ||
| 162 | -} | ||
| 163 | - | ||
| 164 | Status Routing::SelectWorkers(const std::vector<std::string> &keys, DataPlacementPolicy policy, | 141 | Status Routing::SelectWorkers(const std::vector<std::string> &keys, DataPlacementPolicy policy, |
| 165 | std::unordered_map<HostPort, std::vector<std::string>> &groups, | 142 | std::unordered_map<HostPort, std::vector<std::string>> &groups, |
| 166 | const std::vector<HostPort> &exclude) | 143 | const std::vector<HostPort> &exclude) |
| @@ -28,6 +28,7 @@ | |||
| 28 | 28 | ||
| 29 | 29 | ||
| 30 | 30 | ||
| 31 | + | ||
| 31 | 32 | ||
| 32 | 33 | ||
| 33 | 34 | ||
| @@ -39,12 +40,6 @@ | |||
| 39 | namespace datasystem { | 40 | namespace datasystem { |
| 40 | namespace client { | 41 | namespace client { |
| 41 | 42 | ||
| 42 | -enum class DataPlacementPolicy : uint8_t { | ||
| 43 | - PREFERRED_SAME_NODE, // Prefer a same-node worker, then fall back to the metadata owner. | ||
| 44 | - REQUIRED_SAME_NODE, // Select only from same-node workers. | ||
| 45 | - PREFERRED_META_OWNER, // Prefer the metadata owner selected by the hash ring. | ||
| 46 | -}; | ||
| 47 | - | ||
| 48 | class Routing { | 43 | class Routing { |
| 49 | public: | 44 | public: |
| 50 | Routing(BrpcChannelConfig channelConfig, std::shared_ptr<Signature> signature, | 45 | Routing(BrpcChannelConfig channelConfig, std::shared_ptr<Signature> signature, |
| @@ -59,23 +54,10 @@ public: | |||
| 59 | 54 | ||
| 60 | Status Init(const std::string &hostId, const HostPort &initialWorkerAddr, bool initialWorkerIsLocal = false); | 55 | Status Init(const std::string &hostId, const HostPort &initialWorkerAddr, bool initialWorkerIsLocal = false); |
| 61 | 56 | ||
| 62 | - // Legacy compatibility overload. New callers must use DataPlacementPolicy. | ||
| 63 | - Status SelectWorker(const std::string &key, SelectStrategy strategy, HostPort &worker, | ||
| 64 | - const std::vector<HostPort> &exclude = {}); | ||
| 65 | - | ||
| 66 | // Routing owns only its GetHashRing control channel. The caller owns business RPC execution and retries. | 57 | // Routing owns only its GetHashRing control channel. The caller owns business RPC execution and retries. |
| 67 | Status SelectWorker(const std::string &key, DataPlacementPolicy policy, HostPort &worker, | 58 | Status SelectWorker(const std::string &key, DataPlacementPolicy policy, HostPort &worker, |
| 68 | const std::vector<HostPort> &exclude = {}); | 59 | const std::vector<HostPort> &exclude = {}); |
| 69 | 60 | ||
| 70 | - // Legacy compatibility overload. New callers must use DataPlacementPolicy. | ||
| 71 | - Status SelectWorkers(const std::vector<std::string> &keys, SelectStrategy strategy, | ||
| 72 | - std::unordered_map<HostPort, std::vector<std::string>> &groups); | ||
| 73 | - | ||
| 74 | - // Legacy compatibility overload with per-request exclude list (Exist reroute). | ||
| 75 | - Status SelectWorkers(const std::vector<std::string> &keys, SelectStrategy strategy, | ||
| 76 | - std::unordered_map<HostPort, std::vector<std::string>> &groups, | ||
| 77 | - const std::vector<HostPort> &exclude); | ||
| 78 | - | ||
| 79 | Status SelectWorkers(const std::vector<std::string> &keys, DataPlacementPolicy policy, | 61 | Status SelectWorkers(const std::vector<std::string> &keys, DataPlacementPolicy policy, |
| 80 | std::unordered_map<HostPort, std::vector<std::string>> &groups, | 62 | std::unordered_map<HostPort, std::vector<std::string>> &groups, |
| 81 | const std::vector<HostPort> &exclude = {}); | 63 | const std::vector<HostPort> &exclude = {}); |
| @@ -92,14 +92,6 @@ bool WorkerRouter::IsExcluded(const HostPort &addr, const std::vector<HostPort> | |||
| 92 | [&](const HostPort &e) { return e == addr; }); | 92 | [&](const HostPort &e) { return e == addr; }); |
| 93 | } | 93 | } |
| 94 | 94 | ||
| 95 | -Status WorkerRouter::SelectWorker(const std::string &key, SelectStrategy strategy, | ||
| 96 | - HostPort &worker, const std::vector<HostPort> &exclude) const | ||
| 97 | -{ | ||
| 98 | - auto policy = strategy == SelectStrategy::SAME_NODE_PREFERRED ? DataPlacementPolicy::PREFERRED_SAME_NODE | ||
| 99 | - : DataPlacementPolicy::PREFERRED_META_OWNER; | ||
| 100 | - return SelectWorker(key, policy, worker, exclude); | ||
| 101 | -} | ||
| 102 | - | ||
| 103 | Status WorkerRouter::SelectWorker(const std::string &key, DataPlacementPolicy policy, HostPort &worker, | 95 | Status WorkerRouter::SelectWorker(const std::string &key, DataPlacementPolicy policy, HostPort &worker, |
| 104 | const std::vector<HostPort> &exclude) const | 96 | const std::vector<HostPort> &exclude) const |
| 105 | { | 97 | { |
| @@ -160,14 +152,6 @@ Status WorkerRouter::SelectWorkerFromView(const std::string &key, DataPlacementP | |||
| 160 | return Status(K_NO_AVAILABLE_WORKER, "All workers filtered or excluded"); | 152 | return Status(K_NO_AVAILABLE_WORKER, "All workers filtered or excluded"); |
| 161 | } | 153 | } |
| 162 | 154 | ||
| 163 | -Status WorkerRouter::SelectWorkers(const std::vector<std::string> &keys, SelectStrategy strategy, | ||
| 164 | - std::unordered_map<HostPort, std::vector<std::string>> &groups) const | ||
| 165 | -{ | ||
| 166 | - auto policy = strategy == SelectStrategy::SAME_NODE_PREFERRED ? DataPlacementPolicy::PREFERRED_SAME_NODE | ||
| 167 | - : DataPlacementPolicy::PREFERRED_META_OWNER; | ||
| 168 | - return SelectWorkers(keys, policy, groups); | ||
| 169 | -} | ||
| 170 | - | ||
| 171 | Status WorkerRouter::SelectWorkers(const std::vector<std::string> &keys, DataPlacementPolicy policy, | 155 | Status WorkerRouter::SelectWorkers(const std::vector<std::string> &keys, DataPlacementPolicy policy, |
| 172 | std::unordered_map<HostPort, std::vector<std::string>> &groups, | 156 | std::unordered_map<HostPort, std::vector<std::string>> &groups, |
| 173 | const std::vector<HostPort> &exclude) const | 157 | const std::vector<HostPort> &exclude) const |
| @@ -17,7 +17,7 @@ | |||
| 17 | /** | 17 | /** |
| 18 | * Description: WorkerRouter - core worker selection component. | 18 | * Description: WorkerRouter - core worker selection component. |
| 19 | * Holds hash ring data (written by Refresher), traverses filter chain, | 19 | * Holds hash ring data (written by Refresher), traverses filter chain, |
| 20 | - * selects worker by strategy. Does NOT do RPC calls, retries, or error handling. | 20 | + * selects worker by data placement policy. Does NOT do RPC calls, retries, or error handling. |
| 21 | */ | 21 | */ |
| 22 | 22 | ||
| 23 | 23 | ||
| @@ -29,8 +29,8 @@ | |||
| 29 | 29 | ||
| 30 | 30 | ||
| 31 | 31 | ||
| 32 | + | ||
| 32 | 33 | ||
| 33 | - | ||
| 34 | 34 | ||
| 35 | 35 | ||
| 36 | 36 | ||
| @@ -39,8 +39,6 @@ | |||
| 39 | namespace datasystem { | 39 | namespace datasystem { |
| 40 | namespace client { | 40 | namespace client { |
| 41 | 41 | ||
| 42 | -enum class DataPlacementPolicy : uint8_t; | ||
| 43 | - | ||
| 44 | enum class WorkerRingState { | 42 | enum class WorkerRingState { |
| 45 | UNKNOWN, | 43 | UNKNOWN, |
| 46 | INITIAL, | 44 | INITIAL, |
| @@ -60,18 +58,10 @@ public: | |||
| 60 | void SetHostId(std::string hostId); | 58 | void SetHostId(std::string hostId); |
| 61 | 59 | ||
| 62 | // Single key selection. exclude list avoids specific workers (e.g., LEAVING on write retry). | 60 | // Single key selection. exclude list avoids specific workers (e.g., LEAVING on write retry). |
| 63 | - // Legacy compatibility overload. New callers must use DataPlacementPolicy. | ||
| 64 | - Status SelectWorker(const std::string &key, SelectStrategy strategy, HostPort &worker, | ||
| 65 | - const std::vector<HostPort> &exclude = {}) const; | ||
| 66 | - | ||
| 67 | Status SelectWorker(const std::string &key, DataPlacementPolicy policy, HostPort &worker, | 61 | Status SelectWorker(const std::string &key, DataPlacementPolicy policy, HostPort &worker, |
| 68 | const std::vector<HostPort> &exclude = {}) const; | 62 | const std::vector<HostPort> &exclude = {}) const; |
| 69 | 63 | ||
| 70 | // Batch selection: group keys by owner, return map<worker, keys>. | 64 | // Batch selection: group keys by owner, return map<worker, keys>. |
| 71 | - // Legacy compatibility overload. New callers must use DataPlacementPolicy. | ||
| 72 | - Status SelectWorkers(const std::vector<std::string> &keys, SelectStrategy strategy, | ||
| 73 | - std::unordered_map<HostPort, std::vector<std::string>> &groups) const; | ||
| 74 | - | ||
| 75 | Status SelectWorkers(const std::vector<std::string> &keys, DataPlacementPolicy policy, | 65 | Status SelectWorkers(const std::vector<std::string> &keys, DataPlacementPolicy policy, |
| 76 | std::unordered_map<HostPort, std::vector<std::string>> &groups, | 66 | std::unordered_map<HostPort, std::vector<std::string>> &groups, |
| 77 | const std::vector<HostPort> &exclude = {}) const; | 67 | const std::vector<HostPort> &exclude = {}) const; |
| @@ -111,7 +111,7 @@ protected: | |||
| 111 | for (size_t i = 0; i < KEY_SEARCH_LIMIT; ++i) { | 111 | for (size_t i = 0; i < KEY_SEARCH_LIMIT; ++i) { |
| 112 | std::string key = ROUTE_KEY_PREFIX + std::to_string(i); | 112 | std::string key = ROUTE_KEY_PREFIX + std::to_string(i); |
| 113 | HostPort selectedWorker; | 113 | HostPort selectedWorker; |
| 114 | - Status rc = router.SelectWorker(key, client::SelectStrategy::HASH_RING_AFFINITY, selectedWorker); | 114 | + Status rc = router.SelectWorker(key, client::DataPlacementPolicy::PREFERRED_META_OWNER, selectedWorker); |
| 115 | if (rc.IsOk() && selectedWorker == targetWorker) { | 115 | if (rc.IsOk() && selectedWorker == targetWorker) { |
| 116 | return key; | 116 | return key; |
| 117 | } | 117 | } |
| @@ -126,13 +126,13 @@ protected: | |||
| 126 | const HostPort &targetWorker, const std::string &key) | 126 | const HostPort &targetWorker, const std::string &key) |
| 127 | { | 127 | { |
| 128 | HostPort selectedWorker; | 128 | HostPort selectedWorker; |
| 129 | - DS_ASSERT_OK(router.SelectWorker(key, client::SelectStrategy::HASH_RING_AFFINITY, selectedWorker)); | 129 | + DS_ASSERT_OK(router.SelectWorker(key, client::DataPlacementPolicy::PREFERRED_META_OWNER, selectedWorker)); |
| 130 | EXPECT_EQ(selectedWorker, targetWorker); | 130 | EXPECT_EQ(selectedWorker, targetWorker); |
| 131 | 131 | ||
| 132 | auto leavingTopology = std::make_shared<ClusterTopologyPb>(topology); | 132 | auto leavingTopology = std::make_shared<ClusterTopologyPb>(topology); |
| 133 | (*leavingTopology->mutable_members())[targetWorker.ToString()].set_state(MembershipPb::LEAVING); | 133 | (*leavingTopology->mutable_members())[targetWorker.ToString()].set_state(MembershipPb::LEAVING); |
| 134 | router.UpdateHashRing(leavingTopology, hostIdMap); | 134 | router.UpdateHashRing(leavingTopology, hostIdMap); |
| 135 | - DS_ASSERT_OK(router.SelectWorker(key, client::SelectStrategy::HASH_RING_AFFINITY, selectedWorker)); | 135 | + DS_ASSERT_OK(router.SelectWorker(key, client::DataPlacementPolicy::PREFERRED_META_OWNER, selectedWorker)); |
| 136 | EXPECT_EQ(selectedWorker, remainingWorker); | 136 | EXPECT_EQ(selectedWorker, remainingWorker); |
| 137 | } | 137 | } |
| 138 | 138 | ||
| @@ -145,13 +145,13 @@ protected: | |||
| 145 | router.UpdateState(targetWorker, K_CLIENT_WORKER_DISCONNECT); | 145 | router.UpdateState(targetWorker, K_CLIENT_WORKER_DISCONNECT); |
| 146 | } | 146 | } |
| 147 | HostPort selectedWorker; | 147 | HostPort selectedWorker; |
| 148 | - DS_ASSERT_OK(router.SelectWorker(key, client::SelectStrategy::HASH_RING_AFFINITY, selectedWorker)); | 148 | + DS_ASSERT_OK(router.SelectWorker(key, client::DataPlacementPolicy::PREFERRED_META_OWNER, selectedWorker)); |
| 149 | ASSERT_EQ(selectedWorker, remainingWorker); | 149 | ASSERT_EQ(selectedWorker, remainingWorker); |
| 150 | 150 | ||
| 151 | const auto deadline = | 151 | const auto deadline = |
| 152 | std::chrono::steady_clock::now() + std::chrono::milliseconds(BROKEN_FILTER_RECOVERY_TIMEOUT_MS); | 152 | std::chrono::steady_clock::now() + std::chrono::milliseconds(BROKEN_FILTER_RECOVERY_TIMEOUT_MS); |
| 153 | while (std::chrono::steady_clock::now() < deadline) { | 153 | while (std::chrono::steady_clock::now() < deadline) { |
| 154 | - DS_ASSERT_OK(router.SelectWorker(key, client::SelectStrategy::HASH_RING_AFFINITY, selectedWorker)); | 154 | + DS_ASSERT_OK(router.SelectWorker(key, client::DataPlacementPolicy::PREFERRED_META_OWNER, selectedWorker)); |
| 155 | if (selectedWorker == targetWorker) { | 155 | if (selectedWorker == targetWorker) { |
| 156 | return; | 156 | return; |
| 157 | } | 157 | } |
| @@ -40,12 +40,12 @@ HostPort MakeWorker(int port) | |||
| 40 | 40 | ||
| 41 | class FakeExistRouting : public IExistRouting { | 41 | class FakeExistRouting : public IExistRouting { |
| 42 | public: | 42 | public: |
| 43 | - Status SelectWorkers(const std::vector<std::string> &, client::SelectStrategy strategy, | 43 | + Status SelectWorkers(const std::vector<std::string> &, client::DataPlacementPolicy policy, |
| 44 | std::unordered_map<HostPort, std::vector<std::string>> &output, | 44 | std::unordered_map<HostPort, std::vector<std::string>> &output, |
| 45 | const std::vector<HostPort> &exclude) override | 45 | const std::vector<HostPort> &exclude) override |
| 46 | { | 46 | { |
| 47 | ++selectWorkersCount; | 47 | ++selectWorkersCount; |
| 48 | - selectedStrategy = strategy; | 48 | + selectedPolicy = policy; |
| 49 | excludeHistory.emplace_back(exclude); | 49 | excludeHistory.emplace_back(exclude); |
| 50 | if (!groupSequence.empty()) { | 50 | if (!groupSequence.empty()) { |
| 51 | output = groupSequence.front(); | 51 | output = groupSequence.front(); |
| @@ -63,7 +63,7 @@ public: | |||
| 63 | } | 63 | } |
| 64 | 64 | ||
| 65 | Status selectStatus = Status::OK(); | 65 | Status selectStatus = Status::OK(); |
| 66 | - client::SelectStrategy selectedStrategy = client::SelectStrategy::SAME_NODE_PREFERRED; | 66 | + client::DataPlacementPolicy selectedPolicy = client::DataPlacementPolicy::PREFERRED_SAME_NODE; |
| 67 | std::unordered_map<HostPort, std::vector<std::string>> groups; | 67 | std::unordered_map<HostPort, std::vector<std::string>> groups; |
| 68 | std::vector<std::unordered_map<HostPort, std::vector<std::string>>> groupSequence; | 68 | std::vector<std::unordered_map<HostPort, std::vector<std::string>>> groupSequence; |
| 69 | std::vector<HostPort> updatedWorkers; | 69 | std::vector<HostPort> updatedWorkers; |
| @@ -169,7 +169,7 @@ TEST_F(ExistHandlerTest, ExistUsesRoutingAndTransportAndKeepsInputOrder) | |||
| 169 | ASSERT_TRUE(rc.IsOk()); | 169 | ASSERT_TRUE(rc.IsOk()); |
| 170 | EXPECT_EQ(exists, std::vector<bool>({ true, false })); | 170 | EXPECT_EQ(exists, std::vector<bool>({ true, false })); |
| 171 | EXPECT_EQ(routing_->selectWorkersCount, 1); | 171 | EXPECT_EQ(routing_->selectWorkersCount, 1); |
| 172 | - EXPECT_EQ(routing_->selectedStrategy, client::SelectStrategy::HASH_RING_AFFINITY); | 172 | + EXPECT_EQ(routing_->selectedPolicy, client::DataPlacementPolicy::PREFERRED_META_OWNER); |
| 173 | EXPECT_TRUE(transport_->queryL2Cache[0]); | 173 | EXPECT_TRUE(transport_->queryL2Cache[0]); |
| 174 | EXPECT_FALSE(transport_->isLocal[0]); | 174 | EXPECT_FALSE(transport_->isLocal[0]); |
| 175 | } | 175 | } |
| @@ -91,7 +91,7 @@ TEST_F(HashRingRefresherTest, TestInitialFetchRunsBeforePeriodicThread) | |||
| 91 | EXPECT_EQ(fetchCount, 1); | 91 | EXPECT_EQ(fetchCount, 1); |
| 92 | 92 | ||
| 93 | HostPort selected; | 93 | HostPort selected; |
| 94 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 94 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 95 | EXPECT_EQ(selected.ToString(), "127.0.0.1:1000"); | 95 | EXPECT_EQ(selected.ToString(), "127.0.0.1:1000"); |
| 96 | } | 96 | } |
| 97 | 97 | ||
| @@ -404,7 +404,7 @@ TEST_F(HashRingRefresherTest, TestRingUpdateHookRunsBeforeRoutePublication) | |||
| 404 | EXPECT_EQ(ring.members_size(), 1); | 404 | EXPECT_EQ(ring.members_size(), 1); |
| 405 | HostPort selected; | 405 | HostPort selected; |
| 406 | routeWasUnpublished = | 406 | routeWasUnpublished = |
| 407 | - router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected).IsError(); | 407 | + router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected).IsError(); |
| 408 | return Status::OK(); | 408 | return Status::OK(); |
| 409 | }; | 409 | }; |
| 410 | client::HashRingRefresher refresher(router, fetch, hook); | 410 | client::HashRingRefresher refresher(router, fetch, hook); |
| @@ -413,7 +413,7 @@ TEST_F(HashRingRefresherTest, TestRingUpdateHookRunsBeforeRoutePublication) | |||
| 413 | EXPECT_EQ(hookVersion, 5u); | 413 | EXPECT_EQ(hookVersion, 5u); |
| 414 | EXPECT_TRUE(routeWasUnpublished); | 414 | EXPECT_TRUE(routeWasUnpublished); |
| 415 | HostPort selected; | 415 | HostPort selected; |
| 416 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 416 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 417 | EXPECT_EQ(selected.ToString(), "127.0.0.1:1000"); | 417 | EXPECT_EQ(selected.ToString(), "127.0.0.1:1000"); |
| 418 | } | 418 | } |
| 419 | 419 | ||
| @@ -448,7 +448,7 @@ TEST_F(HashRingRefresherTest, TestFailedRingUpdateHookRetainsVersionAndRetries) | |||
| 448 | 448 | ||
| 449 | EXPECT_EQ(refresher.InitialFetch(HostPort("127.0.0.1", 1000)).GetCode(), K_RUNTIME_ERROR); | 449 | EXPECT_EQ(refresher.InitialFetch(HostPort("127.0.0.1", 1000)).GetCode(), K_RUNTIME_ERROR); |
| 450 | HostPort selected; | 450 | HostPort selected; |
| 451 | - EXPECT_TRUE(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected).IsError()); | 451 | + EXPECT_TRUE(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected).IsError()); |
| 452 | DS_ASSERT_OK(refresher.StartPeriodicRefresh(60'000)); | 452 | DS_ASSERT_OK(refresher.StartPeriodicRefresh(60'000)); |
| 453 | bool retried = false; | 453 | bool retried = false; |
| 454 | { | 454 | { |
| @@ -461,7 +461,7 @@ TEST_F(HashRingRefresherTest, TestFailedRingUpdateHookRetainsVersionAndRetries) | |||
| 461 | ASSERT_GE(requestedVersions.size(), 2u); | 461 | ASSERT_GE(requestedVersions.size(), 2u); |
| 462 | EXPECT_EQ(requestedVersions[0], 0u); | 462 | EXPECT_EQ(requestedVersions[0], 0u); |
| 463 | EXPECT_EQ(requestedVersions[1], 0u); | 463 | EXPECT_EQ(requestedVersions[1], 0u); |
| 464 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 464 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 465 | } | 465 | } |
| 466 | 466 | ||
| 467 | TEST_F(HashRingRefresherTest, TestStaleVersionDoesNotReplaceCurrentRing) | 467 | TEST_F(HashRingRefresherTest, TestStaleVersionDoesNotReplaceCurrentRing) |
| @@ -504,7 +504,7 @@ TEST_F(HashRingRefresherTest, TestStaleVersionDoesNotReplaceCurrentRing) | |||
| 504 | EXPECT_EQ(requestedVersions[0], 0u); | 504 | EXPECT_EQ(requestedVersions[0], 0u); |
| 505 | EXPECT_EQ(requestedVersions[1], 2u); | 505 | EXPECT_EQ(requestedVersions[1], 2u); |
| 506 | HostPort selected; | 506 | HostPort selected; |
| 507 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 507 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 508 | EXPECT_EQ(selected.ToString(), "127.0.0.1:1000"); | 508 | EXPECT_EQ(selected.ToString(), "127.0.0.1:1000"); |
| 509 | } | 509 | } |
| 510 | 510 | ||
| @@ -537,7 +537,7 @@ TEST_F(HashRingRefresherTest, TestUnchangedResponseKeepsCurrentRing) | |||
| 537 | 537 | ||
| 538 | EXPECT_EQ(requestedVersions[1], 5u); | 538 | EXPECT_EQ(requestedVersions[1], 5u); |
| 539 | HostPort selected; | 539 | HostPort selected; |
| 540 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 540 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 541 | EXPECT_EQ(selected.ToString(), "127.0.0.1:1000"); | 541 | EXPECT_EQ(selected.ToString(), "127.0.0.1:1000"); |
| 542 | } | 542 | } |
| 543 | 543 | ||
| @@ -571,7 +571,7 @@ TEST_F(HashRingRefresherTest, TestForcedRefreshRetriesUntilRingChanges) | |||
| 571 | refresher.Stop(); | 571 | refresher.Stop(); |
| 572 | 572 | ||
| 573 | HostPort selected; | 573 | HostPort selected; |
| 574 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 574 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 575 | EXPECT_EQ(selected.ToString(), "127.0.0.1:2000"); | 575 | EXPECT_EQ(selected.ToString(), "127.0.0.1:2000"); |
| 576 | } | 576 | } |
| 577 | 577 | ||
| @@ -615,7 +615,7 @@ TEST_F(HashRingRefresherTest, RepeatedFailureExtendsForcedRefreshUntilIsolationP | |||
| 615 | ASSERT_TRUE(isolationPublished); | 615 | ASSERT_TRUE(isolationPublished); |
| 616 | 616 | ||
| 617 | HostPort selected; | 617 | HostPort selected; |
| 618 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 618 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 619 | EXPECT_EQ(selected.ToString(), "127.0.0.1:2000"); | 619 | EXPECT_EQ(selected.ToString(), "127.0.0.1:2000"); |
| 620 | } | 620 | } |
| 621 | 621 | ||
| @@ -644,7 +644,7 @@ TEST_F(HashRingRefresherTest, TestAllWorkersUnreachableKeepsCurrentRing) | |||
| 644 | 644 | ||
| 645 | DS_ASSERT_OK(refresher.InitialFetch(HostPort("127.0.0.1", 1000))); | 645 | DS_ASSERT_OK(refresher.InitialFetch(HostPort("127.0.0.1", 1000))); |
| 646 | HostPort before; | 646 | HostPort before; |
| 647 | - DS_ASSERT_OK(router->SelectWorker("stable-key", client::SelectStrategy::HASH_RING_AFFINITY, before)); | 647 | + DS_ASSERT_OK(router->SelectWorker("stable-key", client::DataPlacementPolicy::PREFERRED_META_OWNER, before)); |
| 648 | DS_ASSERT_OK(refresher.StartPeriodicRefresh(60'000)); | 648 | DS_ASSERT_OK(refresher.StartPeriodicRefresh(60'000)); |
| 649 | { | 649 | { |
| 650 | std::unique_lock<std::mutex> lock(mutex); | 650 | std::unique_lock<std::mutex> lock(mutex); |
| @@ -653,7 +653,7 @@ TEST_F(HashRingRefresherTest, TestAllWorkersUnreachableKeepsCurrentRing) | |||
| 653 | refresher.Stop(); | 653 | refresher.Stop(); |
| 654 | 654 | ||
| 655 | HostPort after; | 655 | HostPort after; |
| 656 | - DS_ASSERT_OK(router->SelectWorker("stable-key", client::SelectStrategy::HASH_RING_AFFINITY, after)); | 656 | + DS_ASSERT_OK(router->SelectWorker("stable-key", client::DataPlacementPolicy::PREFERRED_META_OWNER, after)); |
| 657 | EXPECT_EQ(after, before); | 657 | EXPECT_EQ(after, before); |
| 658 | } | 658 | } |
| 659 | 659 | ||
| @@ -698,7 +698,7 @@ TEST_F(HashRingRefresherTest, TestFilterNotifiedOnlyWhenRingChanges) | |||
| 698 | 698 | ||
| 699 | EXPECT_EQ(filter->UpdateCount(), 2); | 699 | EXPECT_EQ(filter->UpdateCount(), 2); |
| 700 | HostPort selected; | 700 | HostPort selected; |
| 701 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 701 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 702 | EXPECT_EQ(selected.ToString(), "127.0.0.1:2000"); | 702 | EXPECT_EQ(selected.ToString(), "127.0.0.1:2000"); |
| 703 | } | 703 | } |
| 704 | 704 | ||
| @@ -742,7 +742,7 @@ TEST_F(HashRingRefresherTest, BackgroundRefreshContinuesPastReachableUnchangedWo | |||
| 742 | refresher.Stop(); | 742 | refresher.Stop(); |
| 743 | 743 | ||
| 744 | HostPort selected; | 744 | HostPort selected; |
| 745 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 745 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 746 | EXPECT_EQ(selected, HostPort("127.0.0.1", 2000)); | 746 | EXPECT_EQ(selected, HostPort("127.0.0.1", 2000)); |
| 747 | } | 747 | } |
| 748 | 748 | ||
| @@ -87,7 +87,7 @@ TEST_F(RoutingTest, TestSelectWorkerEmptyRing) | |||
| 87 | { | 87 | { |
| 88 | auto router = CreateRouter(); | 88 | auto router = CreateRouter(); |
| 89 | HostPort worker; | 89 | HostPort worker; |
| 90 | - auto st = router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, worker); | 90 | + auto st = router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, worker); |
| 91 | EXPECT_FALSE(st.IsOk()); | 91 | EXPECT_FALSE(st.IsOk()); |
| 92 | } | 92 | } |
| 93 | 93 | ||
| @@ -97,7 +97,7 @@ TEST_F(RoutingTest, TestSelectWorkerReturnsActive) | |||
| 97 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); | 97 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); |
| 98 | 98 | ||
| 99 | HostPort worker; | 99 | HostPort worker; |
| 100 | - DS_ASSERT_OK(router->SelectWorker("key", client::SelectStrategy::HASH_RING_AFFINITY, worker)); | 100 | + DS_ASSERT_OK(router->SelectWorker("key", client::DataPlacementPolicy::PREFERRED_META_OWNER, worker)); |
| 101 | std::string addr = worker.ToString(); | 101 | std::string addr = worker.ToString(); |
| 102 | EXPECT_TRUE(addr == "127.0.0.1:1000" || addr == "127.0.0.1:2000"); | 102 | EXPECT_TRUE(addr == "127.0.0.1:1000" || addr == "127.0.0.1:2000"); |
| 103 | } | 103 | } |
| @@ -108,8 +108,8 @@ TEST_F(RoutingTest, TestSelectWorkerConsistency) | |||
| 108 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); | 108 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); |
| 109 | 109 | ||
| 110 | HostPort w1, w2; | 110 | HostPort w1, w2; |
| 111 | - DS_ASSERT_OK(router->SelectWorker("consistency_key", client::SelectStrategy::HASH_RING_AFFINITY, w1)); | 111 | + DS_ASSERT_OK(router->SelectWorker("consistency_key", client::DataPlacementPolicy::PREFERRED_META_OWNER, w1)); |
| 112 | - DS_ASSERT_OK(router->SelectWorker("consistency_key", client::SelectStrategy::HASH_RING_AFFINITY, w2)); | 112 | + DS_ASSERT_OK(router->SelectWorker("consistency_key", client::DataPlacementPolicy::PREFERRED_META_OWNER, w2)); |
| 113 | EXPECT_EQ(w1.ToString(), w2.ToString()); | 113 | EXPECT_EQ(w1.ToString(), w2.ToString()); |
| 114 | } | 114 | } |
| 115 | 115 | ||
| @@ -119,11 +119,12 @@ TEST_F(RoutingTest, TestSelectWorkerExclude) | |||
| 119 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); | 119 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); |
| 120 | 120 | ||
| 121 | HostPort first; | 121 | HostPort first; |
| 122 | - DS_ASSERT_OK(router->SelectWorker("exclude_key", client::SelectStrategy::HASH_RING_AFFINITY, first)); | 122 | + DS_ASSERT_OK(router->SelectWorker("exclude_key", client::DataPlacementPolicy::PREFERRED_META_OWNER, first)); |
| 123 | 123 | ||
| 124 | // Exclude the first result, should get a different one | 124 | // Exclude the first result, should get a different one |
| 125 | HostPort second; | 125 | HostPort second; |
| 126 | - DS_ASSERT_OK(router->SelectWorker("exclude_key", client::SelectStrategy::HASH_RING_AFFINITY, second, {first})); | 126 | + DS_ASSERT_OK(router->SelectWorker("exclude_key", client::DataPlacementPolicy::PREFERRED_META_OWNER, second, |
| 127 | + { first })); | ||
| 127 | EXPECT_NE(first.ToString(), second.ToString()); | 128 | EXPECT_NE(first.ToString(), second.ToString()); |
| 128 | } | 129 | } |
| 129 | 130 | ||
| @@ -134,7 +135,7 @@ TEST_F(RoutingTest, TestSelectWorkersBatch) | |||
| 134 | 135 | ||
| 135 | std::vector<std::string> keys = { "k1", "k2", "k3", "k4", "k5" }; | 136 | std::vector<std::string> keys = { "k1", "k2", "k3", "k4", "k5" }; |
| 136 | std::unordered_map<HostPort, std::vector<std::string>> groups; | 137 | std::unordered_map<HostPort, std::vector<std::string>> groups; |
| 137 | - DS_ASSERT_OK(router->SelectWorkers(keys, client::SelectStrategy::HASH_RING_AFFINITY, groups)); | 138 | + DS_ASSERT_OK(router->SelectWorkers(keys, client::DataPlacementPolicy::PREFERRED_META_OWNER, groups)); |
| 138 | 139 | ||
| 139 | // All keys should be grouped | 140 | // All keys should be grouped |
| 140 | size_t totalKeys = 0; | 141 | size_t totalKeys = 0; |
| @@ -152,7 +153,7 @@ TEST_F(RoutingTest, TestSelectWorkersFailureDoesNotMutateOutput) | |||
| 152 | 153 | ||
| 153 | HostPort existing("127.0.0.1", 3000); | 154 | HostPort existing("127.0.0.1", 3000); |
| 154 | std::unordered_map<HostPort, std::vector<std::string>> groups{ { existing, { "existing" } } }; | 155 | std::unordered_map<HostPort, std::vector<std::string>> groups{ { existing, { "existing" } } }; |
| 155 | - auto rc = router->SelectWorkers({ "k1", "k2" }, client::SelectStrategy::HASH_RING_AFFINITY, groups); | 156 | + auto rc = router->SelectWorkers({ "k1", "k2" }, client::DataPlacementPolicy::PREFERRED_META_OWNER, groups); |
| 156 | 157 | ||
| 157 | EXPECT_TRUE(rc.IsError()); | 158 | EXPECT_TRUE(rc.IsError()); |
| 158 | ASSERT_EQ(groups.size(), 1u); | 159 | ASSERT_EQ(groups.size(), 1u); |
| @@ -165,7 +166,7 @@ TEST_F(RoutingTest, TestSelectWorkersEmptyInputClearsOutput) | |||
| 165 | HostPort existing("127.0.0.1", 3000); | 166 | HostPort existing("127.0.0.1", 3000); |
| 166 | std::unordered_map<HostPort, std::vector<std::string>> groups{ { existing, { "existing" } } }; | 167 | std::unordered_map<HostPort, std::vector<std::string>> groups{ { existing, { "existing" } } }; |
| 167 | 168 | ||
| 168 | - DS_ASSERT_OK(router->SelectWorkers({}, client::SelectStrategy::HASH_RING_AFFINITY, groups)); | 169 | + DS_ASSERT_OK(router->SelectWorkers({}, client::DataPlacementPolicy::PREFERRED_META_OWNER, groups)); |
| 169 | EXPECT_TRUE(groups.empty()); | 170 | EXPECT_TRUE(groups.empty()); |
| 170 | } | 171 | } |
| 171 | 172 | ||
| @@ -190,7 +191,7 @@ TEST_F(RoutingTest, TestSameNodePreferred) | |||
| 190 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); | 191 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); |
| 191 | 192 | ||
| 192 | HostPort worker; | 193 | HostPort worker; |
| 193 | - DS_ASSERT_OK(router->SelectWorker("samenode_key", client::SelectStrategy::SAME_NODE_PREFERRED, worker)); | 194 | + DS_ASSERT_OK(router->SelectWorker("samenode_key", client::DataPlacementPolicy::PREFERRED_SAME_NODE, worker)); |
| 194 | EXPECT_EQ(worker.ToString(), "127.0.0.1:1000"); // host-a's worker | 195 | EXPECT_EQ(worker.ToString(), "127.0.0.1:1000"); // host-a's worker |
| 195 | } | 196 | } |
| 196 | 197 | ||
| @@ -266,7 +267,7 @@ TEST_F(RoutingTest, TestSameNodePreferredDistributesByKey) | |||
| 266 | for (int i = 0; i < 64; ++i) { | 267 | for (int i = 0; i < 64; ++i) { |
| 267 | HostPort worker; | 268 | HostPort worker; |
| 268 | DS_ASSERT_OK(router->SelectWorker("same-node-" + std::to_string(i), | 269 | DS_ASSERT_OK(router->SelectWorker("same-node-" + std::to_string(i), |
| 269 | - client::SelectStrategy::SAME_NODE_PREFERRED, worker)); | 270 | + client::DataPlacementPolicy::PREFERRED_SAME_NODE, worker)); |
| 270 | selected.emplace(worker.ToString()); | 271 | selected.emplace(worker.ToString()); |
| 271 | } | 272 | } |
| 272 | EXPECT_EQ(selected.size(), 2u); | 273 | EXPECT_EQ(selected.size(), 2u); |
| @@ -279,20 +280,20 @@ TEST_F(RoutingTest, TestEmptyHostIdDoesNotTreatMissingWorkerHostIdAsSameNode) | |||
| 279 | auto router = CreateRouter(""); | 280 | auto router = CreateRouter(""); |
| 280 | router->UpdateHashRing(BuildRing(), hostIdMap); | 281 | router->UpdateHashRing(BuildRing(), hostIdMap); |
| 281 | 282 | ||
| 282 | - // When client hostId is empty, SAME_NODE_PREFERRED should not treat | 283 | + // When client hostId is empty, PREFERRED_SAME_NODE should not treat |
| 283 | // workers with empty hostId as same-node. It should behave identically | 284 | // workers with empty hostId as same-node. It should behave identically |
| 284 | - // to HASH_RING_AFFINITY (no same-node bias). | 285 | + // to PREFERRED_META_OWNER (no same-node bias). |
| 285 | // Verify with multiple keys: every key should return the same worker | 286 | // Verify with multiple keys: every key should return the same worker |
| 286 | // regardless of strategy. | 287 | // regardless of strategy. |
| 287 | for (int i = 0; i < 100; ++i) { | 288 | for (int i = 0; i < 100; ++i) { |
| 288 | std::string key = "empty-hostid-key-" + std::to_string(i); | 289 | std::string key = "empty-hostid-key-" + std::to_string(i); |
| 289 | HostPort hashOwner; | 290 | HostPort hashOwner; |
| 290 | - DS_ASSERT_OK(router->SelectWorker(key, client::SelectStrategy::HASH_RING_AFFINITY, hashOwner)); | 291 | + DS_ASSERT_OK(router->SelectWorker(key, client::DataPlacementPolicy::PREFERRED_META_OWNER, hashOwner)); |
| 291 | 292 | ||
| 292 | HostPort selected; | 293 | HostPort selected; |
| 293 | - DS_ASSERT_OK(router->SelectWorker(key, client::SelectStrategy::SAME_NODE_PREFERRED, selected)); | 294 | + DS_ASSERT_OK(router->SelectWorker(key, client::DataPlacementPolicy::PREFERRED_SAME_NODE, selected)); |
| 294 | EXPECT_EQ(selected, hashOwner) | 295 | EXPECT_EQ(selected, hashOwner) |
| 295 | - << "Key " << key << ": SAME_NODE_PREFERRED diverged from HASH_RING_AFFINITY"; | 296 | + << "Key " << key << ": PREFERRED_SAME_NODE diverged from PREFERRED_META_OWNER"; |
| 296 | } | 297 | } |
| 297 | } | 298 | } |
| 298 | 299 | ||
| @@ -304,14 +305,14 @@ TEST_F(RoutingTest, TestStateFilterRejectsLeavingWorker) | |||
| 304 | 305 | ||
| 305 | const std::string key = "leaving-owner"; | 306 | const std::string key = "leaving-owner"; |
| 306 | HostPort original; | 307 | HostPort original; |
| 307 | - DS_ASSERT_OK(router->SelectWorker(key, client::SelectStrategy::HASH_RING_AFFINITY, original)); | 308 | + DS_ASSERT_OK(router->SelectWorker(key, client::DataPlacementPolicy::PREFERRED_META_OWNER, original)); |
| 308 | 309 | ||
| 309 | auto updatedRing = BuildRing(); | 310 | auto updatedRing = BuildRing(); |
| 310 | (*updatedRing->mutable_members())[original.ToString()].set_state(::datasystem::MembershipPb::LEAVING); | 311 | (*updatedRing->mutable_members())[original.ToString()].set_state(::datasystem::MembershipPb::LEAVING); |
| 311 | router->UpdateHashRing(updatedRing, BuildHostIdMap()); | 312 | router->UpdateHashRing(updatedRing, BuildHostIdMap()); |
| 312 | 313 | ||
| 313 | HostPort selected; | 314 | HostPort selected; |
| 314 | - DS_ASSERT_OK(router->SelectWorker(key, client::SelectStrategy::HASH_RING_AFFINITY, selected)); | 315 | + DS_ASSERT_OK(router->SelectWorker(key, client::DataPlacementPolicy::PREFERRED_META_OWNER, selected)); |
| 315 | EXPECT_NE(selected, original); | 316 | EXPECT_NE(selected, original); |
| 316 | } | 317 | } |
| 317 | 318 | ||
| @@ -447,14 +448,14 @@ TEST_F(RoutingTest, TestBrokenFilterIntegrationWithRouter) | |||
| 447 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); | 448 | router->UpdateHashRing(BuildRing(), BuildHostIdMap()); |
| 448 | 449 | ||
| 449 | HostPort first; | 450 | HostPort first; |
| 450 | - DS_ASSERT_OK(router->SelectWorker("broken_key", client::SelectStrategy::HASH_RING_AFFINITY, first)); | 451 | + DS_ASSERT_OK(router->SelectWorker("broken_key", client::DataPlacementPolicy::PREFERRED_META_OWNER, first)); |
| 451 | 452 | ||
| 452 | // Mark first worker as broken (reach the debounce threshold) | 453 | // Mark first worker as broken (reach the debounce threshold) |
| 453 | MarkWorkerBroken(*router, first); | 454 | MarkWorkerBroken(*router, first); |
| 454 | 455 | ||
| 455 | // Subsequent SelectWorker should skip broken worker | 456 | // Subsequent SelectWorker should skip broken worker |
| 456 | HostPort second; | 457 | HostPort second; |
| 457 | - DS_ASSERT_OK(router->SelectWorker("broken_key", client::SelectStrategy::HASH_RING_AFFINITY, second)); | 458 | + DS_ASSERT_OK(router->SelectWorker("broken_key", client::DataPlacementPolicy::PREFERRED_META_OWNER, second)); |
| 458 | EXPECT_NE(first.ToString(), second.ToString()); | 459 | EXPECT_NE(first.ToString(), second.ToString()); |
| 459 | } | 460 | } |
| 460 | 461 | ||
| @@ -484,7 +485,7 @@ TEST_F(RoutingTest, U7RoutesWithFiveThousandWorkerSnapshot) | |||
| 484 | keys.emplace_back("u7-scale-key-" + std::to_string(i)); | 485 | keys.emplace_back("u7-scale-key-" + std::to_string(i)); |
| 485 | } | 486 | } |
| 486 | std::unordered_map<HostPort, std::vector<std::string>> groups; | 487 | std::unordered_map<HostPort, std::vector<std::string>> groups; |
| 487 | - DS_ASSERT_OK(router->SelectWorkers(keys, client::SelectStrategy::HASH_RING_AFFINITY, groups)); | 488 | + DS_ASSERT_OK(router->SelectWorkers(keys, client::DataPlacementPolicy::PREFERRED_META_OWNER, groups)); |
| 488 | 489 | ||
| 489 | size_t selectedKeyCount = 0; | 490 | size_t selectedKeyCount = 0; |
| 490 | for (const auto &group : groups) { | 491 | for (const auto &group : groups) { |
| @@ -546,7 +547,7 @@ TEST_F(RoutingTest, U7BatchSelectionNeverMixesConcurrentSnapshots) | |||
| 546 | for (size_t iteration = 0; iteration < 200; ++iteration) { | 547 | for (size_t iteration = 0; iteration < 200; ++iteration) { |
| 547 | std::unordered_map<HostPort, std::vector<std::string>> groups; | 548 | std::unordered_map<HostPort, std::vector<std::string>> groups; |
| 548 | selectionStatus = | 549 | selectionStatus = |
| 549 | - router->SelectWorkers(keys, client::SelectStrategy::HASH_RING_AFFINITY, groups); | 550 | + router->SelectWorkers(keys, client::DataPlacementPolicy::PREFERRED_META_OWNER, groups); |
| 550 | if (selectionStatus.IsError() || groups.size() != 2u) { | 551 | if (selectionStatus.IsError() || groups.size() != 2u) { |
| 551 | snapshotsConsistent = false; | 552 | snapshotsConsistent = false; |
| 552 | break; | 553 | break; |
| @@ -587,7 +588,7 @@ TEST_F(RoutingTest, AllWorkersBrokenExhaustsRingAndReturnsCode37) | |||
| 587 | 588 | ||
| 588 | // Sanity: routing succeeds before any worker is marked broken. | 589 | // Sanity: routing succeeds before any worker is marked broken. |
| 589 | HostPort healthy; | 590 | HostPort healthy; |
| 590 | - DS_ASSERT_OK(router->SelectWorker("k", client::SelectStrategy::HASH_RING_AFFINITY, healthy)); | 591 | + DS_ASSERT_OK(router->SelectWorker("k", client::DataPlacementPolicy::PREFERRED_META_OWNER, healthy)); |
| 591 | 592 | ||
| 592 | const std::vector<HostPort> allWorkers{ HostPort("127.0.0.1", 1000), HostPort("127.0.0.1", 2000) }; | 593 | const std::vector<HostPort> allWorkers{ HostPort("127.0.0.1", 1000), HostPort("127.0.0.1", 2000) }; |
| 593 | // Mark every candidate broken via the genuine-disconnect path that HandleSetRouteFailure | 594 | // Mark every candidate broken via the genuine-disconnect path that HandleSetRouteFailure |
| @@ -598,7 +599,7 @@ TEST_F(RoutingTest, AllWorkersBrokenExhaustsRingAndReturnsCode37) | |||
| 598 | 599 | ||
| 599 | // Every candidate filtered by BrokenFilter -> K_NO_AVAILABLE_WORKER (code 37). | 600 | // Every candidate filtered by BrokenFilter -> K_NO_AVAILABLE_WORKER (code 37). |
| 600 | HostPort selected; | 601 | HostPort selected; |
| 601 | - auto rc = router->SelectWorker("k", client::SelectStrategy::HASH_RING_AFFINITY, selected); | 602 | + auto rc = router->SelectWorker("k", client::DataPlacementPolicy::PREFERRED_META_OWNER, selected); |
| 602 | EXPECT_EQ(rc.GetCode(), K_NO_AVAILABLE_WORKER); | 603 | EXPECT_EQ(rc.GetCode(), K_NO_AVAILABLE_WORKER); |
| 603 | } | 604 | } |
| 604 | 605 | ||