diff --git a/Package.swift b/Package.swift index 5e15d8e55af..0465d912aa5 100644 --- a/Package.swift +++ b/Package.swift @@ -1896,7 +1896,6 @@ let package = Package( "src/core/load_balancing/rls/rls.h", "src/core/load_balancing/round_robin/round_robin.cc", "src/core/load_balancing/subchannel_interface.h", - "src/core/load_balancing/subchannel_list.h", "src/core/load_balancing/weighted_round_robin/static_stride_scheduler.cc", "src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h", "src/core/load_balancing/weighted_round_robin/weighted_round_robin.cc", diff --git a/bazel/experiments.bzl b/bazel/experiments.bzl index 7edd8cd75ea..79de231736b 100644 --- a/bazel/experiments.bzl +++ b/bazel/experiments.bzl @@ -38,7 +38,6 @@ EXPERIMENT_ENABLES = { "chaotic_good": "chaotic_good,event_engine_client,event_engine_listener,promise_based_client_call,promise_based_server_call", "registered_method_lookup_in_transport": "registered_method_lookup_in_transport", "promise_based_inproc_transport": "event_engine_client,event_engine_listener,promise_based_client_call,promise_based_inproc_transport,promise_based_server_call,registered_method_lookup_in_transport", - "round_robin_delegate_to_pick_first": "round_robin_delegate_to_pick_first", "rstpit": "rstpit", "schedule_cancellation_over_write": "schedule_cancellation_over_write", "server_privacy": "server_privacy", @@ -52,7 +51,6 @@ EXPERIMENT_ENABLES = { "v3_server_auth_filter": "v3_server_auth_filter", "work_serializer_clears_time_cache": "work_serializer_clears_time_cache", "work_serializer_dispatch": "event_engine_client,work_serializer_dispatch", - "wrr_delegate_to_pick_first": "wrr_delegate_to_pick_first", } EXPERIMENT_POLLERS = [ @@ -101,27 +99,15 @@ EXPERIMENTS = { "core_end2end_test": [ "event_engine_listener", ], - "cpp_lb_end2end_test": [ - "round_robin_delegate_to_pick_first", - "wrr_delegate_to_pick_first", - ], "credential_token_tests": [ "absl_base64", ], "event_engine_listener_test": [ "event_engine_listener", ], - "lb_unit_test": [ - "round_robin_delegate_to_pick_first", - "wrr_delegate_to_pick_first", - ], "surface_registered_method_lookup": [ "registered_method_lookup_in_transport", ], - "xds_end2end_test": [ - "round_robin_delegate_to_pick_first", - "wrr_delegate_to_pick_first", - ], }, }, "ios": { @@ -160,24 +146,12 @@ EXPERIMENTS = { ], }, "on": { - "cpp_lb_end2end_test": [ - "round_robin_delegate_to_pick_first", - "wrr_delegate_to_pick_first", - ], "credential_token_tests": [ "absl_base64", ], - "lb_unit_test": [ - "round_robin_delegate_to_pick_first", - "wrr_delegate_to_pick_first", - ], "surface_registered_method_lookup": [ "registered_method_lookup_in_transport", ], - "xds_end2end_test": [ - "round_robin_delegate_to_pick_first", - "wrr_delegate_to_pick_first", - ], }, }, "posix": { @@ -235,10 +209,6 @@ EXPERIMENTS = { "cpp_end2end_test": [ "work_serializer_dispatch", ], - "cpp_lb_end2end_test": [ - "round_robin_delegate_to_pick_first", - "wrr_delegate_to_pick_first", - ], "credential_token_tests": [ "absl_base64", ], @@ -246,9 +216,7 @@ EXPERIMENTS = { "event_engine_listener", ], "lb_unit_test": [ - "round_robin_delegate_to_pick_first", "work_serializer_dispatch", - "wrr_delegate_to_pick_first", ], "resolver_component_tests_runner_invoker": [ "event_engine_dns", @@ -257,9 +225,7 @@ EXPERIMENTS = { "registered_method_lookup_in_transport", ], "xds_end2end_test": [ - "round_robin_delegate_to_pick_first", "work_serializer_dispatch", - "wrr_delegate_to_pick_first", ], }, }, diff --git a/build_autogenerated.yaml b/build_autogenerated.yaml index a3f46868f21..91e82652419 100644 --- a/build_autogenerated.yaml +++ b/build_autogenerated.yaml @@ -1183,7 +1183,6 @@ libs: - src/core/load_balancing/ring_hash/ring_hash.h - src/core/load_balancing/rls/rls.h - src/core/load_balancing/subchannel_interface.h - - src/core/load_balancing/subchannel_list.h - src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h - src/core/load_balancing/weighted_target/weighted_target.h - src/core/load_balancing/xds/xds_channel_args.h @@ -2655,7 +2654,6 @@ libs: - src/core/load_balancing/pick_first/pick_first.h - src/core/load_balancing/rls/rls.h - src/core/load_balancing/subchannel_interface.h - - src/core/load_balancing/subchannel_list.h - src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h - src/core/load_balancing/weighted_target/weighted_target.h - src/core/resolver/dns/c_ares/dns_resolver_ares.h diff --git a/gRPC-C++.podspec b/gRPC-C++.podspec index 898b33d3655..f459d1b7b85 100644 --- a/gRPC-C++.podspec +++ b/gRPC-C++.podspec @@ -1289,7 +1289,6 @@ Pod::Spec.new do |s| 'src/core/load_balancing/ring_hash/ring_hash.h', 'src/core/load_balancing/rls/rls.h', 'src/core/load_balancing/subchannel_interface.h', - 'src/core/load_balancing/subchannel_list.h', 'src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h', 'src/core/load_balancing/weighted_target/weighted_target.h', 'src/core/load_balancing/xds/xds_channel_args.h', @@ -2554,7 +2553,6 @@ Pod::Spec.new do |s| 'src/core/load_balancing/ring_hash/ring_hash.h', 'src/core/load_balancing/rls/rls.h', 'src/core/load_balancing/subchannel_interface.h', - 'src/core/load_balancing/subchannel_list.h', 'src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h', 'src/core/load_balancing/weighted_target/weighted_target.h', 'src/core/load_balancing/xds/xds_channel_args.h', diff --git a/gRPC-Core.podspec b/gRPC-Core.podspec index 8326e6b78bd..dd30ea8829c 100644 --- a/gRPC-Core.podspec +++ b/gRPC-Core.podspec @@ -2006,7 +2006,6 @@ Pod::Spec.new do |s| 'src/core/load_balancing/rls/rls.h', 'src/core/load_balancing/round_robin/round_robin.cc', 'src/core/load_balancing/subchannel_interface.h', - 'src/core/load_balancing/subchannel_list.h', 'src/core/load_balancing/weighted_round_robin/static_stride_scheduler.cc', 'src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h', 'src/core/load_balancing/weighted_round_robin/weighted_round_robin.cc', @@ -3336,7 +3335,6 @@ Pod::Spec.new do |s| 'src/core/load_balancing/ring_hash/ring_hash.h', 'src/core/load_balancing/rls/rls.h', 'src/core/load_balancing/subchannel_interface.h', - 'src/core/load_balancing/subchannel_list.h', 'src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h', 'src/core/load_balancing/weighted_target/weighted_target.h', 'src/core/load_balancing/xds/xds_channel_args.h', diff --git a/grpc.gemspec b/grpc.gemspec index 8b33f403ed2..64b5cf62e4b 100644 --- a/grpc.gemspec +++ b/grpc.gemspec @@ -1898,7 +1898,6 @@ Gem::Specification.new do |s| s.files += %w( src/core/load_balancing/rls/rls.h ) s.files += %w( src/core/load_balancing/round_robin/round_robin.cc ) s.files += %w( src/core/load_balancing/subchannel_interface.h ) - s.files += %w( src/core/load_balancing/subchannel_list.h ) s.files += %w( src/core/load_balancing/weighted_round_robin/static_stride_scheduler.cc ) s.files += %w( src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h ) s.files += %w( src/core/load_balancing/weighted_round_robin/weighted_round_robin.cc ) diff --git a/package.xml b/package.xml index cfe611d7459..84621351c14 100644 --- a/package.xml +++ b/package.xml @@ -1880,7 +1880,6 @@ - diff --git a/src/core/BUILD b/src/core/BUILD index c67173ba22b..857484df4f1 100644 --- a/src/core/BUILD +++ b/src/core/BUILD @@ -5470,35 +5470,6 @@ grpc_cc_library( ], ) -grpc_cc_library( - name = "grpc_lb_subchannel_list", - hdrs = [ - "load_balancing/subchannel_list.h", - ], - external_deps = [ - "absl/status", - "absl/types:optional", - ], - language = "c++", - deps = [ - "channel_args", - "connectivity_state", - "dual_ref_counted", - "gpr_manual_constructor", - "health_check_client", - "iomgr_fwd", - "lb_policy", - "subchannel_interface", - "//:debug_location", - "//:endpoint_addresses", - "//:gpr", - "//:grpc_base", - "//:ref_counted_ptr", - "//:server_address", - "//:work_serializer", - ], -) - grpc_cc_library( name = "lb_endpoint_list", srcs = [ @@ -5721,13 +5692,10 @@ grpc_cc_library( deps = [ "channel_args", "connectivity_state", - "experiments", - "grpc_lb_subchannel_list", "json", "lb_endpoint_list", "lb_policy", "lb_policy_factory", - "subchannel_interface", "//:config", "//:debug_location", "//:endpoint_addresses", @@ -5736,7 +5704,6 @@ grpc_cc_library( "//:grpc_trace", "//:orphanable", "//:ref_counted_ptr", - "//:server_address", "//:work_serializer", ], ) @@ -5780,7 +5747,6 @@ grpc_cc_library( "experiments", "grpc_backend_metric_data", "grpc_lb_policy_weighted_target", - "grpc_lb_subchannel_list", "json", "json_args", "json_object_loader", @@ -5806,8 +5772,6 @@ grpc_cc_library( "//:oob_backend_metric", "//:orphanable", "//:ref_counted_ptr", - "//:server_address", - "//:sockaddr_utils", "//:stats", "//:work_serializer", ], diff --git a/src/core/lib/experiments/experiments.cc b/src/core/lib/experiments/experiments.cc index 495570303b8..41307cf8843 100644 --- a/src/core/lib/experiments/experiments.cc +++ b/src/core/lib/experiments/experiments.cc @@ -110,11 +110,6 @@ const uint8_t required_experiments_promise_based_inproc_transport[] = { static_cast(grpc_core::kExperimentIdPromiseBasedServerCall), static_cast( grpc_core::kExperimentIdRegisteredMethodLookupInTransport)}; -const char* const description_round_robin_delegate_to_pick_first = - "Change round_robin code to delegate to pick_first as per dualstack " - "backend design."; -const char* const additional_constraints_round_robin_delegate_to_pick_first = - "{}"; const char* const description_rstpit = "On RST_STREAM on a server, reduce MAX_CONCURRENT_STREAMS for a short " "duration"; @@ -164,10 +159,6 @@ const char* const description_work_serializer_dispatch = const char* const additional_constraints_work_serializer_dispatch = "{}"; const uint8_t required_experiments_work_serializer_dispatch[] = { static_cast(grpc_core::kExperimentIdEventEngineClient)}; -const char* const description_wrr_delegate_to_pick_first = - "Change WRR code to delegate to pick_first as per dualstack backend " - "design."; -const char* const additional_constraints_wrr_delegate_to_pick_first = "{}"; #ifdef NDEBUG const bool kDefaultForDebugOnly = false; #else @@ -228,10 +219,6 @@ const ExperimentMetadata g_experiment_metadata[] = { description_promise_based_inproc_transport, additional_constraints_promise_based_inproc_transport, required_experiments_promise_based_inproc_transport, 3, false, false}, - {"round_robin_delegate_to_pick_first", - description_round_robin_delegate_to_pick_first, - additional_constraints_round_robin_delegate_to_pick_first, nullptr, 0, - true, true}, {"rstpit", description_rstpit, additional_constraints_rstpit, nullptr, 0, false, true}, {"schedule_cancellation_over_write", @@ -265,8 +252,6 @@ const ExperimentMetadata g_experiment_metadata[] = { {"work_serializer_dispatch", description_work_serializer_dispatch, additional_constraints_work_serializer_dispatch, required_experiments_work_serializer_dispatch, 1, false, true}, - {"wrr_delegate_to_pick_first", description_wrr_delegate_to_pick_first, - additional_constraints_wrr_delegate_to_pick_first, nullptr, 0, true, true}, }; } // namespace grpc_core @@ -359,11 +344,6 @@ const uint8_t required_experiments_promise_based_inproc_transport[] = { static_cast(grpc_core::kExperimentIdPromiseBasedServerCall), static_cast( grpc_core::kExperimentIdRegisteredMethodLookupInTransport)}; -const char* const description_round_robin_delegate_to_pick_first = - "Change round_robin code to delegate to pick_first as per dualstack " - "backend design."; -const char* const additional_constraints_round_robin_delegate_to_pick_first = - "{}"; const char* const description_rstpit = "On RST_STREAM on a server, reduce MAX_CONCURRENT_STREAMS for a short " "duration"; @@ -413,10 +393,6 @@ const char* const description_work_serializer_dispatch = const char* const additional_constraints_work_serializer_dispatch = "{}"; const uint8_t required_experiments_work_serializer_dispatch[] = { static_cast(grpc_core::kExperimentIdEventEngineClient)}; -const char* const description_wrr_delegate_to_pick_first = - "Change WRR code to delegate to pick_first as per dualstack backend " - "design."; -const char* const additional_constraints_wrr_delegate_to_pick_first = "{}"; #ifdef NDEBUG const bool kDefaultForDebugOnly = false; #else @@ -477,10 +453,6 @@ const ExperimentMetadata g_experiment_metadata[] = { description_promise_based_inproc_transport, additional_constraints_promise_based_inproc_transport, required_experiments_promise_based_inproc_transport, 3, false, false}, - {"round_robin_delegate_to_pick_first", - description_round_robin_delegate_to_pick_first, - additional_constraints_round_robin_delegate_to_pick_first, nullptr, 0, - true, true}, {"rstpit", description_rstpit, additional_constraints_rstpit, nullptr, 0, false, true}, {"schedule_cancellation_over_write", @@ -514,8 +486,6 @@ const ExperimentMetadata g_experiment_metadata[] = { {"work_serializer_dispatch", description_work_serializer_dispatch, additional_constraints_work_serializer_dispatch, required_experiments_work_serializer_dispatch, 1, false, true}, - {"wrr_delegate_to_pick_first", description_wrr_delegate_to_pick_first, - additional_constraints_wrr_delegate_to_pick_first, nullptr, 0, true, true}, }; } // namespace grpc_core @@ -608,11 +578,6 @@ const uint8_t required_experiments_promise_based_inproc_transport[] = { static_cast(grpc_core::kExperimentIdPromiseBasedServerCall), static_cast( grpc_core::kExperimentIdRegisteredMethodLookupInTransport)}; -const char* const description_round_robin_delegate_to_pick_first = - "Change round_robin code to delegate to pick_first as per dualstack " - "backend design."; -const char* const additional_constraints_round_robin_delegate_to_pick_first = - "{}"; const char* const description_rstpit = "On RST_STREAM on a server, reduce MAX_CONCURRENT_STREAMS for a short " "duration"; @@ -662,10 +627,6 @@ const char* const description_work_serializer_dispatch = const char* const additional_constraints_work_serializer_dispatch = "{}"; const uint8_t required_experiments_work_serializer_dispatch[] = { static_cast(grpc_core::kExperimentIdEventEngineClient)}; -const char* const description_wrr_delegate_to_pick_first = - "Change WRR code to delegate to pick_first as per dualstack backend " - "design."; -const char* const additional_constraints_wrr_delegate_to_pick_first = "{}"; #ifdef NDEBUG const bool kDefaultForDebugOnly = false; #else @@ -726,10 +687,6 @@ const ExperimentMetadata g_experiment_metadata[] = { description_promise_based_inproc_transport, additional_constraints_promise_based_inproc_transport, required_experiments_promise_based_inproc_transport, 3, false, false}, - {"round_robin_delegate_to_pick_first", - description_round_robin_delegate_to_pick_first, - additional_constraints_round_robin_delegate_to_pick_first, nullptr, 0, - true, true}, {"rstpit", description_rstpit, additional_constraints_rstpit, nullptr, 0, false, true}, {"schedule_cancellation_over_write", @@ -763,8 +720,6 @@ const ExperimentMetadata g_experiment_metadata[] = { {"work_serializer_dispatch", description_work_serializer_dispatch, additional_constraints_work_serializer_dispatch, required_experiments_work_serializer_dispatch, 1, true, true}, - {"wrr_delegate_to_pick_first", description_wrr_delegate_to_pick_first, - additional_constraints_wrr_delegate_to_pick_first, nullptr, 0, true, true}, }; } // namespace grpc_core diff --git a/src/core/lib/experiments/experiments.h b/src/core/lib/experiments/experiments.h index 70027ee9d99..90e66aee6f5 100644 --- a/src/core/lib/experiments/experiments.h +++ b/src/core/lib/experiments/experiments.h @@ -92,8 +92,6 @@ inline bool IsChaoticGoodEnabled() { return false; } #define GRPC_EXPERIMENT_IS_INCLUDED_REGISTERED_METHOD_LOOKUP_IN_TRANSPORT inline bool IsRegisteredMethodLookupInTransportEnabled() { return true; } inline bool IsPromiseBasedInprocTransportEnabled() { return false; } -#define GRPC_EXPERIMENT_IS_INCLUDED_ROUND_ROBIN_DELEGATE_TO_PICK_FIRST -inline bool IsRoundRobinDelegateToPickFirstEnabled() { return true; } inline bool IsRstpitEnabled() { return false; } inline bool IsScheduleCancellationOverWriteEnabled() { return false; } inline bool IsServerPrivacyEnabled() { return false; } @@ -108,8 +106,6 @@ inline bool IsV3ServerAuthFilterEnabled() { return false; } #define GRPC_EXPERIMENT_IS_INCLUDED_WORK_SERIALIZER_CLEARS_TIME_CACHE inline bool IsWorkSerializerClearsTimeCacheEnabled() { return true; } inline bool IsWorkSerializerDispatchEnabled() { return false; } -#define GRPC_EXPERIMENT_IS_INCLUDED_WRR_DELEGATE_TO_PICK_FIRST -inline bool IsWrrDelegateToPickFirstEnabled() { return true; } #elif defined(GPR_WINDOWS) #define GRPC_EXPERIMENT_IS_INCLUDED_ABSL_BASE64 @@ -148,8 +144,6 @@ inline bool IsChaoticGoodEnabled() { return false; } #define GRPC_EXPERIMENT_IS_INCLUDED_REGISTERED_METHOD_LOOKUP_IN_TRANSPORT inline bool IsRegisteredMethodLookupInTransportEnabled() { return true; } inline bool IsPromiseBasedInprocTransportEnabled() { return false; } -#define GRPC_EXPERIMENT_IS_INCLUDED_ROUND_ROBIN_DELEGATE_TO_PICK_FIRST -inline bool IsRoundRobinDelegateToPickFirstEnabled() { return true; } inline bool IsRstpitEnabled() { return false; } inline bool IsScheduleCancellationOverWriteEnabled() { return false; } inline bool IsServerPrivacyEnabled() { return false; } @@ -164,8 +158,6 @@ inline bool IsV3ServerAuthFilterEnabled() { return false; } #define GRPC_EXPERIMENT_IS_INCLUDED_WORK_SERIALIZER_CLEARS_TIME_CACHE inline bool IsWorkSerializerClearsTimeCacheEnabled() { return true; } inline bool IsWorkSerializerDispatchEnabled() { return false; } -#define GRPC_EXPERIMENT_IS_INCLUDED_WRR_DELEGATE_TO_PICK_FIRST -inline bool IsWrrDelegateToPickFirstEnabled() { return true; } #else #define GRPC_EXPERIMENT_IS_INCLUDED_ABSL_BASE64 @@ -205,8 +197,6 @@ inline bool IsChaoticGoodEnabled() { return false; } #define GRPC_EXPERIMENT_IS_INCLUDED_REGISTERED_METHOD_LOOKUP_IN_TRANSPORT inline bool IsRegisteredMethodLookupInTransportEnabled() { return true; } inline bool IsPromiseBasedInprocTransportEnabled() { return false; } -#define GRPC_EXPERIMENT_IS_INCLUDED_ROUND_ROBIN_DELEGATE_TO_PICK_FIRST -inline bool IsRoundRobinDelegateToPickFirstEnabled() { return true; } inline bool IsRstpitEnabled() { return false; } inline bool IsScheduleCancellationOverWriteEnabled() { return false; } inline bool IsServerPrivacyEnabled() { return false; } @@ -222,8 +212,6 @@ inline bool IsV3ServerAuthFilterEnabled() { return false; } inline bool IsWorkSerializerClearsTimeCacheEnabled() { return true; } #define GRPC_EXPERIMENT_IS_INCLUDED_WORK_SERIALIZER_DISPATCH inline bool IsWorkSerializerDispatchEnabled() { return true; } -#define GRPC_EXPERIMENT_IS_INCLUDED_WRR_DELEGATE_TO_PICK_FIRST -inline bool IsWrrDelegateToPickFirstEnabled() { return true; } #endif #else @@ -249,7 +237,6 @@ enum ExperimentIds { kExperimentIdChaoticGood, kExperimentIdRegisteredMethodLookupInTransport, kExperimentIdPromiseBasedInprocTransport, - kExperimentIdRoundRobinDelegateToPickFirst, kExperimentIdRstpit, kExperimentIdScheduleCancellationOverWrite, kExperimentIdServerPrivacy, @@ -263,7 +250,6 @@ enum ExperimentIds { kExperimentIdV3ServerAuthFilter, kExperimentIdWorkSerializerClearsTimeCache, kExperimentIdWorkSerializerDispatch, - kExperimentIdWrrDelegateToPickFirst, kNumExperiments }; #define GRPC_EXPERIMENT_IS_INCLUDED_ABSL_BASE64 @@ -350,10 +336,6 @@ inline bool IsRegisteredMethodLookupInTransportEnabled() { inline bool IsPromiseBasedInprocTransportEnabled() { return IsExperimentEnabled(kExperimentIdPromiseBasedInprocTransport); } -#define GRPC_EXPERIMENT_IS_INCLUDED_ROUND_ROBIN_DELEGATE_TO_PICK_FIRST -inline bool IsRoundRobinDelegateToPickFirstEnabled() { - return IsExperimentEnabled(kExperimentIdRoundRobinDelegateToPickFirst); -} #define GRPC_EXPERIMENT_IS_INCLUDED_RSTPIT inline bool IsRstpitEnabled() { return IsExperimentEnabled(kExperimentIdRstpit); @@ -406,10 +388,6 @@ inline bool IsWorkSerializerClearsTimeCacheEnabled() { inline bool IsWorkSerializerDispatchEnabled() { return IsExperimentEnabled(kExperimentIdWorkSerializerDispatch); } -#define GRPC_EXPERIMENT_IS_INCLUDED_WRR_DELEGATE_TO_PICK_FIRST -inline bool IsWrrDelegateToPickFirstEnabled() { - return IsExperimentEnabled(kExperimentIdWrrDelegateToPickFirst); -} extern const ExperimentMetadata g_experiment_metadata[kNumExperiments]; diff --git a/src/core/lib/experiments/experiments.yaml b/src/core/lib/experiments/experiments.yaml index b86fde5cbdd..f1c1f4a11fc 100644 --- a/src/core/lib/experiments/experiments.yaml +++ b/src/core/lib/experiments/experiments.yaml @@ -184,13 +184,6 @@ expiry: 2024/03/31 owner: yashkt@google.com test_tags: ["surface_registered_method_lookup"] -- name: round_robin_delegate_to_pick_first - description: - Change round_robin code to delegate to pick_first as per dualstack - backend design. - expiry: 2024/03/15 - owner: roth@google.com - test_tags: ["lb_unit_test", "cpp_lb_end2end_test", "xds_end2end_test"] - name: rstpit description: On RST_STREAM on a server, reduce MAX_CONCURRENT_STREAMS for a short duration @@ -271,10 +264,3 @@ expiry: 2024/03/31 owner: ysseung@google.com test_tags: ["core_end2end_test", "cpp_end2end_test", "xds_end2end_test", "lb_unit_test"] -- name: wrr_delegate_to_pick_first - description: - Change WRR code to delegate to pick_first as per dualstack - backend design. - expiry: 2024/03/15 - owner: roth@google.com - test_tags: ["lb_unit_test", "cpp_lb_end2end_test", "xds_end2end_test"] diff --git a/src/core/lib/experiments/rollouts.yaml b/src/core/lib/experiments/rollouts.yaml index 8bbd66af0ff..b6cbcd468a2 100644 --- a/src/core/lib/experiments/rollouts.yaml +++ b/src/core/lib/experiments/rollouts.yaml @@ -100,8 +100,6 @@ default: false - name: registered_method_lookup_in_transport default: true -- name: round_robin_delegate_to_pick_first - default: true - name: rstpit default: false - name: schedule_cancellation_over_write @@ -126,5 +124,3 @@ posix: true # TODO(ysseung): Test flakes not fully resolved. windows: broken -- name: wrr_delegate_to_pick_first - default: true diff --git a/src/core/load_balancing/round_robin/round_robin.cc b/src/core/load_balancing/round_robin/round_robin.cc index 66065aa06bf..1bfa55ccbaa 100644 --- a/src/core/load_balancing/round_robin/round_robin.cc +++ b/src/core/load_balancing/round_robin/round_robin.cc @@ -37,23 +37,19 @@ #include #include -#include "src/core/load_balancing/endpoint_list.h" -#include "src/core/load_balancing/subchannel_list.h" #include "src/core/lib/channel/channel_args.h" #include "src/core/lib/config/core_configuration.h" #include "src/core/lib/debug/trace.h" -#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/gprpp/debug_location.h" #include "src/core/lib/gprpp/orphanable.h" #include "src/core/lib/gprpp/ref_counted_ptr.h" #include "src/core/lib/gprpp/work_serializer.h" #include "src/core/lib/json/json.h" #include "src/core/lib/transport/connectivity_state.h" +#include "src/core/load_balancing/endpoint_list.h" #include "src/core/load_balancing/lb_policy.h" #include "src/core/load_balancing/lb_policy_factory.h" -#include "src/core/load_balancing/subchannel_interface.h" #include "src/core/resolver/endpoint_addresses.h" -#include "src/core/resolver/server_address.h" namespace grpc_core { @@ -61,456 +57,8 @@ TraceFlag grpc_lb_round_robin_trace(false, "round_robin"); namespace { -// -// legacy round_robin LB policy (before dualstack support) -// - constexpr absl::string_view kRoundRobin = "round_robin"; -class OldRoundRobin : public LoadBalancingPolicy { - public: - explicit OldRoundRobin(Args args); - - absl::string_view name() const override { return kRoundRobin; } - - absl::Status UpdateLocked(UpdateArgs args) override; - void ResetBackoffLocked() override; - - private: - ~OldRoundRobin() override; - - // Forward declaration. - class RoundRobinSubchannelList; - - // Data for a particular subchannel in a subchannel list. - // This subclass adds the following functionality: - // - Tracks the previous connectivity state of the subchannel, so that - // we know how many subchannels are in each state. - class RoundRobinSubchannelData - : public SubchannelData { - public: - RoundRobinSubchannelData( - SubchannelList* - subchannel_list, - const ServerAddress& address, - RefCountedPtr subchannel) - : SubchannelData(subchannel_list, address, std::move(subchannel)) {} - - absl::optional connectivity_state() const { - return logical_connectivity_state_; - } - - private: - // Performs connectivity state updates that need to be done only - // after we have started watching. - void ProcessConnectivityChangeLocked( - absl::optional old_state, - grpc_connectivity_state new_state) override; - - // Updates the logical connectivity state. - void UpdateLogicalConnectivityStateLocked( - grpc_connectivity_state connectivity_state); - - // The logical connectivity state of the subchannel. - // Note that the logical connectivity state may differ from the - // actual reported state in some cases (e.g., after we see - // TRANSIENT_FAILURE, we ignore any subsequent state changes until - // we see READY). - absl::optional logical_connectivity_state_; - }; - - // A list of subchannels. - class RoundRobinSubchannelList - : public SubchannelList { - public: - RoundRobinSubchannelList(OldRoundRobin* policy, - EndpointAddressesIterator* addresses, - const ChannelArgs& args) - : SubchannelList(policy, - (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace) - ? "RoundRobinSubchannelList" - : nullptr), - addresses, policy->channel_control_helper(), args) { - // Need to maintain a ref to the LB policy as long as we maintain - // any references to subchannels, since the subchannels' - // pollset_sets will include the LB policy's pollset_set. - policy->Ref(DEBUG_LOCATION, "subchannel_list").release(); - } - - ~RoundRobinSubchannelList() override { - OldRoundRobin* p = static_cast(policy()); - p->Unref(DEBUG_LOCATION, "subchannel_list"); - } - - // Updates the counters of subchannels in each state when a - // subchannel transitions from old_state to new_state. - void UpdateStateCountersLocked( - absl::optional old_state, - grpc_connectivity_state new_state); - - // Ensures that the right subchannel list is used and then updates - // the RR policy's connectivity state based on the subchannel list's - // state counters. - void MaybeUpdateRoundRobinConnectivityStateLocked( - absl::Status status_for_tf); - - private: - std::shared_ptr work_serializer() const override { - return static_cast(policy())->work_serializer(); - } - - std::string CountersString() const { - return absl::StrCat("num_subchannels=", num_subchannels(), - " num_ready=", num_ready_, - " num_connecting=", num_connecting_, - " num_transient_failure=", num_transient_failure_); - } - - size_t num_ready_ = 0; - size_t num_connecting_ = 0; - size_t num_transient_failure_ = 0; - - absl::Status last_failure_; - }; - - class Picker : public SubchannelPicker { - public: - Picker(OldRoundRobin* parent, RoundRobinSubchannelList* subchannel_list); - - PickResult Pick(PickArgs args) override; - - private: - // Using pointer value only, no ref held -- do not dereference! - OldRoundRobin* parent_; - - std::atomic last_picked_index_; - std::vector> subchannels_; - }; - - void ShutdownLocked() override; - - // List of subchannels. - RefCountedPtr subchannel_list_; - // Latest pending subchannel list. - // When we get an updated address list, we create a new subchannel list - // for it here, and we wait to swap it into subchannel_list_ until the new - // list becomes READY. - RefCountedPtr latest_pending_subchannel_list_; - - bool shutdown_ = false; - - absl::BitGen bit_gen_; -}; - -// -// OldRoundRobin::Picker -// - -OldRoundRobin::Picker::Picker(OldRoundRobin* parent, - RoundRobinSubchannelList* subchannel_list) - : parent_(parent) { - for (size_t i = 0; i < subchannel_list->num_subchannels(); ++i) { - RoundRobinSubchannelData* sd = subchannel_list->subchannel(i); - if (sd->connectivity_state().value_or(GRPC_CHANNEL_IDLE) == - GRPC_CHANNEL_READY) { - subchannels_.push_back(sd->subchannel()->Ref()); - } - } - // For discussion on why we generate a random starting index for - // the picker, see https://github.com/grpc/grpc-go/issues/2580. - size_t index = - absl::Uniform(parent->bit_gen_, 0, subchannels_.size()); - last_picked_index_.store(index, std::memory_order_relaxed); - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, - "[RR %p picker %p] created picker from subchannel_list=%p " - "with %" PRIuPTR " READY subchannels; last_picked_index_=%" PRIuPTR, - parent_, this, subchannel_list, subchannels_.size(), index); - } -} - -OldRoundRobin::PickResult OldRoundRobin::Picker::Pick(PickArgs /*args*/) { - size_t index = last_picked_index_.fetch_add(1, std::memory_order_relaxed) % - subchannels_.size(); - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, - "[RR %p picker %p] returning index %" PRIuPTR ", subchannel=%p", - parent_, this, index, subchannels_[index].get()); - } - return PickResult::Complete(subchannels_[index]); -} - -// -// RoundRobin -// - -OldRoundRobin::OldRoundRobin(Args args) : LoadBalancingPolicy(std::move(args)) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, "[RR %p] Created", this); - } -} - -OldRoundRobin::~OldRoundRobin() { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, "[RR %p] Destroying Round Robin policy", this); - } - GPR_ASSERT(subchannel_list_ == nullptr); - GPR_ASSERT(latest_pending_subchannel_list_ == nullptr); -} - -void OldRoundRobin::ShutdownLocked() { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, "[RR %p] Shutting down", this); - } - shutdown_ = true; - subchannel_list_.reset(); - latest_pending_subchannel_list_.reset(); -} - -void OldRoundRobin::ResetBackoffLocked() { - subchannel_list_->ResetBackoffLocked(); - if (latest_pending_subchannel_list_ != nullptr) { - latest_pending_subchannel_list_->ResetBackoffLocked(); - } -} - -absl::Status OldRoundRobin::UpdateLocked(UpdateArgs args) { - EndpointAddressesIterator* addresses = nullptr; - if (args.addresses.ok()) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, "[RR %p] received update", this); - } - addresses = args.addresses->get(); - } else { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, "[RR %p] received update with address error: %s", this, - args.addresses.status().ToString().c_str()); - } - // If we already have a subchannel list, then keep using the existing - // list, but still report back that the update was not accepted. - if (subchannel_list_ != nullptr) return args.addresses.status(); - } - // Create new subchannel list, replacing the previous pending list, if any. - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace) && - latest_pending_subchannel_list_ != nullptr) { - gpr_log(GPR_INFO, "[RR %p] replacing previous pending subchannel list %p", - this, latest_pending_subchannel_list_.get()); - } - latest_pending_subchannel_list_ = - MakeRefCounted(this, addresses, args.args); - latest_pending_subchannel_list_->StartWatchingLocked(args.args); - // If the new list is empty, immediately promote it to - // subchannel_list_ and report TRANSIENT_FAILURE. - if (latest_pending_subchannel_list_->num_subchannels() == 0) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace) && - subchannel_list_ != nullptr) { - gpr_log(GPR_INFO, "[RR %p] replacing previous subchannel list %p", this, - subchannel_list_.get()); - } - subchannel_list_ = std::move(latest_pending_subchannel_list_); - absl::Status status = - args.addresses.ok() ? absl::UnavailableError(absl::StrCat( - "empty address list: ", args.resolution_note)) - : args.addresses.status(); - channel_control_helper()->UpdateState( - GRPC_CHANNEL_TRANSIENT_FAILURE, status, - MakeRefCounted(status)); - return status; - } - // Otherwise, if this is the initial update, immediately promote it to - // subchannel_list_. - if (subchannel_list_.get() == nullptr) { - subchannel_list_ = std::move(latest_pending_subchannel_list_); - } - return absl::OkStatus(); -} - -// -// RoundRobinSubchannelList -// - -void OldRoundRobin::RoundRobinSubchannelList::UpdateStateCountersLocked( - absl::optional old_state, - grpc_connectivity_state new_state) { - if (old_state.has_value()) { - GPR_ASSERT(*old_state != GRPC_CHANNEL_SHUTDOWN); - if (*old_state == GRPC_CHANNEL_READY) { - GPR_ASSERT(num_ready_ > 0); - --num_ready_; - } else if (*old_state == GRPC_CHANNEL_CONNECTING) { - GPR_ASSERT(num_connecting_ > 0); - --num_connecting_; - } else if (*old_state == GRPC_CHANNEL_TRANSIENT_FAILURE) { - GPR_ASSERT(num_transient_failure_ > 0); - --num_transient_failure_; - } - } - GPR_ASSERT(new_state != GRPC_CHANNEL_SHUTDOWN); - if (new_state == GRPC_CHANNEL_READY) { - ++num_ready_; - } else if (new_state == GRPC_CHANNEL_CONNECTING) { - ++num_connecting_; - } else if (new_state == GRPC_CHANNEL_TRANSIENT_FAILURE) { - ++num_transient_failure_; - } -} - -void OldRoundRobin::RoundRobinSubchannelList:: - MaybeUpdateRoundRobinConnectivityStateLocked(absl::Status status_for_tf) { - OldRoundRobin* p = static_cast(policy()); - // If this is latest_pending_subchannel_list_, then swap it into - // subchannel_list_ in the following cases: - // - subchannel_list_ has no READY subchannels. - // - This list has at least one READY subchannel and we have seen the - // initial connectivity state notification for all subchannels. - // - All of the subchannels in this list are in TRANSIENT_FAILURE. - // (This may cause the channel to go from READY to TRANSIENT_FAILURE, - // but we're doing what the control plane told us to do.) - if (p->latest_pending_subchannel_list_.get() == this && - (p->subchannel_list_->num_ready_ == 0 || - (num_ready_ > 0 && AllSubchannelsSeenInitialState()) || - num_transient_failure_ == num_subchannels())) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - const std::string old_counters_string = - p->subchannel_list_ != nullptr ? p->subchannel_list_->CountersString() - : ""; - gpr_log( - GPR_INFO, - "[RR %p] swapping out subchannel list %p (%s) in favor of %p (%s)", p, - p->subchannel_list_.get(), old_counters_string.c_str(), this, - CountersString().c_str()); - } - p->subchannel_list_ = std::move(p->latest_pending_subchannel_list_); - } - // Only set connectivity state if this is the current subchannel list. - if (p->subchannel_list_.get() != this) return; - // First matching rule wins: - // 1) ANY subchannel is READY => policy is READY. - // 2) ANY subchannel is CONNECTING => policy is CONNECTING. - // 3) ALL subchannels are TRANSIENT_FAILURE => policy is TRANSIENT_FAILURE. - if (num_ready_ > 0) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, "[RR %p] reporting READY with subchannel list %p", p, - this); - } - p->channel_control_helper()->UpdateState(GRPC_CHANNEL_READY, absl::Status(), - MakeRefCounted(p, this)); - } else if (num_connecting_ > 0) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, "[RR %p] reporting CONNECTING with subchannel list %p", - p, this); - } - p->channel_control_helper()->UpdateState( - GRPC_CHANNEL_CONNECTING, absl::Status(), - MakeRefCounted( - p->RefAsSubclass(DEBUG_LOCATION, "QueuePicker"))); - } else if (num_transient_failure_ == num_subchannels()) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, - "[RR %p] reporting TRANSIENT_FAILURE with subchannel list %p: %s", - p, this, status_for_tf.ToString().c_str()); - } - if (!status_for_tf.ok()) { - last_failure_ = absl::UnavailableError( - absl::StrCat("connections to all backends failing; last error: ", - status_for_tf.ToString())); - } - p->channel_control_helper()->UpdateState( - GRPC_CHANNEL_TRANSIENT_FAILURE, last_failure_, - MakeRefCounted(last_failure_)); - } -} - -// -// RoundRobinSubchannelData -// - -void OldRoundRobin::RoundRobinSubchannelData::ProcessConnectivityChangeLocked( - absl::optional old_state, - grpc_connectivity_state new_state) { - OldRoundRobin* p = static_cast(subchannel_list()->policy()); - GPR_ASSERT(subchannel() != nullptr); - // If this is not the initial state notification and the new state is - // TRANSIENT_FAILURE or IDLE, re-resolve. - // Note that we don't want to do this on the initial state notification, - // because that would result in an endless loop of re-resolution. - if (old_state.has_value() && (new_state == GRPC_CHANNEL_TRANSIENT_FAILURE || - new_state == GRPC_CHANNEL_IDLE)) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, - "[RR %p] Subchannel %p reported %s; requesting re-resolution", p, - subchannel(), ConnectivityStateName(new_state)); - } - p->channel_control_helper()->RequestReresolution(); - } - if (new_state == GRPC_CHANNEL_IDLE) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, - "[RR %p] Subchannel %p reported IDLE; requesting connection", p, - subchannel()); - } - subchannel()->RequestConnection(); - } - // Update logical connectivity state. - UpdateLogicalConnectivityStateLocked(new_state); - // Update the policy state. - subchannel_list()->MaybeUpdateRoundRobinConnectivityStateLocked( - connectivity_status()); -} - -void OldRoundRobin::RoundRobinSubchannelData:: - UpdateLogicalConnectivityStateLocked( - grpc_connectivity_state connectivity_state) { - OldRoundRobin* p = static_cast(subchannel_list()->policy()); - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log( - GPR_INFO, - "[RR %p] connectivity changed for subchannel %p, subchannel_list %p " - "(index %" PRIuPTR " of %" PRIuPTR "): prev_state=%s new_state=%s", - p, subchannel(), subchannel_list(), Index(), - subchannel_list()->num_subchannels(), - (logical_connectivity_state_.has_value() - ? ConnectivityStateName(*logical_connectivity_state_) - : "N/A"), - ConnectivityStateName(connectivity_state)); - } - // Decide what state to report for aggregation purposes. - // If the last logical state was TRANSIENT_FAILURE, then ignore the - // state change unless the new state is READY. - if (logical_connectivity_state_.has_value() && - *logical_connectivity_state_ == GRPC_CHANNEL_TRANSIENT_FAILURE && - connectivity_state != GRPC_CHANNEL_READY) { - return; - } - // If the new state is IDLE, treat it as CONNECTING, since it will - // immediately transition into CONNECTING anyway. - if (connectivity_state == GRPC_CHANNEL_IDLE) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_round_robin_trace)) { - gpr_log(GPR_INFO, - "[RR %p] subchannel %p, subchannel_list %p (index %" PRIuPTR - " of %" PRIuPTR "): treating IDLE as CONNECTING", - p, subchannel(), subchannel_list(), Index(), - subchannel_list()->num_subchannels()); - } - connectivity_state = GRPC_CHANNEL_CONNECTING; - } - // If no change, return false. - if (logical_connectivity_state_.has_value() && - *logical_connectivity_state_ == connectivity_state) { - return; - } - // Otherwise, update counters and logical state. - subchannel_list()->UpdateStateCountersLocked(logical_connectivity_state_, - connectivity_state); - logical_connectivity_state_ = connectivity_state; -} - -// -// round_robin LB policy (with dualstack changes) -// - class RoundRobin : public LoadBalancingPolicy { public: explicit RoundRobin(Args args); @@ -892,9 +440,6 @@ class RoundRobinFactory : public LoadBalancingPolicyFactory { public: OrphanablePtr CreateLoadBalancingPolicy( LoadBalancingPolicy::Args args) const override { - if (!IsRoundRobinDelegateToPickFirstEnabled()) { - return MakeOrphanable(std::move(args)); - } return MakeOrphanable(std::move(args)); } diff --git a/src/core/load_balancing/subchannel_list.h b/src/core/load_balancing/subchannel_list.h deleted file mode 100644 index da9e35272a6..00000000000 --- a/src/core/load_balancing/subchannel_list.h +++ /dev/null @@ -1,455 +0,0 @@ -// -// Copyright 2015 gRPC authors. -// -// Licensed 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. -// - -#ifndef GRPC_SRC_CORE_LOAD_BALANCING_SUBCHANNEL_LIST_H -#define GRPC_SRC_CORE_LOAD_BALANCING_SUBCHANNEL_LIST_H - -#include - -#include -#include - -#include -#include -#include - -#include "absl/status/status.h" -#include "absl/types/optional.h" - -#include -#include - -#include "src/core/load_balancing/health_check_client.h" -#include "src/core/lib/channel/channel_args.h" -#include "src/core/lib/gprpp/debug_location.h" -#include "src/core/lib/gprpp/dual_ref_counted.h" -#include "src/core/lib/gprpp/manual_constructor.h" -#include "src/core/lib/gprpp/ref_counted_ptr.h" -#include "src/core/lib/gprpp/work_serializer.h" -#include "src/core/lib/iomgr/iomgr_fwd.h" -#include "src/core/lib/transport/connectivity_state.h" -#include "src/core/load_balancing/lb_policy.h" -#include "src/core/load_balancing/subchannel_interface.h" -#include "src/core/resolver/endpoint_addresses.h" -#include "src/core/resolver/server_address.h" - -// Code for maintaining a list of subchannels within an LB policy. -// -// To use this, callers must create their own subclasses, like so: -// - -// class MySubchannelList; // Forward declaration. - -// class MySubchannelData -// : public SubchannelData { -// public: -// void ProcessConnectivityChangeLocked( -// absl::optional old_state, -// grpc_connectivity_state new_state) override { -// // ...code to handle connectivity changes... -// } -// }; - -// class MySubchannelList -// : public SubchannelList { -// }; - -// -// All methods will be called from within the client_channel work serializer. - -namespace grpc_core { - -// Forward declaration. -template -class SubchannelList; - -// Stores data for a particular subchannel in a subchannel list. -// Callers must create a subclass that implements the -// ProcessConnectivityChangeLocked() method. -template -class SubchannelData { - public: - // Returns a pointer to the subchannel list containing this object. - SubchannelListType* subchannel_list() const { - return static_cast(subchannel_list_); - } - - // Returns the index into the subchannel list of this object. - size_t Index() const { - return static_cast(static_cast(this) - - subchannel_list_->subchannel(0)); - } - - // Returns a pointer to the subchannel. - SubchannelInterface* subchannel() const { return subchannel_.get(); } - - // Returns the cached connectivity state, if any. - absl::optional connectivity_state() { - return connectivity_state_; - } - absl::Status connectivity_status() { return connectivity_status_; } - - // Resets the connection backoff. - void ResetBackoffLocked(); - - // Cancels any pending connectivity watch and unrefs the subchannel. - void ShutdownLocked(); - - protected: - SubchannelData( - SubchannelList* subchannel_list, - const ServerAddress& address, - RefCountedPtr subchannel); - - virtual ~SubchannelData(); - - // This method will be invoked once soon after instantiation to report - // the current connectivity state, and it will then be invoked again - // whenever the connectivity state changes. - virtual void ProcessConnectivityChangeLocked( - absl::optional old_state, - grpc_connectivity_state new_state) = 0; - - private: - // For accessing StartConnectivityWatchLocked(). - friend class SubchannelList; - - // Watcher for subchannel connectivity state. - class Watcher - : public SubchannelInterface::ConnectivityStateWatcherInterface { - public: - Watcher( - SubchannelData* subchannel_data, - WeakRefCountedPtr subchannel_list) - : subchannel_data_(subchannel_data), - subchannel_list_(std::move(subchannel_list)) {} - - ~Watcher() override { - subchannel_list_.reset(DEBUG_LOCATION, "Watcher dtor"); - } - - void OnConnectivityStateChange(grpc_connectivity_state new_state, - absl::Status status) override; - - grpc_pollset_set* interested_parties() override { - return subchannel_list_->policy()->interested_parties(); - } - - private: - SubchannelData* subchannel_data_; - WeakRefCountedPtr subchannel_list_; - }; - - // Starts watching the connectivity state of the subchannel. - // ProcessConnectivityChangeLocked() will be called whenever the - // connectivity state changes. - void StartConnectivityWatchLocked(const ChannelArgs& args); - - // Cancels watching the connectivity state of the subchannel. - void CancelConnectivityWatchLocked(const char* reason); - - // Unrefs the subchannel. - void UnrefSubchannelLocked(const char* reason); - - // Backpointer to owning subchannel list. Not owned. - SubchannelList* subchannel_list_; - // The subchannel. - RefCountedPtr subchannel_; - // Will be non-null when the subchannel's state is being watched. - SubchannelInterface::DataWatcherInterface* health_watcher_ = nullptr; - // Data updated by the watcher. - absl::optional connectivity_state_; - absl::Status connectivity_status_; -}; - -// A list of subchannels. -template -class SubchannelList : public DualRefCounted { - public: - // Starts watching the connectivity state of all subchannels. - // Must be called immediately after instantiation. - void StartWatchingLocked(const ChannelArgs& args); - - // The number of subchannels in the list. - size_t num_subchannels() const { return subchannels_.size(); } - - // The data for the subchannel at a particular index. - SubchannelDataType* subchannel(size_t index) { - return subchannels_[index].get(); - } - - // Returns true if the subchannel list is shutting down. - bool shutting_down() const { return shutting_down_; } - - // Accessors. - LoadBalancingPolicy* policy() const { return policy_; } - const char* tracer() const { return tracer_; } - - // Resets connection backoff of all subchannels. - void ResetBackoffLocked(); - - // Returns true if all subchannels have seen their initial - // connectivity state notifications. - bool AllSubchannelsSeenInitialState(); - - void Orphan() override; - - protected: - SubchannelList(LoadBalancingPolicy* policy, const char* tracer, - EndpointAddressesIterator* addresses, - LoadBalancingPolicy::ChannelControlHelper* helper, - const ChannelArgs& args); - - virtual ~SubchannelList(); - - private: - // For accessing Ref() and Unref(). - friend class SubchannelData; - - virtual std::shared_ptr work_serializer() const = 0; - - // Backpointer to owning policy. - LoadBalancingPolicy* policy_; - - const char* tracer_; - - // The list of subchannels. - // We use ManualConstructor here to support SubchannelDataType classes - // that are not copyable. - std::vector> subchannels_; - - // Is this list shutting down? This may be true due to the shutdown of the - // policy itself or because a newer update has arrived while this one hadn't - // finished processing. - bool shutting_down_ = false; -}; - -// -// implementation -- no user-servicable parts below -// - -// -// SubchannelData::Watcher -// - -template -void SubchannelData::Watcher:: - OnConnectivityStateChange(grpc_connectivity_state new_state, - absl::Status status) { - if (GPR_UNLIKELY(subchannel_list_->tracer() != nullptr)) { - gpr_log( - GPR_INFO, - "[%s %p] subchannel list %p index %" PRIuPTR " of %" PRIuPTR - " (subchannel %p): connectivity changed: old_state=%s, new_state=%s, " - "status=%s, shutting_down=%d, health_watcher=%p", - subchannel_list_->tracer(), subchannel_list_->policy(), - subchannel_list_.get(), subchannel_data_->Index(), - subchannel_list_->num_subchannels(), - subchannel_data_->subchannel_.get(), - (subchannel_data_->connectivity_state_.has_value() - ? ConnectivityStateName(*subchannel_data_->connectivity_state_) - : "N/A"), - ConnectivityStateName(new_state), status.ToString().c_str(), - subchannel_list_->shutting_down(), subchannel_data_->health_watcher_); - } - if (!subchannel_list_->shutting_down() && - subchannel_data_->health_watcher_ != nullptr) { - absl::optional old_state = - subchannel_data_->connectivity_state_; - subchannel_data_->connectivity_state_ = new_state; - subchannel_data_->connectivity_status_ = status; - // Call the subclass's ProcessConnectivityChangeLocked() method. - subchannel_data_->ProcessConnectivityChangeLocked(old_state, new_state); - } -} - -// -// SubchannelData -// - -template -SubchannelData::SubchannelData( - SubchannelList* subchannel_list, - const ServerAddress& /*address*/, - RefCountedPtr subchannel) - : subchannel_list_(subchannel_list), subchannel_(std::move(subchannel)) {} - -template -SubchannelData::~SubchannelData() { - GPR_ASSERT(subchannel_ == nullptr); -} - -template -void SubchannelData:: - UnrefSubchannelLocked(const char* reason) { - if (subchannel_ != nullptr) { - if (GPR_UNLIKELY(subchannel_list_->tracer() != nullptr)) { - gpr_log(GPR_INFO, - "[%s %p] subchannel list %p index %" PRIuPTR " of %" PRIuPTR - " (subchannel %p): unreffing subchannel (%s)", - subchannel_list_->tracer(), subchannel_list_->policy(), - subchannel_list_, Index(), subchannel_list_->num_subchannels(), - subchannel_.get(), reason); - } - subchannel_.reset(); - } -} - -template -void SubchannelData::ResetBackoffLocked() { - if (subchannel_ != nullptr) { - subchannel_->ResetBackoff(); - } -} - -template -void SubchannelData:: - StartConnectivityWatchLocked(const ChannelArgs& args) { - if (GPR_UNLIKELY(subchannel_list_->tracer() != nullptr)) { - gpr_log(GPR_INFO, - "[%s %p] subchannel list %p index %" PRIuPTR " of %" PRIuPTR - " (subchannel %p): starting watch", - subchannel_list_->tracer(), subchannel_list_->policy(), - subchannel_list_, Index(), subchannel_list_->num_subchannels(), - subchannel_.get()); - } - GPR_ASSERT(health_watcher_ == nullptr); - auto watcher = std::make_unique( - this, subchannel_list()->WeakRef(DEBUG_LOCATION, "Watcher")); - auto health_watcher = MakeHealthCheckWatcher( - subchannel_list_->work_serializer(), args, std::move(watcher)); - health_watcher_ = health_watcher.get(); - subchannel_->AddDataWatcher(std::move(health_watcher)); -} - -template -void SubchannelData:: - CancelConnectivityWatchLocked(const char* reason) { - if (health_watcher_ != nullptr) { - if (GPR_UNLIKELY(subchannel_list_->tracer() != nullptr)) { - gpr_log(GPR_INFO, - "[%s %p] subchannel list %p index %" PRIuPTR " of %" PRIuPTR - " (subchannel %p): canceling health watch (%s)", - subchannel_list_->tracer(), subchannel_list_->policy(), - subchannel_list_, Index(), subchannel_list_->num_subchannels(), - subchannel_.get(), reason); - } - subchannel_->CancelDataWatcher(health_watcher_); - health_watcher_ = nullptr; - } -} - -template -void SubchannelData::ShutdownLocked() { - CancelConnectivityWatchLocked("shutdown"); - UnrefSubchannelLocked("shutdown"); -} - -// -// SubchannelList -// - -template -SubchannelList::SubchannelList( - LoadBalancingPolicy* policy, const char* tracer, - EndpointAddressesIterator* addresses, - LoadBalancingPolicy::ChannelControlHelper* helper, const ChannelArgs& args) - : DualRefCounted(tracer), - policy_(policy), - tracer_(tracer) { - if (GPR_UNLIKELY(tracer_ != nullptr)) { - gpr_log(GPR_INFO, "[%s %p] Creating subchannel list %p", tracer_, policy, - this); - } - if (addresses == nullptr) return; - // Create a subchannel for each address. - addresses->ForEach([&](const EndpointAddresses& address) { - RefCountedPtr subchannel = - helper->CreateSubchannel(address.address(), address.args(), args); - if (subchannel == nullptr) { - // Subchannel could not be created. - if (GPR_UNLIKELY(tracer_ != nullptr)) { - gpr_log(GPR_INFO, - "[%s %p] could not create subchannel for address %s, ignoring", - tracer_, policy_, address.ToString().c_str()); - } - return; - } - if (GPR_UNLIKELY(tracer_ != nullptr)) { - gpr_log(GPR_INFO, - "[%s %p] subchannel list %p index %" PRIuPTR - ": Created subchannel %p for address %s", - tracer_, policy_, this, subchannels_.size(), subchannel.get(), - address.ToString().c_str()); - } - subchannels_.emplace_back(); - subchannels_.back().Init(this, address, std::move(subchannel)); - }); -} - -template -SubchannelList::~SubchannelList() { - if (GPR_UNLIKELY(tracer_ != nullptr)) { - gpr_log(GPR_INFO, "[%s %p] Destroying subchannel_list %p", tracer_, policy_, - this); - } - for (auto& sd : subchannels_) { - sd.Destroy(); - } -} - -template -void SubchannelList:: - StartWatchingLocked(const ChannelArgs& args) { - for (auto& sd : subchannels_) { - sd->StartConnectivityWatchLocked(args); - } -} - -template -void SubchannelList::Orphan() { - if (GPR_UNLIKELY(tracer_ != nullptr)) { - gpr_log(GPR_INFO, "[%s %p] Shutting down subchannel_list %p", tracer_, - policy_, this); - } - GPR_ASSERT(!shutting_down_); - shutting_down_ = true; - for (auto& sd : subchannels_) { - sd->ShutdownLocked(); - } -} - -template -void SubchannelList::ResetBackoffLocked() { - for (auto& sd : subchannels_) { - sd->ResetBackoffLocked(); - } -} - -template -bool SubchannelList::AllSubchannelsSeenInitialState() { - for (size_t i = 0; i < num_subchannels(); ++i) { - if (!subchannel(i)->connectivity_state().has_value()) return false; - } - return true; -} - -} // namespace grpc_core - -#endif // GRPC_SRC_CORE_LOAD_BALANCING_SUBCHANNEL_LIST_H diff --git a/src/core/load_balancing/weighted_round_robin/weighted_round_robin.cc b/src/core/load_balancing/weighted_round_robin/weighted_round_robin.cc index 9a049a5a03e..14f01fb267e 100644 --- a/src/core/load_balancing/weighted_round_robin/weighted_round_robin.cc +++ b/src/core/load_balancing/weighted_round_robin/weighted_round_robin.cc @@ -18,11 +18,9 @@ #include #include -#include #include #include -#include #include #include #include @@ -46,13 +44,6 @@ #include #include -#include "src/core/load_balancing/backend_metric_data.h" -#include "src/core/load_balancing/endpoint_list.h" -#include "src/core/load_balancing/oob_backend_metric.h" -#include "src/core/load_balancing/subchannel_list.h" -#include "src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h" -#include "src/core/load_balancing/weighted_target/weighted_target.h" -#include "src/core/lib/address_utils/sockaddr_utils.h" #include "src/core/lib/channel/channel_args.h" #include "src/core/lib/channel/metrics.h" #include "src/core/lib/config/core_configuration.h" @@ -74,11 +65,15 @@ #include "src/core/lib/json/json_args.h" #include "src/core/lib/json/json_object_loader.h" #include "src/core/lib/transport/connectivity_state.h" +#include "src/core/load_balancing/backend_metric_data.h" +#include "src/core/load_balancing/endpoint_list.h" +#include "src/core/load_balancing/oob_backend_metric.h" +#include "src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h" +#include "src/core/load_balancing/weighted_target/weighted_target.h" #include "src/core/load_balancing/lb_policy.h" #include "src/core/load_balancing/lb_policy_factory.h" #include "src/core/load_balancing/subchannel_interface.h" #include "src/core/resolver/endpoint_addresses.h" -#include "src/core/resolver/server_address.h" namespace grpc_core { @@ -184,858 +179,7 @@ class WeightedRoundRobinConfig : public LoadBalancingPolicy::Config { float error_utilization_penalty_ = 1.0; }; -// Legacy WRR LB policy (not delegating to pick_first) -class OldWeightedRoundRobin : public LoadBalancingPolicy { - public: - explicit OldWeightedRoundRobin(Args args); - - absl::string_view name() const override { return kWeightedRoundRobin; } - - absl::Status UpdateLocked(UpdateArgs args) override; - void ResetBackoffLocked() override; - - private: - // Represents the weight for a given address. - class AddressWeight : public RefCounted { - public: - AddressWeight(RefCountedPtr wrr, std::string key) - : wrr_(std::move(wrr)), key_(std::move(key)) {} - ~AddressWeight() override; - - void MaybeUpdateWeight(double qps, double eps, double utilization, - float error_utilization_penalty); - - float GetWeight(Timestamp now, Duration weight_expiration_period, - Duration blackout_period); - - void ResetNonEmptySince(); - - private: - RefCountedPtr wrr_; - const std::string key_; - - Mutex mu_; - float weight_ ABSL_GUARDED_BY(&mu_) = 0; - Timestamp non_empty_since_ ABSL_GUARDED_BY(&mu_) = Timestamp::InfFuture(); - Timestamp last_update_time_ ABSL_GUARDED_BY(&mu_) = Timestamp::InfPast(); - }; - - // Forward declaration. - class WeightedRoundRobinSubchannelList; - - // Data for a particular subchannel in a subchannel list. - // This subclass adds the following functionality: - // - Tracks the previous connectivity state of the subchannel, so that - // we know how many subchannels are in each state. - class WeightedRoundRobinSubchannelData - : public SubchannelData { - public: - WeightedRoundRobinSubchannelData( - SubchannelList* subchannel_list, - const ServerAddress& address, RefCountedPtr sc); - - absl::optional connectivity_state() const { - return logical_connectivity_state_; - } - - RefCountedPtr weight() const { return weight_; } - - private: - class OobWatcher : public OobBackendMetricWatcher { - public: - OobWatcher(RefCountedPtr weight, - float error_utilization_penalty) - : weight_(std::move(weight)), - error_utilization_penalty_(error_utilization_penalty) {} - - void OnBackendMetricReport( - const BackendMetricData& backend_metric_data) override; - - private: - RefCountedPtr weight_; - const float error_utilization_penalty_; - }; - - // Performs connectivity state updates that need to be done only - // after we have started watching. - void ProcessConnectivityChangeLocked( - absl::optional old_state, - grpc_connectivity_state new_state) override; - - // Updates the logical connectivity state. - void UpdateLogicalConnectivityStateLocked( - grpc_connectivity_state connectivity_state); - - // The logical connectivity state of the subchannel. - // Note that the logical connectivity state may differ from the - // actual reported state in some cases (e.g., after we see - // TRANSIENT_FAILURE, we ignore any subsequent state changes until - // we see READY). - absl::optional logical_connectivity_state_; - - RefCountedPtr weight_; - }; - - // A list of subchannels. - class WeightedRoundRobinSubchannelList - : public SubchannelList { - public: - WeightedRoundRobinSubchannelList(OldWeightedRoundRobin* policy, - EndpointAddressesIterator* addresses, - const ChannelArgs& args) - : SubchannelList(policy, - (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace) - ? "WeightedRoundRobinSubchannelList" - : nullptr), - addresses, policy->channel_control_helper(), args) { - // Need to maintain a ref to the LB policy as long as we maintain - // any references to subchannels, since the subchannels' - // pollset_sets will include the LB policy's pollset_set. - policy->Ref(DEBUG_LOCATION, "subchannel_list").release(); - } - - ~WeightedRoundRobinSubchannelList() override { - OldWeightedRoundRobin* p = static_cast(policy()); - p->Unref(DEBUG_LOCATION, "subchannel_list"); - } - - // Updates the counters of subchannels in each state when a - // subchannel transitions from old_state to new_state. - void UpdateStateCountersLocked( - absl::optional old_state, - grpc_connectivity_state new_state); - - // Ensures that the right subchannel list is used and then updates - // the aggregated connectivity state based on the subchannel list's - // state counters. - void MaybeUpdateAggregatedConnectivityStateLocked( - absl::Status status_for_tf); - - private: - std::shared_ptr work_serializer() const override { - return static_cast(policy())->work_serializer(); - } - - std::string CountersString() const { - return absl::StrCat("num_subchannels=", num_subchannels(), - " num_ready=", num_ready_, - " num_connecting=", num_connecting_, - " num_transient_failure=", num_transient_failure_); - } - - size_t num_ready_ = 0; - size_t num_connecting_ = 0; - size_t num_transient_failure_ = 0; - - absl::Status last_failure_; - }; - - // A picker that performs WRR picks with weights based on - // endpoint-reported utilization and QPS. - class Picker : public SubchannelPicker { - public: - Picker(RefCountedPtr wrr, - WeightedRoundRobinSubchannelList* subchannel_list); - - ~Picker() override; - - PickResult Pick(PickArgs args) override; - - void Orphan() override; - - private: - // A call tracker that collects per-call endpoint utilization reports. - class SubchannelCallTracker : public SubchannelCallTrackerInterface { - public: - SubchannelCallTracker(RefCountedPtr weight, - float error_utilization_penalty) - : weight_(std::move(weight)), - error_utilization_penalty_(error_utilization_penalty) {} - - void Start() override {} - - void Finish(FinishArgs args) override; - - private: - RefCountedPtr weight_; - const float error_utilization_penalty_; - }; - - // Info stored about each subchannel. - struct SubchannelInfo { - SubchannelInfo(RefCountedPtr subchannel, - RefCountedPtr weight) - : subchannel(std::move(subchannel)), weight(std::move(weight)) {} - - RefCountedPtr subchannel; - RefCountedPtr weight; - }; - - // Returns the index into subchannels_ to be picked. - size_t PickIndex(); - - // Builds a new scheduler and swaps it into place, then starts a - // timer for the next update. - void BuildSchedulerAndStartTimerLocked() - ABSL_EXCLUSIVE_LOCKS_REQUIRED(&timer_mu_); - - RefCountedPtr wrr_; - RefCountedPtr config_; - std::vector subchannels_; - - Mutex scheduler_mu_; - std::shared_ptr scheduler_ - ABSL_GUARDED_BY(&scheduler_mu_); - - Mutex timer_mu_ ABSL_ACQUIRED_BEFORE(&scheduler_mu_); - absl::optional - timer_handle_ ABSL_GUARDED_BY(&timer_mu_); - - // Used when falling back to RR. - std::atomic last_picked_index_; - }; - - ~OldWeightedRoundRobin() override; - - void ShutdownLocked() override; - - RefCountedPtr GetOrCreateWeight( - const grpc_resolved_address& address); - - RefCountedPtr config_; - - // List of subchannels. - RefCountedPtr subchannel_list_; - // Latest pending subchannel list. - // When we get an updated address list, we create a new subchannel list - // for it here, and we wait to swap it into subchannel_list_ until the new - // list becomes READY. - RefCountedPtr - latest_pending_subchannel_list_; - - Mutex address_weight_map_mu_; - std::map> address_weight_map_ - ABSL_GUARDED_BY(&address_weight_map_mu_); - - bool shutdown_ = false; - - absl::BitGen bit_gen_; - - // Accessed by picker. - std::atomic scheduler_state_{absl::Uniform(bit_gen_)}; -}; - -// -// OldWeightedRoundRobin::AddressWeight -// - -OldWeightedRoundRobin::AddressWeight::~AddressWeight() { - MutexLock lock(&wrr_->address_weight_map_mu_); - auto it = wrr_->address_weight_map_.find(key_); - if (it != wrr_->address_weight_map_.end() && it->second == this) { - wrr_->address_weight_map_.erase(it); - } -} - -void OldWeightedRoundRobin::AddressWeight::MaybeUpdateWeight( - double qps, double eps, double utilization, - float error_utilization_penalty) { - // Compute weight. - float weight = 0; - if (qps > 0 && utilization > 0) { - double penalty = 0.0; - if (eps > 0 && error_utilization_penalty > 0) { - penalty = eps / qps * error_utilization_penalty; - } - weight = qps / (utilization + penalty); - } - if (weight == 0) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, - "[WRR %p] subchannel %s: qps=%f, eps=%f, utilization=%f: " - "error_util_penalty=%f, weight=%f (not updating)", - wrr_.get(), key_.c_str(), qps, eps, utilization, - error_utilization_penalty, weight); - } - return; - } - Timestamp now = Timestamp::Now(); - // Grab the lock and update the data. - MutexLock lock(&mu_); - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, - "[WRR %p] subchannel %s: qps=%f, eps=%f, utilization=%f " - "error_util_penalty=%f : setting weight=%f weight_=%f now=%s " - "last_update_time_=%s non_empty_since_=%s", - wrr_.get(), key_.c_str(), qps, eps, utilization, - error_utilization_penalty, weight, weight_, now.ToString().c_str(), - last_update_time_.ToString().c_str(), - non_empty_since_.ToString().c_str()); - } - if (non_empty_since_ == Timestamp::InfFuture()) non_empty_since_ = now; - weight_ = weight; - last_update_time_ = now; -} - -float OldWeightedRoundRobin::AddressWeight::GetWeight( - Timestamp now, Duration weight_expiration_period, - Duration blackout_period) { - MutexLock lock(&mu_); - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, - "[WRR %p] subchannel %s: getting weight: now=%s " - "weight_expiration_period=%s blackout_period=%s " - "last_update_time_=%s non_empty_since_=%s weight_=%f", - wrr_.get(), key_.c_str(), now.ToString().c_str(), - weight_expiration_period.ToString().c_str(), - blackout_period.ToString().c_str(), - last_update_time_.ToString().c_str(), - non_empty_since_.ToString().c_str(), weight_); - } - // If the most recent update was longer ago than the expiration - // period, reset non_empty_since_ so that we apply the blackout period - // again if we start getting data again in the future, and return 0. - if (now - last_update_time_ >= weight_expiration_period) { - non_empty_since_ = Timestamp::InfFuture(); - return 0; - } - // If we don't have at least blackout_period worth of data, return 0. - if (blackout_period > Duration::Zero() && - now - non_empty_since_ < blackout_period) { - return 0; - } - // Otherwise, return the weight. - return weight_; -} - -void OldWeightedRoundRobin::AddressWeight::ResetNonEmptySince() { - MutexLock lock(&mu_); - non_empty_since_ = Timestamp::InfFuture(); -} - -// -// OldWeightedRoundRobin::Picker::SubchannelCallTracker -// - -void OldWeightedRoundRobin::Picker::SubchannelCallTracker::Finish( - FinishArgs args) { - auto* backend_metric_data = - args.backend_metric_accessor->GetBackendMetricData(); - double qps = 0; - double eps = 0; - double utilization = 0; - if (backend_metric_data != nullptr) { - qps = backend_metric_data->qps; - eps = backend_metric_data->eps; - utilization = backend_metric_data->application_utilization; - if (utilization <= 0) { - utilization = backend_metric_data->cpu_utilization; - } - } - weight_->MaybeUpdateWeight(qps, eps, utilization, error_utilization_penalty_); -} - -// -// OldWeightedRoundRobin::Picker -// - -OldWeightedRoundRobin::Picker::Picker( - RefCountedPtr wrr, - WeightedRoundRobinSubchannelList* subchannel_list) - : wrr_(std::move(wrr)), - config_(wrr_->config_), - last_picked_index_(absl::Uniform(wrr_->bit_gen_)) { - for (size_t i = 0; i < subchannel_list->num_subchannels(); ++i) { - WeightedRoundRobinSubchannelData* sd = subchannel_list->subchannel(i); - if (sd->connectivity_state() == GRPC_CHANNEL_READY) { - subchannels_.emplace_back(sd->subchannel()->Ref(), sd->weight()); - } - } - global_stats().IncrementWrrSubchannelListSize( - subchannel_list->num_subchannels()); - global_stats().IncrementWrrSubchannelReadySize(subchannels_.size()); - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, - "[WRR %p picker %p] created picker from subchannel_list=%p " - "with %" PRIuPTR " subchannels", - wrr_.get(), this, subchannel_list, subchannels_.size()); - } - BuildSchedulerAndStartTimerLocked(); -} - -OldWeightedRoundRobin::Picker::~Picker() { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p picker %p] destroying picker", wrr_.get(), this); - } -} - -void OldWeightedRoundRobin::Picker::Orphan() { - MutexLock lock(&timer_mu_); - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p picker %p] cancelling timer", wrr_.get(), this); - } - wrr_->channel_control_helper()->GetEventEngine()->Cancel(*timer_handle_); - timer_handle_.reset(); -} - -OldWeightedRoundRobin::PickResult OldWeightedRoundRobin::Picker::Pick( - PickArgs /*args*/) { - size_t index = PickIndex(); - GPR_ASSERT(index < subchannels_.size()); - auto& subchannel_info = subchannels_[index]; - // Collect per-call utilization data if needed. - std::unique_ptr subchannel_call_tracker; - if (!config_->enable_oob_load_report()) { - subchannel_call_tracker = std::make_unique( - subchannel_info.weight, config_->error_utilization_penalty()); - } - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, - "[WRR %p picker %p] returning index %" PRIuPTR ", subchannel=%p", - wrr_.get(), this, index, subchannel_info.subchannel.get()); - } - return PickResult::Complete(subchannel_info.subchannel, - std::move(subchannel_call_tracker)); -} - -size_t OldWeightedRoundRobin::Picker::PickIndex() { - // Grab a ref to the scheduler. - std::shared_ptr scheduler; - { - MutexLock lock(&scheduler_mu_); - scheduler = scheduler_; - } - // If we have a scheduler, use it to do a WRR pick. - if (scheduler != nullptr) return scheduler->Pick(); - // We don't have a scheduler (i.e., either all of the weights are 0 or - // there is only one subchannel), so fall back to RR. - return last_picked_index_.fetch_add(1) % subchannels_.size(); -} - -void OldWeightedRoundRobin::Picker::BuildSchedulerAndStartTimerLocked() { - // Build scheduler. - const Timestamp now = Timestamp::Now(); - std::vector weights; - weights.reserve(subchannels_.size()); - for (const auto& subchannel : subchannels_) { - weights.push_back(subchannel.weight->GetWeight( - now, config_->weight_expiration_period(), config_->blackout_period())); - } - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p picker %p] new weights: %s", wrr_.get(), this, - absl::StrJoin(weights, " ").c_str()); - } - auto scheduler_or = StaticStrideScheduler::Make( - weights, [this]() { return wrr_->scheduler_state_.fetch_add(1); }); - std::shared_ptr scheduler; - if (scheduler_or.has_value()) { - scheduler = - std::make_shared(std::move(*scheduler_or)); - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p picker %p] new scheduler: %p", wrr_.get(), - this, scheduler.get()); - } - } else if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p picker %p] no scheduler, falling back to RR", - wrr_.get(), this); - } - { - MutexLock lock(&scheduler_mu_); - scheduler_ = std::move(scheduler); - } - // Start timer. - timer_handle_ = wrr_->channel_control_helper()->GetEventEngine()->RunAfter( - config_->weight_update_period(), - [self = WeakRefAsSubclass(), - work_serializer = wrr_->work_serializer()]() mutable { - ApplicationCallbackExecCtx callback_exec_ctx; - ExecCtx exec_ctx; - { - MutexLock lock(&self->timer_mu_); - if (self->timer_handle_.has_value()) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p picker %p] timer fired", - self->wrr_.get(), self.get()); - } - self->BuildSchedulerAndStartTimerLocked(); - } - } - if (!IsWorkSerializerDispatchEnabled()) { - // Release the picker ref inside the WorkSerializer. - work_serializer->Run([self = std::move(self)]() {}, DEBUG_LOCATION); - return; - } - self.reset(); - }); -} - -// -// WeightedRoundRobin -// - -OldWeightedRoundRobin::OldWeightedRoundRobin(Args args) - : LoadBalancingPolicy(std::move(args)) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p] Created", this); - } -} - -OldWeightedRoundRobin::~OldWeightedRoundRobin() { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p] Destroying Round Robin policy", this); - } - GPR_ASSERT(subchannel_list_ == nullptr); - GPR_ASSERT(latest_pending_subchannel_list_ == nullptr); -} - -void OldWeightedRoundRobin::ShutdownLocked() { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p] Shutting down", this); - } - shutdown_ = true; - subchannel_list_.reset(); - latest_pending_subchannel_list_.reset(); -} - -void OldWeightedRoundRobin::ResetBackoffLocked() { - subchannel_list_->ResetBackoffLocked(); - if (latest_pending_subchannel_list_ != nullptr) { - latest_pending_subchannel_list_->ResetBackoffLocked(); - } -} - -absl::Status OldWeightedRoundRobin::UpdateLocked(UpdateArgs args) { - global_stats().IncrementWrrUpdates(); - config_ = args.config.TakeAsSubclass(); - std::shared_ptr addresses; - if (args.addresses.ok()) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p] received update", this); - } - // Weed out duplicate addresses. Also sort the addresses so that if - // the set of the addresses don't change, their indexes in the - // subchannel list don't change, since this avoids unnecessary churn - // in the picker. Note that this does not ensure that if a given - // address remains present that it will have the same index; if, - // for example, an address at the end of the list is replaced with one - // that sorts much earlier in the list, then all of the addresses in - // between those two positions will have changed indexes. - struct AddressLessThan { - bool operator()(const ServerAddress& address1, - const ServerAddress& address2) const { - const grpc_resolved_address& addr1 = address1.address(); - const grpc_resolved_address& addr2 = address2.address(); - if (addr1.len != addr2.len) return addr1.len < addr2.len; - return memcmp(addr1.addr, addr2.addr, addr1.len) < 0; - } - }; - std::set ordered_addresses; - (*args.addresses)->ForEach([&](const EndpointAddresses& endpoint) { - ordered_addresses.insert(endpoint); - }); - addresses = std::make_shared( - ServerAddressList(ordered_addresses.begin(), ordered_addresses.end())); - } else { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p] received update with address error: %s", this, - args.addresses.status().ToString().c_str()); - } - // If we already have a subchannel list, then keep using the existing - // list, but still report back that the update was not accepted. - if (subchannel_list_ != nullptr) return args.addresses.status(); - } - // Create new subchannel list, replacing the previous pending list, if any. - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace) && - latest_pending_subchannel_list_ != nullptr) { - gpr_log(GPR_INFO, "[WRR %p] replacing previous pending subchannel list %p", - this, latest_pending_subchannel_list_.get()); - } - latest_pending_subchannel_list_ = - MakeRefCounted(this, addresses.get(), - args.args); - latest_pending_subchannel_list_->StartWatchingLocked(args.args); - // If the new list is empty, immediately promote it to - // subchannel_list_ and report TRANSIENT_FAILURE. - if (latest_pending_subchannel_list_->num_subchannels() == 0) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace) && - subchannel_list_ != nullptr) { - gpr_log(GPR_INFO, "[WRR %p] replacing previous subchannel list %p", this, - subchannel_list_.get()); - } - subchannel_list_ = std::move(latest_pending_subchannel_list_); - absl::Status status = - args.addresses.ok() ? absl::UnavailableError(absl::StrCat( - "empty address list: ", args.resolution_note)) - : args.addresses.status(); - channel_control_helper()->UpdateState( - GRPC_CHANNEL_TRANSIENT_FAILURE, status, - MakeRefCounted(status)); - return status; - } - // Otherwise, if this is the initial update, immediately promote it to - // subchannel_list_. - if (subchannel_list_.get() == nullptr) { - subchannel_list_ = std::move(latest_pending_subchannel_list_); - } - return absl::OkStatus(); -} - -RefCountedPtr -OldWeightedRoundRobin::GetOrCreateWeight(const grpc_resolved_address& address) { - auto key = grpc_sockaddr_to_uri(&address); - if (!key.ok()) return nullptr; - MutexLock lock(&address_weight_map_mu_); - auto it = address_weight_map_.find(*key); - if (it != address_weight_map_.end()) { - auto weight = it->second->RefIfNonZero(); - if (weight != nullptr) return weight; - } - auto weight = MakeRefCounted( - RefAsSubclass(DEBUG_LOCATION, "AddressWeight"), - *key); - address_weight_map_.emplace(*key, weight.get()); - return weight; -} - -// -// OldWeightedRoundRobin::WeightedRoundRobinSubchannelList -// - -void OldWeightedRoundRobin::WeightedRoundRobinSubchannelList:: - UpdateStateCountersLocked(absl::optional old_state, - grpc_connectivity_state new_state) { - if (old_state.has_value()) { - GPR_ASSERT(*old_state != GRPC_CHANNEL_SHUTDOWN); - if (*old_state == GRPC_CHANNEL_READY) { - GPR_ASSERT(num_ready_ > 0); - --num_ready_; - } else if (*old_state == GRPC_CHANNEL_CONNECTING) { - GPR_ASSERT(num_connecting_ > 0); - --num_connecting_; - } else if (*old_state == GRPC_CHANNEL_TRANSIENT_FAILURE) { - GPR_ASSERT(num_transient_failure_ > 0); - --num_transient_failure_; - } - } - GPR_ASSERT(new_state != GRPC_CHANNEL_SHUTDOWN); - if (new_state == GRPC_CHANNEL_READY) { - ++num_ready_; - } else if (new_state == GRPC_CHANNEL_CONNECTING) { - ++num_connecting_; - } else if (new_state == GRPC_CHANNEL_TRANSIENT_FAILURE) { - ++num_transient_failure_; - } -} - -void OldWeightedRoundRobin::WeightedRoundRobinSubchannelList:: - MaybeUpdateAggregatedConnectivityStateLocked(absl::Status status_for_tf) { - OldWeightedRoundRobin* p = static_cast(policy()); - // If this is latest_pending_subchannel_list_, then swap it into - // subchannel_list_ in the following cases: - // - subchannel_list_ has no READY subchannels. - // - This list has at least one READY subchannel and we have seen the - // initial connectivity state notification for all subchannels. - // - All of the subchannels in this list are in TRANSIENT_FAILURE. - // (This may cause the channel to go from READY to TRANSIENT_FAILURE, - // but we're doing what the control plane told us to do.) - if (p->latest_pending_subchannel_list_.get() == this && - (p->subchannel_list_->num_ready_ == 0 || - (num_ready_ > 0 && AllSubchannelsSeenInitialState()) || - num_transient_failure_ == num_subchannels())) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - const std::string old_counters_string = - p->subchannel_list_ != nullptr ? p->subchannel_list_->CountersString() - : ""; - gpr_log( - GPR_INFO, - "[WRR %p] swapping out subchannel list %p (%s) in favor of %p (%s)", - p, p->subchannel_list_.get(), old_counters_string.c_str(), this, - CountersString().c_str()); - } - p->subchannel_list_ = std::move(p->latest_pending_subchannel_list_); - } - // Only set connectivity state if this is the current subchannel list. - if (p->subchannel_list_.get() != this) return; - // First matching rule wins: - // 1) ANY subchannel is READY => policy is READY. - // 2) ANY subchannel is CONNECTING => policy is CONNECTING. - // 3) ALL subchannels are TRANSIENT_FAILURE => policy is TRANSIENT_FAILURE. - if (num_ready_ > 0) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p] reporting READY with subchannel list %p", p, - this); - } - p->channel_control_helper()->UpdateState( - GRPC_CHANNEL_READY, absl::Status(), - MakeRefCounted(p->RefAsSubclass(), - this)); - } else if (num_connecting_ > 0) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, "[WRR %p] reporting CONNECTING with subchannel list %p", - p, this); - } - p->channel_control_helper()->UpdateState( - GRPC_CHANNEL_CONNECTING, absl::Status(), - MakeRefCounted(p->Ref(DEBUG_LOCATION, "QueuePicker"))); - } else if (num_transient_failure_ == num_subchannels()) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log( - GPR_INFO, - "[WRR %p] reporting TRANSIENT_FAILURE with subchannel list %p: %s", p, - this, status_for_tf.ToString().c_str()); - } - if (!status_for_tf.ok()) { - last_failure_ = absl::UnavailableError( - absl::StrCat("connections to all backends failing; last error: ", - status_for_tf.ToString())); - } - p->channel_control_helper()->UpdateState( - GRPC_CHANNEL_TRANSIENT_FAILURE, last_failure_, - MakeRefCounted(last_failure_)); - } -} - -// -// OldWeightedRoundRobin::WeightedRoundRobinSubchannelData::OobWatcher -// - -void OldWeightedRoundRobin::WeightedRoundRobinSubchannelData::OobWatcher:: - OnBackendMetricReport(const BackendMetricData& backend_metric_data) { - double utilization = backend_metric_data.application_utilization; - if (utilization <= 0) { - utilization = backend_metric_data.cpu_utilization; - } - weight_->MaybeUpdateWeight(backend_metric_data.qps, backend_metric_data.eps, - utilization, error_utilization_penalty_); -} - -// -// OldWeightedRoundRobin::WeightedRoundRobinSubchannelData -// - -OldWeightedRoundRobin::WeightedRoundRobinSubchannelData:: - WeightedRoundRobinSubchannelData( - SubchannelList* subchannel_list, - const ServerAddress& address, RefCountedPtr sc) - : SubchannelData(subchannel_list, address, std::move(sc)), - weight_(static_cast(subchannel_list->policy()) - ->GetOrCreateWeight(address.address())) { - // Start OOB watch if configured. - OldWeightedRoundRobin* p = - static_cast(subchannel_list->policy()); - if (p->config_->enable_oob_load_report()) { - subchannel()->AddDataWatcher(MakeOobBackendMetricWatcher( - p->config_->oob_reporting_period(), - std::make_unique(weight_, - p->config_->error_utilization_penalty()))); - } -} - -void OldWeightedRoundRobin::WeightedRoundRobinSubchannelData:: - ProcessConnectivityChangeLocked( - absl::optional old_state, - grpc_connectivity_state new_state) { - OldWeightedRoundRobin* p = - static_cast(subchannel_list()->policy()); - GPR_ASSERT(subchannel() != nullptr); - // If this is not the initial state notification and the new state is - // TRANSIENT_FAILURE or IDLE, re-resolve. - // Note that we don't want to do this on the initial state notification, - // because that would result in an endless loop of re-resolution. - if (old_state.has_value() && (new_state == GRPC_CHANNEL_TRANSIENT_FAILURE || - new_state == GRPC_CHANNEL_IDLE)) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, - "[WRR %p] Subchannel %p reported %s; requesting re-resolution", p, - subchannel(), ConnectivityStateName(new_state)); - } - p->channel_control_helper()->RequestReresolution(); - } - if (new_state == GRPC_CHANNEL_IDLE) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, - "[WRR %p] Subchannel %p reported IDLE; requesting connection", p, - subchannel()); - } - subchannel()->RequestConnection(); - } else if (new_state == GRPC_CHANNEL_READY) { - // If we transition back to READY state, restart the blackout period. - // Skip this if this is the initial notification for this - // subchannel (which happens whenever we get updated addresses and - // create a new endpoint list). Also skip it if the previous state - // was READY (which should never happen in practice, but we've seen - // at least one bug that caused this in the outlier_detection - // policy, so let's be defensive here). - // - // Note that we cannot guarantee that we will never receive - // lingering callbacks for backend metric reports from the previous - // connection after the new connection has been established, but they - // should be masked by new backend metric reports from the new - // connection by the time the blackout period ends. - if (old_state.has_value() && old_state != GRPC_CHANNEL_READY) { - weight_->ResetNonEmptySince(); - } - } - // Update logical connectivity state. - UpdateLogicalConnectivityStateLocked(new_state); - // Update the policy state. - subchannel_list()->MaybeUpdateAggregatedConnectivityStateLocked( - connectivity_status()); -} - -void OldWeightedRoundRobin::WeightedRoundRobinSubchannelData:: - UpdateLogicalConnectivityStateLocked( - grpc_connectivity_state connectivity_state) { - OldWeightedRoundRobin* p = - static_cast(subchannel_list()->policy()); - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log( - GPR_INFO, - "[WRR %p] connectivity changed for subchannel %p, subchannel_list %p " - "(index %" PRIuPTR " of %" PRIuPTR "): prev_state=%s new_state=%s", - p, subchannel(), subchannel_list(), Index(), - subchannel_list()->num_subchannels(), - (logical_connectivity_state_.has_value() - ? ConnectivityStateName(*logical_connectivity_state_) - : "N/A"), - ConnectivityStateName(connectivity_state)); - } - // Decide what state to report for aggregation purposes. - // If the last logical state was TRANSIENT_FAILURE, then ignore the - // state change unless the new state is READY. - if (logical_connectivity_state_.has_value() && - *logical_connectivity_state_ == GRPC_CHANNEL_TRANSIENT_FAILURE && - connectivity_state != GRPC_CHANNEL_READY) { - return; - } - // If the new state is IDLE, treat it as CONNECTING, since it will - // immediately transition into CONNECTING anyway. - if (connectivity_state == GRPC_CHANNEL_IDLE) { - if (GRPC_TRACE_FLAG_ENABLED(grpc_lb_wrr_trace)) { - gpr_log(GPR_INFO, - "[WRR %p] subchannel %p, subchannel_list %p (index %" PRIuPTR - " of %" PRIuPTR "): treating IDLE as CONNECTING", - p, subchannel(), subchannel_list(), Index(), - subchannel_list()->num_subchannels()); - } - connectivity_state = GRPC_CHANNEL_CONNECTING; - } - // If no change, return false. - if (logical_connectivity_state_.has_value() && - *logical_connectivity_state_ == connectivity_state) { - return; - } - // Otherwise, update counters and logical state. - subchannel_list()->UpdateStateCountersLocked(logical_connectivity_state_, - connectivity_state); - logical_connectivity_state_ = connectivity_state; -} - -// New WRR LB policy (with delegation to pick_first) +// WRR LB policy class WeightedRoundRobin : public LoadBalancingPolicy { public: explicit WeightedRoundRobin(Args args); @@ -1858,9 +1002,6 @@ class WeightedRoundRobinFactory : public LoadBalancingPolicyFactory { public: OrphanablePtr CreateLoadBalancingPolicy( LoadBalancingPolicy::Args args) const override { - if (!IsWrrDelegateToPickFirstEnabled()) { - return MakeOrphanable(std::move(args)); - } return MakeOrphanable(std::move(args)); } diff --git a/test/core/client_channel/lb_policy/outlier_detection_test.cc b/test/core/client_channel/lb_policy/outlier_detection_test.cc index 8e93de35f0f..1dbb5476037 100644 --- a/test/core/client_channel/lb_policy/outlier_detection_test.cc +++ b/test/core/client_channel/lb_policy/outlier_detection_test.cc @@ -35,7 +35,6 @@ #include #include -#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/gprpp/orphanable.h" #include "src/core/lib/gprpp/ref_counted_ptr.h" #include "src/core/lib/gprpp/time.h" @@ -224,7 +223,6 @@ TEST_F(OutlierDetectionTest, FailurePercentage) { // Advance time and run the timer callback to trigger ejection. IncrementTimeBy(Duration::Seconds(10)); gpr_log(GPR_INFO, "### ejection complete"); - if (!IsRoundRobinDelegateToPickFirstEnabled()) ExpectReresolutionRequest(); // Expect a picker update. std::vector remaining_addresses; for (const auto& addr : kAddresses) { @@ -239,7 +237,6 @@ TEST_F(OutlierDetectionTest, FailurePercentage) { } TEST_F(OutlierDetectionTest, MultipleAddressesPerEndpoint) { - if (!IsRoundRobinDelegateToPickFirstEnabled()) return; // Can't use timer duration expectation here, because the Happy // Eyeballs timer inside pick_first will use a different duration than // the timer in outlier_detection. @@ -338,7 +335,6 @@ TEST_F(OutlierDetectionTest, MultipleAddressesPerEndpoint) { } TEST_F(OutlierDetectionTest, EjectionStateResetsWhenEndpointAddressesChange) { - if (!IsRoundRobinDelegateToPickFirstEnabled()) return; // Can't use timer duration expectation here, because the Happy // Eyeballs timer inside pick_first will use a different duration than // the timer in outlier_detection. diff --git a/test/core/client_channel/lb_policy/round_robin_test.cc b/test/core/client_channel/lb_policy/round_robin_test.cc index ce49f7ef6b4..9af35ed951d 100644 --- a/test/core/client_channel/lb_policy/round_robin_test.cc +++ b/test/core/client_channel/lb_policy/round_robin_test.cc @@ -24,7 +24,6 @@ #include -#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/gprpp/orphanable.h" #include "src/core/lib/gprpp/ref_counted_ptr.h" #include "src/core/resolver/endpoint_addresses.h" @@ -70,7 +69,6 @@ TEST_F(RoundRobinTest, AddressUpdates) { } TEST_F(RoundRobinTest, MultipleAddressesPerEndpoint) { - if (!IsRoundRobinDelegateToPickFirstEnabled()) return; constexpr std::array kEndpoint1Addresses = { "ipv4:127.0.0.1:443", "ipv4:127.0.0.1:444"}; constexpr std::array kEndpoint2Addresses = { diff --git a/test/core/client_channel/lb_policy/weighted_round_robin_test.cc b/test/core/client_channel/lb_policy/weighted_round_robin_test.cc index 4572429513e..edc89309eca 100644 --- a/test/core/client_channel/lb_policy/weighted_round_robin_test.cc +++ b/test/core/client_channel/lb_policy/weighted_round_robin_test.cc @@ -40,7 +40,6 @@ #include #include -#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/gprpp/debug_location.h" #include "src/core/lib/gprpp/orphanable.h" #include "src/core/lib/gprpp/ref_counted_ptr.h" @@ -858,7 +857,6 @@ TEST_F(WeightedRoundRobinTest, ZeroErrorUtilPenalty) { } TEST_F(WeightedRoundRobinTest, MultipleAddressesPerEndpoint) { - if (!IsWrrDelegateToPickFirstEnabled()) return; // Can't use timer duration expectation here, because the Happy // Eyeballs timer inside pick_first will use a different duration than // the timer in WRR. @@ -1066,7 +1064,6 @@ TEST_F(WeightedRoundRobinTest, MetricDefinitionEndpointWeights) { } TEST_F(WeightedRoundRobinTest, MetricValues) { - if (!IsWrrDelegateToPickFirstEnabled()) return; const auto kRrFallback = GlobalInstrumentsRegistryTestPeer::FindUInt64CounterHandleByName( "grpc.lb.wrr.rr_fallback") diff --git a/test/core/client_channel/lb_policy/xds_override_host_test.cc b/test/core/client_channel/lb_policy/xds_override_host_test.cc index d1ddeae05c7..e6c74b7b2f1 100644 --- a/test/core/client_channel/lb_policy/xds_override_host_test.cc +++ b/test/core/client_channel/lb_policy/xds_override_host_test.cc @@ -40,7 +40,6 @@ #include "src/core/ext/filters/stateful_session/stateful_session_filter.h" #include "src/core/ext/xds/xds_health_status.h" #include "src/core/lib/channel/channel_args.h" -#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/gprpp/debug_location.h" #include "src/core/lib/gprpp/ref_counted_ptr.h" #include "src/core/lib/json/json.h" @@ -511,7 +510,6 @@ TEST_F(XdsOverrideHostTest, OverrideHostStatus) { } TEST_F(XdsOverrideHostTest, MultipleAddressesPerEndpoint) { - if (!IsRoundRobinDelegateToPickFirstEnabled()) return; constexpr std::array kEndpoint1Addresses = { "ipv4:127.0.0.1:443", "ipv4:127.0.0.1:444"}; constexpr std::array kEndpoint2Addresses = { diff --git a/test/cpp/end2end/client_lb_end2end_test.cc b/test/cpp/end2end/client_lb_end2end_test.cc index 11bc2492cdd..d5ee65ad9d0 100644 --- a/test/cpp/end2end/client_lb_end2end_test.cc +++ b/test/cpp/end2end/client_lb_end2end_test.cc @@ -56,7 +56,6 @@ #include "src/core/lib/backoff/backoff.h" #include "src/core/lib/channel/channel_args.h" #include "src/core/lib/config/config_vars.h" -#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/gprpp/crash.h" #include "src/core/lib/gprpp/debug_location.h" #include "src/core/lib/gprpp/env.h" @@ -1982,13 +1981,8 @@ TEST_F(RoundRobinTest, HealthChecking) { EXPECT_THAT( status.error_message(), ::testing::MatchesRegex( - grpc_core::IsRoundRobinDelegateToPickFirstEnabled() - ? "connections to all backends failing; last error: " - "(ipv6:%5B::1%5D|ipv4:127.0.0.1):[0-9]+: " - "backend unhealthy" - : "connections to all backends failing; last error: " - "UNAVAILABLE: (ipv6:%5B::1%5D|ipv4:127.0.0.1):[0-9]+: " - "backend unhealthy")); + "connections to all backends failing; last error: " + "(ipv6:%5B::1%5D|ipv4:127.0.0.1):[0-9]+: backend unhealthy")); return false; }); // Clean up. @@ -2048,13 +2042,8 @@ TEST_F(RoundRobinTest, WithHealthCheckingInhibitPerChannel) { EXPECT_FALSE(WaitForChannelReady(channel1.get(), 1)); CheckRpcSendFailure( DEBUG_LOCATION, stub1, StatusCode::UNAVAILABLE, - grpc_core::IsRoundRobinDelegateToPickFirstEnabled() - ? "connections to all backends failing; last error: " - "(ipv6:%5B::1%5D|ipv4:127.0.0.1):[0-9]+: " - "backend unhealthy" - : "connections to all backends failing; last error: " - "UNAVAILABLE: (ipv6:%5B::1%5D|ipv4:127.0.0.1):[0-9]+: " - "backend unhealthy"); + "connections to all backends failing; last error: " + "(ipv6:%5B::1%5D|ipv4:127.0.0.1):[0-9]+: backend unhealthy"); // Second channel should be READY. EXPECT_TRUE(WaitForChannelReady(channel2.get(), 1)); CheckRpcSendOk(DEBUG_LOCATION, stub2); @@ -2099,13 +2088,8 @@ TEST_F(RoundRobinTest, HealthCheckingServiceNamePerChannel) { EXPECT_FALSE(WaitForChannelReady(channel1.get(), 1)); CheckRpcSendFailure( DEBUG_LOCATION, stub1, StatusCode::UNAVAILABLE, - grpc_core::IsRoundRobinDelegateToPickFirstEnabled() - ? "connections to all backends failing; last error: " - "(ipv6:%5B::1%5D|ipv4:127.0.0.1):[0-9]+: " - "backend unhealthy" - : "connections to all backends failing; last error: " - "UNAVAILABLE: (ipv6:%5B::1%5D|ipv4:127.0.0.1):[0-9]+: " - "backend unhealthy"); + "connections to all backends failing; last error: " + "(ipv6:%5B::1%5D|ipv4:127.0.0.1):[0-9]+: backend unhealthy"); // Second channel should be READY. EXPECT_TRUE(WaitForChannelReady(channel2.get(), 1)); CheckRpcSendOk(DEBUG_LOCATION, stub2); diff --git a/test/cpp/end2end/xds/xds_cluster_end2end_test.cc b/test/cpp/end2end/xds/xds_cluster_end2end_test.cc index 7322f06c9f6..2772a2370cf 100644 --- a/test/cpp/end2end/xds/xds_cluster_end2end_test.cc +++ b/test/cpp/end2end/xds/xds_cluster_end2end_test.cc @@ -27,7 +27,6 @@ #include "src/core/lib/address_utils/sockaddr_utils.h" #include "src/core/lib/channel/call_tracer.h" #include "src/core/lib/config/config_vars.h" -#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/surface/call.h" #include "src/proto/grpc/testing/xds/v3/orca_load_report.pb.h" #include "test/core/util/fake_stats_plugin.h" @@ -482,7 +481,6 @@ TEST_P(EdsTest, Vanilla) { } TEST_P(EdsTest, MultipleAddressesPerEndpoint) { - if (!grpc_core::IsRoundRobinDelegateToPickFirstEnabled()) return; grpc_core::testing::ScopedExperimentalEnvVar env( "GRPC_EXPERIMENTAL_XDS_DUALSTACK_ENDPOINTS"); const size_t kNumRpcsPerAddress = 10; diff --git a/test/cpp/end2end/xds/xds_override_host_end2end_test.cc b/test/cpp/end2end/xds/xds_override_host_end2end_test.cc index 65d5c43f9bb..fbd3793b206 100644 --- a/test/cpp/end2end/xds/xds_override_host_end2end_test.cc +++ b/test/cpp/end2end/xds/xds_override_host_end2end_test.cc @@ -23,7 +23,6 @@ #include "absl/strings/str_split.h" #include "src/core/lib/config/config_vars.h" -#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/gprpp/time.h" #include "src/proto/grpc/testing/xds/v3/stateful_session.pb.h" #include "src/proto/grpc/testing/xds/v3/stateful_session_cookie.pb.h" @@ -703,7 +702,6 @@ TEST_P(OverrideHostTest, TTLSetsMaxAge) { } TEST_P(OverrideHostTest, MultipleAddressesPerEndpoint) { - if (!grpc_core::IsRoundRobinDelegateToPickFirstEnabled()) return; grpc_core::testing::ScopedExperimentalEnvVar env( "GRPC_EXPERIMENTAL_XDS_DUALSTACK_ENDPOINTS"); // Create 3 backends, but leave backend 0 unstarted. diff --git a/test/cpp/end2end/xds/xds_wrr_end2end_test.cc b/test/cpp/end2end/xds/xds_wrr_end2end_test.cc index a89a3021e73..88f7513b532 100644 --- a/test/cpp/end2end/xds/xds_wrr_end2end_test.cc +++ b/test/cpp/end2end/xds/xds_wrr_end2end_test.cc @@ -27,7 +27,6 @@ #include "src/core/client_channel/backup_poller.h" #include "src/core/lib/config/config_vars.h" -#include "src/core/lib/experiments/experiments.h" #include "src/proto/grpc/testing/xds/v3/client_side_weighted_round_robin.grpc.pb.h" #include "src/proto/grpc/testing/xds/v3/wrr_locality.grpc.pb.h" #include "test/core/util/fake_stats_plugin.h" @@ -102,7 +101,6 @@ TEST_P(WrrTest, Basic) { } TEST_P(WrrTest, MetricsHaveLocalityLabel) { - if (!grpc_core::IsWrrDelegateToPickFirstEnabled()) return; const auto kEndpointWeights = grpc_core::GlobalInstrumentsRegistryTestPeer:: FindDoubleHistogramHandleByName("grpc.lb.wrr.endpoint_weights") diff --git a/tools/doxygen/Doxyfile.c++.internal b/tools/doxygen/Doxyfile.c++.internal index 35da52b9935..ba4cbaec282 100644 --- a/tools/doxygen/Doxyfile.c++.internal +++ b/tools/doxygen/Doxyfile.c++.internal @@ -2897,7 +2897,6 @@ src/core/load_balancing/rls/rls.cc \ src/core/load_balancing/rls/rls.h \ src/core/load_balancing/round_robin/round_robin.cc \ src/core/load_balancing/subchannel_interface.h \ -src/core/load_balancing/subchannel_list.h \ src/core/load_balancing/weighted_round_robin/static_stride_scheduler.cc \ src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h \ src/core/load_balancing/weighted_round_robin/weighted_round_robin.cc \ diff --git a/tools/doxygen/Doxyfile.core.internal b/tools/doxygen/Doxyfile.core.internal index 1f4c698de8b..8796adfd4c6 100644 --- a/tools/doxygen/Doxyfile.core.internal +++ b/tools/doxygen/Doxyfile.core.internal @@ -2674,7 +2674,6 @@ src/core/load_balancing/rls/rls.cc \ src/core/load_balancing/rls/rls.h \ src/core/load_balancing/round_robin/round_robin.cc \ src/core/load_balancing/subchannel_interface.h \ -src/core/load_balancing/subchannel_list.h \ src/core/load_balancing/weighted_round_robin/static_stride_scheduler.cc \ src/core/load_balancing/weighted_round_robin/static_stride_scheduler.h \ src/core/load_balancing/weighted_round_robin/weighted_round_robin.cc \