From a7bf07e86a5a5e37ab0889b72055c4d3bb0b998e Mon Sep 17 00:00:00 2001 From: Yijie Ma Date: Fri, 21 Jul 2023 13:24:16 -0700 Subject: [PATCH] [EventEngine] PosixEventEngine DNS Resolver (#32701) This PR implements a c-ares based DNS resolver for EventEngine with the reference from the original [grpc_ares_wrapper.h](../blob/master/src/core/ext/filters/client_channel/resolver/dns/c_ares/grpc_ares_wrapper.h). The PosixEventEngine DNSResolver is implemented on top of that. Tests which use the client channel resolver API ([resolver.h](../blob/master/src/core/lib/resolver/resolver.h#L54)) are ported, namely the [resolver_component_test.cc](../blob/master/test/cpp/naming/resolver_component_test.cc) and the [cancel_ares_query_test.cc](../blob/master/test/cpp/naming/cancel_ares_query_test.cc). The WindowsEventEngine DNSResolver will use the same EventEngine's grpc_ares_wrapper and will be worked on next. The [resolve_address_test.cc](https://github.com/grpc/grpc/blob/master/test/core/iomgr/resolve_address_test.cc) which uses the iomgr [DNSResolver](../blob/master/src/core/lib/iomgr/resolve_address.h#L44) API has been ported to EventEngine's dns_test.cc. That leaves only 2 tests which use iomgr's API, notably the [dns_resolver_cooldown_test.cc](../blob/master/test/core/client_channel/resolvers/dns_resolver_cooldown_test.cc) and the [goaway_server_test.cc](../blob/master/test/core/end2end/goaway_server_test.cc) which probably need to be restructured to use EventEngine DNSResolver (for one thing they override the original grpc_ares_wrapper's free functions). I will try to tackle these in the next step. --- CMakeLists.txt | 11 +- Makefile | 2 + Package.swift | 4 + bazel/experiments.bzl | 6 + build_autogenerated.yaml | 24 +- config.m4 | 1 + config.w32 | 1 + gRPC-C++.podspec | 6 + gRPC-Core.podspec | 7 + grpc.gemspec | 4 + grpc.gyp | 3 + include/grpc/event_engine/event_engine.h | 1 - package.xml | 4 + src/core/BUILD | 46 ++ .../resolver/dns/c_ares/grpc_ares_wrapper.cc | 6 +- .../resolver/dns/c_ares/grpc_ares_wrapper.h | 2 +- .../event_engine_client_channel_resolver.cc | 53 +- .../resolver/polling_resolver.cc | 5 +- src/core/lib/event_engine/ares_resolver.cc | 705 ++++++++++++++++++ src/core/lib/event_engine/ares_resolver.h | 147 ++++ src/core/lib/event_engine/grpc_polled_fd.h | 73 ++ .../posix_engine/grpc_polled_fd_posix.h | 112 +++ .../event_engine/posix_engine/posix_engine.cc | 44 +- .../event_engine/posix_engine/posix_engine.h | 13 +- src/core/lib/experiments/experiments.yaml | 2 +- src/python/grpcio/grpc_core_dependencies.py | 1 + templates/CMakeLists.txt.template | 4 + .../resolver_component_tests_defs.include | 5 +- test/core/event_engine/test_suite/BUILD | 1 + test/core/event_engine/test_suite/tests/BUILD | 18 + .../event_engine/test_suite/tests/dns_test.cc | 520 ++++++++++++- .../tests/dns_test_record_groups.yaml | 21 + test/core/http/httpcli_test.cc | 7 +- test/core/iomgr/resolve_address_test.cc | 6 +- test/cpp/naming/BUILD | 1 + test/cpp/naming/cancel_ares_query_test.cc | 14 +- .../generate_resolver_component_tests.bzl | 9 +- test/cpp/naming/resolver_component_test.cc | 39 +- .../naming/resolver_component_tests_runner.py | 5 +- test/cpp/naming/utils/BUILD | 6 + test/cpp/naming/utils/dns_resolver.py | 2 +- test/cpp/naming/utils/dns_server.py | 2 +- test/cpp/naming/utils/health_check.py | 121 +++ .../run_dns_server_for_lb_interop_tests.py | 2 +- test/cpp/naming/utils/tcp_connect.py | 2 +- test/cpp/util/BUILD | 16 + test/cpp/util/get_grpc_test_runfile_dir.cc | 29 + test/cpp/util/get_grpc_test_runfile_dir.h | 32 + tools/doxygen/Doxyfile.c++.internal | 4 + tools/doxygen/Doxyfile.core.internal | 4 + 50 files changed, 2097 insertions(+), 56 deletions(-) create mode 100644 src/core/lib/event_engine/ares_resolver.cc create mode 100644 src/core/lib/event_engine/ares_resolver.h create mode 100644 src/core/lib/event_engine/grpc_polled_fd.h create mode 100644 src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h create mode 100644 test/core/event_engine/test_suite/tests/dns_test_record_groups.yaml create mode 100644 test/cpp/naming/utils/health_check.py create mode 100644 test/cpp/util/get_grpc_test_runfile_dir.cc create mode 100644 test/cpp/util/get_grpc_test_runfile_dir.h diff --git a/CMakeLists.txt b/CMakeLists.txt index ee3538afcc6..e4f151dad52 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -2150,6 +2150,7 @@ add_library(grpc src/core/lib/debug/stats.cc src/core/lib/debug/stats_data.cc src/core/lib/debug/trace.cc + src/core/lib/event_engine/ares_resolver.cc src/core/lib/event_engine/cf_engine/cf_engine.cc src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc src/core/lib/event_engine/channel_args_endpoint_config.cc @@ -2855,6 +2856,7 @@ add_library(grpc_unsecure src/core/lib/debug/stats.cc src/core/lib/debug/stats_data.cc src/core/lib/debug/trace.cc + src/core/lib/event_engine/ares_resolver.cc src/core/lib/event_engine/cf_engine/cf_engine.cc src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc src/core/lib/event_engine/channel_args_endpoint_config.cc @@ -4391,6 +4393,7 @@ add_library(grpc_authorization_provider src/core/lib/debug/stats.cc src/core/lib/debug/stats_data.cc src/core/lib/debug/trace.cc + src/core/lib/event_engine/ares_resolver.cc src/core/lib/event_engine/cf_engine/cf_engine.cc src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc src/core/lib/event_engine/channel_args_endpoint_config.cc @@ -4654,6 +4657,7 @@ target_link_libraries(grpc_authorization_provider ${_gRPC_BASELIB_LIBRARIES} ${_gRPC_PROTOBUF_LIBRARIES} ${_gRPC_ZLIB_LIBRARIES} + ${_gRPC_CARES_LIBRARIES} ${_gRPC_RE2_LIBRARIES} ${_gRPC_ALLTARGETS_LIBRARIES} absl::cleanup @@ -12385,6 +12389,7 @@ add_executable(frame_test src/core/lib/debug/stats.cc src/core/lib/debug/stats_data.cc src/core/lib/debug/trace.cc + src/core/lib/event_engine/ares_resolver.cc src/core/lib/event_engine/cf_engine/cf_engine.cc src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc src/core/lib/event_engine/channel_args_endpoint_config.cc @@ -12608,6 +12613,7 @@ target_link_libraries(frame_test ${_gRPC_BASELIB_LIBRARIES} ${_gRPC_PROTOBUF_LIBRARIES} ${_gRPC_ZLIB_LIBRARIES} + ${_gRPC_CARES_LIBRARIES} ${_gRPC_ALLTARGETS_LIBRARIES} absl::cleanup absl::flat_hash_map @@ -18583,8 +18589,11 @@ if(_gRPC_PLATFORM_LINUX OR _gRPC_PLATFORM_POSIX) test/core/event_engine/test_suite/posix/oracle_event_engine_posix.cc test/core/event_engine/test_suite/posix_event_engine_test.cc test/core/event_engine/test_suite/tests/client_test.cc + test/core/event_engine/test_suite/tests/dns_test.cc test/core/event_engine/test_suite/tests/server_test.cc test/core/event_engine/test_suite/tests/timer_test.cc + test/core/util/fake_udp_and_tcp_server.cc + test/cpp/util/get_grpc_test_runfile_dir.cc third_party/googletest/googletest/src/gtest-all.cc third_party/googletest/googlemock/src/gmock-all.cc ) @@ -18614,7 +18623,7 @@ if(_gRPC_PLATFORM_LINUX OR _gRPC_PLATFORM_POSIX) ${_gRPC_ZLIB_LIBRARIES} ${_gRPC_ALLTARGETS_LIBRARIES} grpc_unsecure - grpc_test_util + grpc++_test_util ) diff --git a/Makefile b/Makefile index 5094236d47e..a9c5977295f 100644 --- a/Makefile +++ b/Makefile @@ -1435,6 +1435,7 @@ LIBGRPC_SRC = \ src/core/lib/debug/stats.cc \ src/core/lib/debug/stats_data.cc \ src/core/lib/debug/trace.cc \ + src/core/lib/event_engine/ares_resolver.cc \ src/core/lib/event_engine/cf_engine/cf_engine.cc \ src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc \ src/core/lib/event_engine/channel_args_endpoint_config.cc \ @@ -1993,6 +1994,7 @@ LIBGRPC_UNSECURE_SRC = \ src/core/lib/debug/stats.cc \ src/core/lib/debug/stats_data.cc \ src/core/lib/debug/trace.cc \ + src/core/lib/event_engine/ares_resolver.cc \ src/core/lib/event_engine/cf_engine/cf_engine.cc \ src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc \ src/core/lib/event_engine/channel_args_endpoint_config.cc \ diff --git a/Package.swift b/Package.swift index 9c40122a673..a3d31ce6d56 100644 --- a/Package.swift +++ b/Package.swift @@ -1057,6 +1057,8 @@ let package = Package( "src/core/lib/debug/stats_data.h", "src/core/lib/debug/trace.cc", "src/core/lib/debug/trace.h", + "src/core/lib/event_engine/ares_resolver.cc", + "src/core/lib/event_engine/ares_resolver.h", "src/core/lib/event_engine/cf_engine/cf_engine.cc", "src/core/lib/event_engine/cf_engine/cf_engine.h", "src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc", @@ -1072,6 +1074,7 @@ let package = Package( "src/core/lib/event_engine/event_engine.cc", "src/core/lib/event_engine/forkable.cc", "src/core/lib/event_engine/forkable.h", + "src/core/lib/event_engine/grpc_polled_fd.h", "src/core/lib/event_engine/handle_containers.h", "src/core/lib/event_engine/memory_allocator.cc", "src/core/lib/event_engine/memory_allocator_factory.h", @@ -1084,6 +1087,7 @@ let package = Package( "src/core/lib/event_engine/posix_engine/event_poller.h", "src/core/lib/event_engine/posix_engine/event_poller_posix_default.cc", "src/core/lib/event_engine/posix_engine/event_poller_posix_default.h", + "src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h", "src/core/lib/event_engine/posix_engine/internal_errqueue.cc", "src/core/lib/event_engine/posix_engine/internal_errqueue.h", "src/core/lib/event_engine/posix_engine/lockfree_event.cc", diff --git a/bazel/experiments.bzl b/bazel/experiments.bzl index a5c78363eb2..e04ba5720ae 100644 --- a/bazel/experiments.bzl +++ b/bazel/experiments.bzl @@ -20,6 +20,9 @@ EXPERIMENTS = { "dbg": { }, "off": { + "cancel_ares_query_test": [ + "event_engine_dns", + ], "census_test": [ "transport_supplies_client_latency", ], @@ -55,6 +58,9 @@ EXPERIMENTS = { "logging_test": [ "promise_based_server_call", ], + "resolver_component_tests_runner_invoker": [ + "event_engine_dns", + ], "resource_quota_test": [ "free_large_allocator", "memory_pressure_controller", diff --git a/build_autogenerated.yaml b/build_autogenerated.yaml index a4ae130713a..891c136a664 100644 --- a/build_autogenerated.yaml +++ b/build_autogenerated.yaml @@ -680,6 +680,7 @@ libs: - src/core/lib/debug/stats.h - src/core/lib/debug/stats_data.h - src/core/lib/debug/trace.h + - src/core/lib/event_engine/ares_resolver.h - src/core/lib/event_engine/cf_engine/cf_engine.h - src/core/lib/event_engine/cf_engine/cfstream_endpoint.h - src/core/lib/event_engine/cf_engine/cftype_unique_ref.h @@ -688,6 +689,7 @@ libs: - src/core/lib/event_engine/default_event_engine.h - src/core/lib/event_engine/default_event_engine_factory.h - src/core/lib/event_engine/forkable.h + - src/core/lib/event_engine/grpc_polled_fd.h - src/core/lib/event_engine/handle_containers.h - src/core/lib/event_engine/memory_allocator_factory.h - src/core/lib/event_engine/poller.h @@ -696,6 +698,7 @@ libs: - src/core/lib/event_engine/posix_engine/ev_poll_posix.h - src/core/lib/event_engine/posix_engine/event_poller.h - src/core/lib/event_engine/posix_engine/event_poller_posix_default.h + - src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h - src/core/lib/event_engine/posix_engine/internal_errqueue.h - src/core/lib/event_engine/posix_engine/lockfree_event.h - src/core/lib/event_engine/posix_engine/posix_endpoint.h @@ -1492,6 +1495,7 @@ libs: - src/core/lib/debug/stats.cc - src/core/lib/debug/stats_data.cc - src/core/lib/debug/trace.cc + - src/core/lib/event_engine/ares_resolver.cc - src/core/lib/event_engine/cf_engine/cf_engine.cc - src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc - src/core/lib/event_engine/channel_args_endpoint_config.cc @@ -2070,6 +2074,7 @@ libs: - src/core/lib/debug/stats.h - src/core/lib/debug/stats_data.h - src/core/lib/debug/trace.h + - src/core/lib/event_engine/ares_resolver.h - src/core/lib/event_engine/cf_engine/cf_engine.h - src/core/lib/event_engine/cf_engine/cfstream_endpoint.h - src/core/lib/event_engine/cf_engine/cftype_unique_ref.h @@ -2078,6 +2083,7 @@ libs: - src/core/lib/event_engine/default_event_engine.h - src/core/lib/event_engine/default_event_engine_factory.h - src/core/lib/event_engine/forkable.h + - src/core/lib/event_engine/grpc_polled_fd.h - src/core/lib/event_engine/handle_containers.h - src/core/lib/event_engine/memory_allocator_factory.h - src/core/lib/event_engine/poller.h @@ -2086,6 +2092,7 @@ libs: - src/core/lib/event_engine/posix_engine/ev_poll_posix.h - src/core/lib/event_engine/posix_engine/event_poller.h - src/core/lib/event_engine/posix_engine/event_poller_posix_default.h + - src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h - src/core/lib/event_engine/posix_engine/internal_errqueue.h - src/core/lib/event_engine/posix_engine/lockfree_event.h - src/core/lib/event_engine/posix_engine/posix_endpoint.h @@ -2489,6 +2496,7 @@ libs: - src/core/lib/debug/stats.cc - src/core/lib/debug/stats_data.cc - src/core/lib/debug/trace.cc + - src/core/lib/event_engine/ares_resolver.cc - src/core/lib/event_engine/cf_engine/cf_engine.cc - src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc - src/core/lib/event_engine/channel_args_endpoint_config.cc @@ -3573,6 +3581,7 @@ libs: - src/core/lib/debug/stats.h - src/core/lib/debug/stats_data.h - src/core/lib/debug/trace.h + - src/core/lib/event_engine/ares_resolver.h - src/core/lib/event_engine/cf_engine/cf_engine.h - src/core/lib/event_engine/cf_engine/cfstream_endpoint.h - src/core/lib/event_engine/cf_engine/cftype_unique_ref.h @@ -3581,6 +3590,7 @@ libs: - src/core/lib/event_engine/default_event_engine.h - src/core/lib/event_engine/default_event_engine_factory.h - src/core/lib/event_engine/forkable.h + - src/core/lib/event_engine/grpc_polled_fd.h - src/core/lib/event_engine/handle_containers.h - src/core/lib/event_engine/memory_allocator_factory.h - src/core/lib/event_engine/poller.h @@ -3589,6 +3599,7 @@ libs: - src/core/lib/event_engine/posix_engine/ev_poll_posix.h - src/core/lib/event_engine/posix_engine/event_poller.h - src/core/lib/event_engine/posix_engine/event_poller_posix_default.h + - src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h - src/core/lib/event_engine/posix_engine/internal_errqueue.h - src/core/lib/event_engine/posix_engine/lockfree_event.h - src/core/lib/event_engine/posix_engine/posix_endpoint.h @@ -3873,6 +3884,7 @@ libs: - src/core/lib/debug/stats.cc - src/core/lib/debug/stats_data.cc - src/core/lib/debug/trace.cc + - src/core/lib/event_engine/ares_resolver.cc - src/core/lib/event_engine/cf_engine/cf_engine.cc - src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc - src/core/lib/event_engine/channel_args_endpoint_config.cc @@ -8071,6 +8083,7 @@ targets: - src/core/lib/debug/stats.h - src/core/lib/debug/stats_data.h - src/core/lib/debug/trace.h + - src/core/lib/event_engine/ares_resolver.h - src/core/lib/event_engine/cf_engine/cf_engine.h - src/core/lib/event_engine/cf_engine/cfstream_endpoint.h - src/core/lib/event_engine/cf_engine/cftype_unique_ref.h @@ -8079,6 +8092,7 @@ targets: - src/core/lib/event_engine/default_event_engine.h - src/core/lib/event_engine/default_event_engine_factory.h - src/core/lib/event_engine/forkable.h + - src/core/lib/event_engine/grpc_polled_fd.h - src/core/lib/event_engine/handle_containers.h - src/core/lib/event_engine/memory_allocator_factory.h - src/core/lib/event_engine/poller.h @@ -8087,6 +8101,7 @@ targets: - src/core/lib/event_engine/posix_engine/ev_poll_posix.h - src/core/lib/event_engine/posix_engine/event_poller.h - src/core/lib/event_engine/posix_engine/event_poller_posix_default.h + - src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h - src/core/lib/event_engine/posix_engine/internal_errqueue.h - src/core/lib/event_engine/posix_engine/lockfree_event.h - src/core/lib/event_engine/posix_engine/posix_endpoint.h @@ -8353,6 +8368,7 @@ targets: - src/core/lib/debug/stats.cc - src/core/lib/debug/stats_data.cc - src/core/lib/debug/trace.cc + - src/core/lib/event_engine/ares_resolver.cc - src/core/lib/event_engine/cf_engine/cf_engine.cc - src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc - src/core/lib/event_engine/channel_args_endpoint_config.cc @@ -11427,19 +11443,25 @@ targets: - test/core/event_engine/test_suite/event_engine_test_framework.h - test/core/event_engine/test_suite/posix/oracle_event_engine_posix.h - test/core/event_engine/test_suite/tests/client_test.h + - test/core/event_engine/test_suite/tests/dns_test.h - test/core/event_engine/test_suite/tests/server_test.h - test/core/event_engine/test_suite/tests/timer_test.h + - test/core/util/fake_udp_and_tcp_server.h + - test/cpp/util/get_grpc_test_runfile_dir.h src: - test/core/event_engine/event_engine_test_utils.cc - test/core/event_engine/test_suite/event_engine_test_framework.cc - test/core/event_engine/test_suite/posix/oracle_event_engine_posix.cc - test/core/event_engine/test_suite/posix_event_engine_test.cc - test/core/event_engine/test_suite/tests/client_test.cc + - test/core/event_engine/test_suite/tests/dns_test.cc - test/core/event_engine/test_suite/tests/server_test.cc - test/core/event_engine/test_suite/tests/timer_test.cc + - test/core/util/fake_udp_and_tcp_server.cc + - test/cpp/util/get_grpc_test_runfile_dir.cc deps: - grpc_unsecure - - grpc_test_util + - grpc++_test_util platforms: - linux - posix diff --git a/config.m4 b/config.m4 index e7ddb42f6c3..44e239db30e 100644 --- a/config.m4 +++ b/config.m4 @@ -517,6 +517,7 @@ if test "$PHP_GRPC" != "no"; then src/core/lib/debug/stats.cc \ src/core/lib/debug/stats_data.cc \ src/core/lib/debug/trace.cc \ + src/core/lib/event_engine/ares_resolver.cc \ src/core/lib/event_engine/cf_engine/cf_engine.cc \ src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc \ src/core/lib/event_engine/channel_args_endpoint_config.cc \ diff --git a/config.w32 b/config.w32 index 64851c79e43..da44d20aedd 100644 --- a/config.w32 +++ b/config.w32 @@ -482,6 +482,7 @@ if (PHP_GRPC != "no") { "src\\core\\lib\\debug\\stats.cc " + "src\\core\\lib\\debug\\stats_data.cc " + "src\\core\\lib\\debug\\trace.cc " + + "src\\core\\lib\\event_engine\\ares_resolver.cc " + "src\\core\\lib\\event_engine\\cf_engine\\cf_engine.cc " + "src\\core\\lib\\event_engine\\cf_engine\\cfstream_endpoint.cc " + "src\\core\\lib\\event_engine\\channel_args_endpoint_config.cc " + diff --git a/gRPC-C++.podspec b/gRPC-C++.podspec index ffc0b135179..94d9b33784e 100644 --- a/gRPC-C++.podspec +++ b/gRPC-C++.podspec @@ -752,6 +752,7 @@ Pod::Spec.new do |s| 'src/core/lib/debug/stats.h', 'src/core/lib/debug/stats_data.h', 'src/core/lib/debug/trace.h', + 'src/core/lib/event_engine/ares_resolver.h', 'src/core/lib/event_engine/cf_engine/cf_engine.h', 'src/core/lib/event_engine/cf_engine/cfstream_endpoint.h', 'src/core/lib/event_engine/cf_engine/cftype_unique_ref.h', @@ -760,6 +761,7 @@ Pod::Spec.new do |s| 'src/core/lib/event_engine/default_event_engine.h', 'src/core/lib/event_engine/default_event_engine_factory.h', 'src/core/lib/event_engine/forkable.h', + 'src/core/lib/event_engine/grpc_polled_fd.h', 'src/core/lib/event_engine/handle_containers.h', 'src/core/lib/event_engine/memory_allocator_factory.h', 'src/core/lib/event_engine/poller.h', @@ -768,6 +770,7 @@ Pod::Spec.new do |s| 'src/core/lib/event_engine/posix_engine/ev_poll_posix.h', 'src/core/lib/event_engine/posix_engine/event_poller.h', 'src/core/lib/event_engine/posix_engine/event_poller_posix_default.h', + 'src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h', 'src/core/lib/event_engine/posix_engine/internal_errqueue.h', 'src/core/lib/event_engine/posix_engine/lockfree_event.h', 'src/core/lib/event_engine/posix_engine/posix_endpoint.h', @@ -1798,6 +1801,7 @@ Pod::Spec.new do |s| 'src/core/lib/debug/stats.h', 'src/core/lib/debug/stats_data.h', 'src/core/lib/debug/trace.h', + 'src/core/lib/event_engine/ares_resolver.h', 'src/core/lib/event_engine/cf_engine/cf_engine.h', 'src/core/lib/event_engine/cf_engine/cfstream_endpoint.h', 'src/core/lib/event_engine/cf_engine/cftype_unique_ref.h', @@ -1806,6 +1810,7 @@ Pod::Spec.new do |s| 'src/core/lib/event_engine/default_event_engine.h', 'src/core/lib/event_engine/default_event_engine_factory.h', 'src/core/lib/event_engine/forkable.h', + 'src/core/lib/event_engine/grpc_polled_fd.h', 'src/core/lib/event_engine/handle_containers.h', 'src/core/lib/event_engine/memory_allocator_factory.h', 'src/core/lib/event_engine/poller.h', @@ -1814,6 +1819,7 @@ Pod::Spec.new do |s| 'src/core/lib/event_engine/posix_engine/ev_poll_posix.h', 'src/core/lib/event_engine/posix_engine/event_poller.h', 'src/core/lib/event_engine/posix_engine/event_poller_posix_default.h', + 'src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h', 'src/core/lib/event_engine/posix_engine/internal_errqueue.h', 'src/core/lib/event_engine/posix_engine/lockfree_event.h', 'src/core/lib/event_engine/posix_engine/posix_endpoint.h', diff --git a/gRPC-Core.podspec b/gRPC-Core.podspec index 36f40caf386..9c1956cbeb0 100644 --- a/gRPC-Core.podspec +++ b/gRPC-Core.podspec @@ -1158,6 +1158,8 @@ Pod::Spec.new do |s| 'src/core/lib/debug/stats_data.h', 'src/core/lib/debug/trace.cc', 'src/core/lib/debug/trace.h', + 'src/core/lib/event_engine/ares_resolver.cc', + 'src/core/lib/event_engine/ares_resolver.h', 'src/core/lib/event_engine/cf_engine/cf_engine.cc', 'src/core/lib/event_engine/cf_engine/cf_engine.h', 'src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc', @@ -1173,6 +1175,7 @@ Pod::Spec.new do |s| 'src/core/lib/event_engine/event_engine.cc', 'src/core/lib/event_engine/forkable.cc', 'src/core/lib/event_engine/forkable.h', + 'src/core/lib/event_engine/grpc_polled_fd.h', 'src/core/lib/event_engine/handle_containers.h', 'src/core/lib/event_engine/memory_allocator.cc', 'src/core/lib/event_engine/memory_allocator_factory.h', @@ -1185,6 +1188,7 @@ Pod::Spec.new do |s| 'src/core/lib/event_engine/posix_engine/event_poller.h', 'src/core/lib/event_engine/posix_engine/event_poller_posix_default.cc', 'src/core/lib/event_engine/posix_engine/event_poller_posix_default.h', + 'src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h', 'src/core/lib/event_engine/posix_engine/internal_errqueue.cc', 'src/core/lib/event_engine/posix_engine/internal_errqueue.h', 'src/core/lib/event_engine/posix_engine/lockfree_event.cc', @@ -2529,6 +2533,7 @@ Pod::Spec.new do |s| 'src/core/lib/debug/stats.h', 'src/core/lib/debug/stats_data.h', 'src/core/lib/debug/trace.h', + 'src/core/lib/event_engine/ares_resolver.h', 'src/core/lib/event_engine/cf_engine/cf_engine.h', 'src/core/lib/event_engine/cf_engine/cfstream_endpoint.h', 'src/core/lib/event_engine/cf_engine/cftype_unique_ref.h', @@ -2537,6 +2542,7 @@ Pod::Spec.new do |s| 'src/core/lib/event_engine/default_event_engine.h', 'src/core/lib/event_engine/default_event_engine_factory.h', 'src/core/lib/event_engine/forkable.h', + 'src/core/lib/event_engine/grpc_polled_fd.h', 'src/core/lib/event_engine/handle_containers.h', 'src/core/lib/event_engine/memory_allocator_factory.h', 'src/core/lib/event_engine/poller.h', @@ -2545,6 +2551,7 @@ Pod::Spec.new do |s| 'src/core/lib/event_engine/posix_engine/ev_poll_posix.h', 'src/core/lib/event_engine/posix_engine/event_poller.h', 'src/core/lib/event_engine/posix_engine/event_poller_posix_default.h', + 'src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h', 'src/core/lib/event_engine/posix_engine/internal_errqueue.h', 'src/core/lib/event_engine/posix_engine/lockfree_event.h', 'src/core/lib/event_engine/posix_engine/posix_endpoint.h', diff --git a/grpc.gemspec b/grpc.gemspec index f3dccb93dc2..4e8db380142 100644 --- a/grpc.gemspec +++ b/grpc.gemspec @@ -1063,6 +1063,8 @@ Gem::Specification.new do |s| s.files += %w( src/core/lib/debug/stats_data.h ) s.files += %w( src/core/lib/debug/trace.cc ) s.files += %w( src/core/lib/debug/trace.h ) + s.files += %w( src/core/lib/event_engine/ares_resolver.cc ) + s.files += %w( src/core/lib/event_engine/ares_resolver.h ) s.files += %w( src/core/lib/event_engine/cf_engine/cf_engine.cc ) s.files += %w( src/core/lib/event_engine/cf_engine/cf_engine.h ) s.files += %w( src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc ) @@ -1078,6 +1080,7 @@ Gem::Specification.new do |s| s.files += %w( src/core/lib/event_engine/event_engine.cc ) s.files += %w( src/core/lib/event_engine/forkable.cc ) s.files += %w( src/core/lib/event_engine/forkable.h ) + s.files += %w( src/core/lib/event_engine/grpc_polled_fd.h ) s.files += %w( src/core/lib/event_engine/handle_containers.h ) s.files += %w( src/core/lib/event_engine/memory_allocator.cc ) s.files += %w( src/core/lib/event_engine/memory_allocator_factory.h ) @@ -1090,6 +1093,7 @@ Gem::Specification.new do |s| s.files += %w( src/core/lib/event_engine/posix_engine/event_poller.h ) s.files += %w( src/core/lib/event_engine/posix_engine/event_poller_posix_default.cc ) s.files += %w( src/core/lib/event_engine/posix_engine/event_poller_posix_default.h ) + s.files += %w( src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h ) s.files += %w( src/core/lib/event_engine/posix_engine/internal_errqueue.cc ) s.files += %w( src/core/lib/event_engine/posix_engine/internal_errqueue.h ) s.files += %w( src/core/lib/event_engine/posix_engine/lockfree_event.cc ) diff --git a/grpc.gyp b/grpc.gyp index 57482cc048e..d82377051ff 100644 --- a/grpc.gyp +++ b/grpc.gyp @@ -739,6 +739,7 @@ 'src/core/lib/debug/stats.cc', 'src/core/lib/debug/stats_data.cc', 'src/core/lib/debug/trace.cc', + 'src/core/lib/event_engine/ares_resolver.cc', 'src/core/lib/event_engine/cf_engine/cf_engine.cc', 'src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc', 'src/core/lib/event_engine/channel_args_endpoint_config.cc', @@ -1237,6 +1238,7 @@ 'src/core/lib/debug/stats.cc', 'src/core/lib/debug/stats_data.cc', 'src/core/lib/debug/trace.cc', + 'src/core/lib/event_engine/ares_resolver.cc', 'src/core/lib/event_engine/cf_engine/cf_engine.cc', 'src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc', 'src/core/lib/event_engine/channel_args_endpoint_config.cc', @@ -1757,6 +1759,7 @@ 'src/core/lib/debug/stats.cc', 'src/core/lib/debug/stats_data.cc', 'src/core/lib/debug/trace.cc', + 'src/core/lib/event_engine/ares_resolver.cc', 'src/core/lib/event_engine/cf_engine/cf_engine.cc', 'src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc', 'src/core/lib/event_engine/channel_args_endpoint_config.cc', diff --git a/include/grpc/event_engine/event_engine.h b/include/grpc/event_engine/event_engine.h index 00b57629e73..4b671d2d3f7 100644 --- a/include/grpc/event_engine/event_engine.h +++ b/include/grpc/event_engine/event_engine.h @@ -365,7 +365,6 @@ class EventEngine : public std::enable_shared_from_this { /// lookup. Implementations should pass the appropriate statuses to the /// callback. For example, callbacks might expect to receive CANCELLED or /// NOT_FOUND. - /// virtual void LookupHostname(LookupHostnameCallback on_resolve, absl::string_view name, absl::string_view default_port) = 0; diff --git a/package.xml b/package.xml index db3a373a327..28a722fe16e 100644 --- a/package.xml +++ b/package.xml @@ -1045,6 +1045,8 @@ + + @@ -1060,6 +1062,7 @@ + @@ -1072,6 +1075,7 @@ + diff --git a/src/core/BUILD b/src/core/BUILD index 34736a7a53f..5b05d7481b4 100644 --- a/src/core/BUILD +++ b/src/core/BUILD @@ -1975,6 +1975,7 @@ grpc_cc_library( "absl/strings", ], deps = [ + "ares_resolver", "event_engine_common", "event_engine_poller", "event_engine_tcp_socket_utils", @@ -1996,6 +1997,7 @@ grpc_cc_library( "//:event_engine_base_hdrs", "//:gpr", "//:grpc_trace", + "//:orphanable", ], ) @@ -2298,6 +2300,49 @@ grpc_cc_library( ], ) +grpc_cc_library( + name = "ares_resolver", + srcs = [ + "lib/event_engine/ares_resolver.cc", + ], + hdrs = [ + "lib/event_engine/ares_resolver.h", + "lib/event_engine/grpc_polled_fd.h", + "lib/event_engine/posix_engine/grpc_polled_fd_posix.h", + ], + external_deps = [ + "absl/base:core_headers", + "absl/container:flat_hash_map", + "absl/functional:any_invocable", + "absl/hash", + "absl/status", + "absl/status:statusor", + "absl/strings", + "absl/strings:str_format", + "absl/types:optional", + "absl/types:variant", + "cares", + ], + deps = [ + "error", + "event_engine_time_util", + "grpc_sockaddr", + "iomgr_port", + "posix_event_engine_closure", + "posix_event_engine_event_poller", + "posix_event_engine_tcp_socket_utils", + "resolved_address", + "//:debug_location", + "//:event_engine_base_hdrs", + "//:gpr", + "//:grpc_trace", + "//:orphanable", + "//:parse_address", + "//:ref_counted_ptr", + "//:sockaddr_utils", + ], +) + grpc_cc_library( name = "channel_args_preconditioning", srcs = [ @@ -5153,6 +5198,7 @@ grpc_cc_library( "validation_errors", "//:backoff", "//:debug_location", + "//:exec_ctx", "//:gpr", "//:gpr_platform", "//:grpc_base", diff --git a/src/core/ext/filters/client_channel/resolver/dns/c_ares/grpc_ares_wrapper.cc b/src/core/ext/filters/client_channel/resolver/dns/c_ares/grpc_ares_wrapper.cc index cee19837b48..f8278cc6fc2 100644 --- a/src/core/ext/filters/client_channel/resolver/dns/c_ares/grpc_ares_wrapper.cc +++ b/src/core/ext/filters/client_channel/resolver/dns/c_ares/grpc_ares_wrapper.cc @@ -510,9 +510,9 @@ void grpc_ares_ev_driver_start_locked(grpc_ares_ev_driver* ev_driver) &ev_driver->on_ares_backup_poll_alarm_locked); } -static void noop_inject_channel_config(ares_channel /*channel*/) {} +static void noop_inject_channel_config(ares_channel* /*channel*/) {} -void (*grpc_ares_test_only_inject_config)(ares_channel channel) = +void (*grpc_ares_test_only_inject_config)(ares_channel* channel) = noop_inject_channel_config; grpc_error_handle grpc_ares_ev_driver_create_locked( @@ -524,7 +524,7 @@ grpc_error_handle grpc_ares_ev_driver_create_locked( memset(&opts, 0, sizeof(opts)); opts.flags |= ARES_FLAG_STAYOPEN; int status = ares_init_options(&(*ev_driver)->channel, &opts, ARES_OPT_FLAGS); - grpc_ares_test_only_inject_config((*ev_driver)->channel); + grpc_ares_test_only_inject_config(&(*ev_driver)->channel); GRPC_CARES_TRACE_LOG("request:%p grpc_ares_ev_driver_create_locked", request); if (status != ARES_SUCCESS) { grpc_error_handle err = GRPC_ERROR_CREATE(absl::StrCat( diff --git a/src/core/ext/filters/client_channel/resolver/dns/c_ares/grpc_ares_wrapper.h b/src/core/ext/filters/client_channel/resolver/dns/c_ares/grpc_ares_wrapper.h index ffe7cf9e01a..69f52bc3df1 100644 --- a/src/core/ext/filters/client_channel/resolver/dns/c_ares/grpc_ares_wrapper.h +++ b/src/core/ext/filters/client_channel/resolver/dns/c_ares/grpc_ares_wrapper.h @@ -131,6 +131,6 @@ void grpc_cares_wrapper_address_sorting_sort( const grpc_ares_request* request, grpc_core::ServerAddressList* addresses); // Exposed in this header for C-core tests only -extern void (*grpc_ares_test_only_inject_config)(ares_channel channel); +extern void (*grpc_ares_test_only_inject_config)(ares_channel* channel); #endif // GRPC_SRC_CORE_EXT_FILTERS_CLIENT_CHANNEL_RESOLVER_DNS_C_ARES_GRPC_ARES_WRAPPER_H diff --git a/src/core/ext/filters/client_channel/resolver/dns/event_engine/event_engine_client_channel_resolver.cc b/src/core/ext/filters/client_channel/resolver/dns/event_engine/event_engine_client_channel_resolver.cc index a8596e7c7ca..b2b03cf5a40 100644 --- a/src/core/ext/filters/client_channel/resolver/dns/event_engine/event_engine_client_channel_resolver.cc +++ b/src/core/ext/filters/client_channel/resolver/dns/event_engine/event_engine_client_channel_resolver.cc @@ -22,7 +22,6 @@ #include #include #include -#include #include #include @@ -51,6 +50,7 @@ #include "src/core/lib/gprpp/sync.h" #include "src/core/lib/gprpp/time.h" #include "src/core/lib/gprpp/validation_errors.h" +#include "src/core/lib/iomgr/exec_ctx.h" #include "src/core/lib/iomgr/resolve_address.h" #include "src/core/lib/resolver/resolver.h" #include "src/core/lib/resolver/resolver_factory.h" @@ -231,8 +231,12 @@ EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: is_hostname_inflight_ = true; event_engine_resolver_->LookupHostname( [self = Ref(DEBUG_LOCATION, "OnHostnameResolved")]( - absl::StatusOr> addresses) { + absl::StatusOr> + addresses) mutable { + ApplicationCallbackExecCtx callback_exec_ctx; + ExecCtx exec_ctx; self->OnHostnameResolved(std::move(addresses)); + self.reset(); }, resolver_->name_to_resolve(), kDefaultSecurePort); if (resolver_->enable_srv_queries_) { @@ -243,8 +247,13 @@ EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: event_engine_resolver_->LookupSRV( [self = Ref(DEBUG_LOCATION, "OnSRVResolved")]( absl::StatusOr> - srv_records) { self->OnSRVResolved(std::move(srv_records)); }, - resolver_->name_to_resolve()); + srv_records) mutable { + ApplicationCallbackExecCtx callback_exec_ctx; + ExecCtx exec_ctx; + self->OnSRVResolved(std::move(srv_records)); + self.reset(); + }, + absl::StrCat("_grpclb._tcp.", resolver_->name_to_resolve())); } if (resolver_->request_service_config_) { GRPC_EVENT_ENGINE_RESOLVER_TRACE( @@ -253,8 +262,11 @@ EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: is_txt_inflight_ = true; event_engine_resolver_->LookupTXT( [self = Ref(DEBUG_LOCATION, "OnTXTResolved")]( - absl::StatusOr> service_config) { + absl::StatusOr> service_config) mutable { + ApplicationCallbackExecCtx callback_exec_ctx; + ExecCtx exec_ctx; self->OnTXTResolved(std::move(service_config)); + self.reset(); }, absl::StrCat("_grpc_config.", resolver_->name_to_resolve())); } @@ -263,8 +275,12 @@ EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: ? EventEngine::Duration::max() : resolver_->query_timeout_ms_; timeout_handle_ = resolver_->event_engine_->RunAfter( - timeout, - [self = Ref(DEBUG_LOCATION, "OnTimeout")]() { self->OnTimeout(); }); + timeout, [self = Ref(DEBUG_LOCATION, "OnTimeout")]() mutable { + ApplicationCallbackExecCtx callback_exec_ctx; + ExecCtx exec_ctx; + self->OnTimeout(); + self.reset(); + }); } EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: @@ -300,10 +316,11 @@ void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: OnHostnameResolved(absl::StatusOr> new_addresses) { - ValidationErrors::ScopedField field(&errors_, "hostname lookup"); absl::optional result; { MutexLock lock(&on_resolved_mu_); + // Make sure field destroys before cleanup. + ValidationErrors::ScopedField field(&errors_, "hostname lookup"); if (orphaned_) return; is_hostname_inflight_ = false; if (!new_addresses.ok()) { @@ -325,7 +342,6 @@ void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: OnSRVResolved( absl::StatusOr> srv_records) { - ValidationErrors::ScopedField field(&errors_, "srv lookup"); absl::optional result; auto cleanup = absl::MakeCleanup([&]() { if (result.has_value()) { @@ -333,6 +349,8 @@ void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: } }); MutexLock lock(&on_resolved_mu_); + // Make sure field destroys before cleanup. + ValidationErrors::ScopedField field(&errors_, "srv lookup"); if (orphaned_) return; is_srv_inflight_ = false; if (!srv_records.ok()) { @@ -359,12 +377,15 @@ void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: resolver_.get(), srv_record.host.c_str(), srv_record.port); ++number_of_balancer_hostnames_initiated_; event_engine_resolver_->LookupHostname( - [host = std::move(srv_record.host), + [host = srv_record.host, self = Ref(DEBUG_LOCATION, "OnBalancerHostnamesResolved")]( absl::StatusOr> new_balancer_addresses) mutable { + ApplicationCallbackExecCtx callback_exec_ctx; + ExecCtx exec_ctx; self->OnBalancerHostnamesResolved(std::move(host), std::move(new_balancer_addresses)); + self.reset(); }, srv_record.host, std::to_string(srv_record.port)); } @@ -375,8 +396,6 @@ void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: std::string authority, absl::StatusOr> new_balancer_addresses) { - ValidationErrors::ScopedField field( - &errors_, absl::StrCat("balancer lookup for ", authority)); absl::optional result; auto cleanup = absl::MakeCleanup([&]() { if (result.has_value()) { @@ -384,6 +403,9 @@ void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: } }); MutexLock lock(&on_resolved_mu_); + // Make sure field destroys before cleanup. + ValidationErrors::ScopedField field( + &errors_, absl::StrCat("balancer lookup for ", authority)); if (orphaned_) return; ++number_of_balancer_hostnames_resolved_; if (!new_balancer_addresses.ok()) { @@ -405,10 +427,11 @@ void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: OnTXTResolved(absl::StatusOr> service_config) { - ValidationErrors::ScopedField field(&errors_, "txt lookup"); absl::optional result; { MutexLock lock(&on_resolved_mu_); + // Make sure field destroys before cleanup. + ValidationErrors::ScopedField field(&errors_, "txt lookup"); if (orphaned_) return; GPR_ASSERT(is_txt_inflight_); is_txt_inflight_ = false; @@ -424,8 +447,12 @@ void EventEngineClientChannelDNSResolver::EventEngineDNSRequestWrapper:: s, kServiceConfigAttributePrefix); }); if (result != service_config->end()) { + // Found a service config record. service_config_json_ = result->substr(kServiceConfigAttributePrefix.size()); + GRPC_EVENT_ENGINE_RESOLVER_TRACE( + "DNSResolver::%p found service config: %s", + event_engine_resolver_.get(), service_config_json_->c_str()); } else { service_config_json_ = absl::UnavailableError(absl::StrCat( "failed to find attribute prefix: ", kServiceConfigAttributePrefix, diff --git a/src/core/ext/filters/client_channel/resolver/polling_resolver.cc b/src/core/ext/filters/client_channel/resolver/polling_resolver.cc index ff94e9834de..bd99aebae53 100644 --- a/src/core/ext/filters/client_channel/resolver/polling_resolver.cc +++ b/src/core/ext/filters/client_channel/resolver/polling_resolver.cc @@ -159,7 +159,7 @@ void PollingResolver::OnRequestCompleteLocked(Result result) { if (GPR_UNLIKELY(tracer_ != nullptr && tracer_->enabled())) { gpr_log(GPR_INFO, "[polling resolver %p] returning result: " - "addresses=%s, service_config=%s", + "addresses=%s, service_config=%s, resolution_note=%s", this, result.addresses.ok() ? absl::StrCat("<", result.addresses->size(), " addresses>") @@ -170,7 +170,8 @@ void PollingResolver::OnRequestCompleteLocked(Result result) { ? "" : std::string((*result.service_config)->json_string()) .c_str()) - : result.service_config.status().ToString().c_str()); + : result.service_config.status().ToString().c_str(), + result.resolution_note.c_str()); } GPR_ASSERT(result.result_health_callback == nullptr); RefCountedPtr self = diff --git a/src/core/lib/event_engine/ares_resolver.cc b/src/core/lib/event_engine/ares_resolver.cc new file mode 100644 index 00000000000..c78028e5be9 --- /dev/null +++ b/src/core/lib/event_engine/ares_resolver.cc @@ -0,0 +1,705 @@ +// Copyright 2023 The 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. +#include + +#include "src/core/lib/event_engine/ares_resolver.h" + +#include + +#include +#include + +#include "src/core/lib/iomgr/port.h" + +// IWYU pragma: no_include +// IWYU pragma: no_include +// IWYU pragma: no_include +// IWYU pragma: no_include +// IWYU pragma: no_include +// IWYU pragma: no_include +// IWYU pragma: no_include +// IWYU pragma: no_include + +#if GRPC_ARES == 1 + +#include +#include + +#include +#include +#include +#include +#include +#include + +#include "absl/functional/any_invocable.h" +#include "absl/hash/hash.h" +#include "absl/strings/match.h" +#include "absl/strings/numbers.h" +#include "absl/strings/str_cat.h" +#include "absl/strings/str_format.h" +#include "absl/types/optional.h" + +#include +#include + +#include "src/core/lib/address_utils/parse_address.h" +#include "src/core/lib/address_utils/sockaddr_utils.h" +#include "src/core/lib/event_engine/grpc_polled_fd.h" +#include "src/core/lib/event_engine/time_util.h" +#include "src/core/lib/gprpp/debug_location.h" +#include "src/core/lib/gprpp/host_port.h" +#include "src/core/lib/gprpp/ref_counted_ptr.h" +#include "src/core/lib/iomgr/resolved_address.h" +#include "src/core/lib/iomgr/sockaddr.h" +#ifdef GRPC_POSIX_SOCKET_ARES_EV_DRIVER +#include "src/core/lib/event_engine/posix_engine/tcp_socket_utils.h" +#endif + +namespace grpc_event_engine { +namespace experimental { + +grpc_core::TraceFlag grpc_trace_ares_resolver(false, "cares_resolver"); + +namespace { + +absl::Status AresStatusToAbslStatus(int status, absl::string_view error_msg) { + switch (status) { + case ARES_ECANCELLED: + return absl::CancelledError(error_msg); + case ARES_ENOTIMP: + return absl::UnimplementedError(error_msg); + case ARES_ENOTFOUND: + return absl::NotFoundError(error_msg); + default: + return absl::UnknownError(error_msg); + } +} + +// An alternative here could be to use ares_timeout to try to be more +// accurate, but that would require using "struct timeval"'s, which just +// makes things a bit more complicated. So just poll every second, as +// suggested by the c-ares code comments. +constexpr EventEngine::Duration kAresBackupPollAlarmDuration = + std::chrono::seconds(1); + +bool IsIpv6LoopbackAvailable() { +#ifdef GRPC_POSIX_SOCKET_ARES_EV_DRIVER + return PosixSocketWrapper::IsIpv6LoopbackAvailable(); +#elif defined(GRPC_WINDOWS_SOCKET_ARES_EV_DRIVER) + // TODO(yijiem): implement this for Windows + return true; +#else +#error "Unsupported platform" +#endif +} + +absl::Status SetRequestDNSServer(absl::string_view dns_server, + ares_channel* channel) { + GRPC_ARES_RESOLVER_TRACE_LOG("Using DNS server %s", dns_server.data()); + grpc_resolved_address addr; + struct ares_addr_port_node dns_server_addr = {}; + if (grpc_parse_ipv4_hostport(dns_server, &addr, /*log_errors=*/false)) { + dns_server_addr.family = AF_INET; + struct sockaddr_in* in = reinterpret_cast(addr.addr); + memcpy(&dns_server_addr.addr.addr4, &in->sin_addr, sizeof(struct in_addr)); + dns_server_addr.tcp_port = grpc_sockaddr_get_port(&addr); + dns_server_addr.udp_port = grpc_sockaddr_get_port(&addr); + } else if (grpc_parse_ipv6_hostport(dns_server, &addr, + /*log_errors=*/false)) { + dns_server_addr.family = AF_INET6; + struct sockaddr_in6* in6 = + reinterpret_cast(addr.addr); + memcpy(&dns_server_addr.addr.addr6, &in6->sin6_addr, + sizeof(struct in6_addr)); + dns_server_addr.tcp_port = grpc_sockaddr_get_port(&addr); + dns_server_addr.udp_port = grpc_sockaddr_get_port(&addr); + } else { + return absl::InvalidArgumentError( + absl::StrCat("Cannot parse authority: ", dns_server)); + } + int status = ares_set_servers_ports(*channel, &dns_server_addr); + if (status != ARES_SUCCESS) { + return AresStatusToAbslStatus(status, ares_strerror(status)); + } + return absl::OkStatus(); +} + +struct QueryArg { + QueryArg(AresResolver* ar, int id, absl::string_view name) + : ares_resolver(ar), callback_map_id(id), query_name(name) {} + AresResolver* ares_resolver; + int callback_map_id; + std::string query_name; +}; + +struct HostnameQueryArg : public QueryArg { + HostnameQueryArg(AresResolver* ar, int id, absl::string_view name, int p) + : QueryArg(ar, id, name), port(p) {} + int port; +}; + +} // namespace + +absl::StatusOr> +AresResolver::CreateAresResolver( + absl::string_view dns_server, + std::unique_ptr polled_fd_factory, + std::shared_ptr event_engine) { + ares_options opts = {}; + opts.flags |= ARES_FLAG_STAYOPEN; + ares_channel channel; + int status = ares_init_options(&channel, &opts, ARES_OPT_FLAGS); + if (status != ARES_SUCCESS) { + gpr_log(GPR_ERROR, "ares_init_options failed, status: %d", status); + return AresStatusToAbslStatus( + status, + absl::StrCat("Failed to init c-ares channel: ", ares_strerror(status))); + } + event_engine_grpc_ares_test_only_inject_config(&channel); + if (!dns_server.empty()) { + absl::Status status = SetRequestDNSServer(dns_server, &channel); + if (!status.ok()) { + return status; + } + } + return grpc_core::MakeOrphanable( + std::move(polled_fd_factory), std::move(event_engine), channel); +} + +AresResolver::~AresResolver() { + GPR_ASSERT(fd_node_list_.empty()); + GPR_ASSERT(callback_map_.empty()); + ares_destroy(channel_); +} + +void AresResolver::Orphan() { + { + grpc_core::MutexLock lock(&mutex_); + shutting_down_ = true; + if (ares_backup_poll_alarm_handle_.has_value()) { + event_engine_->Cancel(*ares_backup_poll_alarm_handle_); + ares_backup_poll_alarm_handle_.reset(); + } + for (const auto& fd_node : fd_node_list_) { + if (!fd_node->already_shutdown) { + GRPC_ARES_RESOLVER_TRACE_LOG("request: %p shutdown fd: %s", this, + fd_node->polled_fd->GetName()); + fd_node->polled_fd->ShutdownLocked( + absl::CancelledError("AresResolver::Orphan")); + fd_node->already_shutdown = true; + } + } + } + Unref(DEBUG_LOCATION, "Orphan"); +} + +void AresResolver::LookupHostname( + absl::string_view name, absl::string_view default_port, + EventEngine::DNSResolver::LookupHostnameCallback callback) { + absl::string_view host; + absl::string_view port_string; + if (!grpc_core::SplitHostPort(name, &host, &port_string)) { + event_engine_->Run( + [callback = std::move(callback), + status = absl::InvalidArgumentError(absl::StrCat( + "Unparseable name: ", name))]() mutable { callback(status); }); + return; + } + GPR_ASSERT(!host.empty()); + if (port_string.empty()) { + if (default_port.empty()) { + event_engine_->Run([callback = std::move(callback), + status = absl::InvalidArgumentError(absl::StrFormat( + "No port in name %s or default_port argument", + name))]() mutable { callback(status); }); + return; + } + port_string = default_port; + } + int port = 0; + if (port_string == "http") { + port = 80; + } else if (port_string == "https") { + port = 443; + } else if (!absl::SimpleAtoi(port_string, &port)) { + event_engine_->Run([callback = std::move(callback), + status = absl::InvalidArgumentError(absl::StrCat( + "Failed to parse port in name: ", + name))]() mutable { callback(status); }); + return; + } + // TODO(yijiem): Change this when refactoring code in + // src/core/lib/address_utils to use EventEngine::ResolvedAddress. + grpc_resolved_address addr; + const std::string hostport = grpc_core::JoinHostPort(host, port); + if (grpc_parse_ipv4_hostport(hostport.c_str(), &addr, + false /* log errors */) || + grpc_parse_ipv6_hostport(hostport.c_str(), &addr, + false /* log errors */)) { + // Early out if the target is an ipv4 or ipv6 literal. + std::vector result; + result.emplace_back(reinterpret_cast(addr.addr), addr.len); + event_engine_->Run( + [callback = std::move(callback), result = std::move(result)]() mutable { + callback(std::move(result)); + }); + return; + } + grpc_core::MutexLock lock(&mutex_); + callback_map_.emplace(++id_, std::move(callback)); + auto* resolver_arg = new HostnameQueryArg(this, id_, name, port); + if (IsIpv6LoopbackAvailable()) { + ares_gethostbyname(channel_, std::string(host).c_str(), AF_UNSPEC, + &AresResolver::OnHostbynameDoneLocked, resolver_arg); + } else { + ares_gethostbyname(channel_, std::string(host).c_str(), AF_INET, + &AresResolver::OnHostbynameDoneLocked, resolver_arg); + } + CheckSocketsLocked(); + MaybeStartTimerLocked(); +} + +void AresResolver::LookupSRV( + absl::string_view name, + EventEngine::DNSResolver::LookupSRVCallback callback) { + absl::string_view host; + absl::string_view port; + if (!grpc_core::SplitHostPort(name, &host, &port)) { + event_engine_->Run( + [callback = std::move(callback), + status = absl::InvalidArgumentError(absl::StrCat( + "Unparseable name: ", name))]() mutable { callback(status); }); + return; + } + GPR_ASSERT(!host.empty()); + // Don't query for SRV records if the target is "localhost" + if (absl::EqualsIgnoreCase(host, "localhost")) { + event_engine_->Run([callback = std::move(callback)]() mutable { + callback(std::vector()); + }); + return; + } + grpc_core::MutexLock lock(&mutex_); + callback_map_.emplace(++id_, std::move(callback)); + auto* resolver_arg = new QueryArg(this, id_, host); + ares_query(channel_, std::string(host).c_str(), ns_c_in, ns_t_srv, + &AresResolver::OnSRVQueryDoneLocked, resolver_arg); + CheckSocketsLocked(); + MaybeStartTimerLocked(); +} + +void AresResolver::LookupTXT( + absl::string_view name, + EventEngine::DNSResolver::LookupTXTCallback callback) { + absl::string_view host; + absl::string_view port; + if (!grpc_core::SplitHostPort(name, &host, &port)) { + event_engine_->Run( + [callback = std::move(callback), + status = absl::InvalidArgumentError(absl::StrCat( + "Unparseable name: ", name))]() mutable { callback(status); }); + return; + } + GPR_ASSERT(!host.empty()); + // Don't query for TXT records if the target is "localhost" + if (absl::EqualsIgnoreCase(host, "localhost")) { + event_engine_->Run([callback = std::move(callback)]() mutable { + callback(std::vector()); + }); + return; + } + grpc_core::MutexLock lock(&mutex_); + callback_map_.emplace(++id_, std::move(callback)); + auto* resolver_arg = new QueryArg(this, id_, host); + ares_search(channel_, std::string(host).c_str(), ns_c_in, ns_t_txt, + &AresResolver::OnTXTDoneLocked, resolver_arg); + CheckSocketsLocked(); + MaybeStartTimerLocked(); +} + +AresResolver::AresResolver( + std::unique_ptr polled_fd_factory, + std::shared_ptr event_engine, ares_channel channel) + : grpc_core::InternallyRefCounted( + GRPC_TRACE_FLAG_ENABLED(grpc_trace_ares_resolver) ? "AresResolver" + : nullptr), + channel_(channel), + polled_fd_factory_(std::move(polled_fd_factory)), + event_engine_(std::move(event_engine)) {} + +void AresResolver::CheckSocketsLocked() { + FdNodeList new_list; + if (!shutting_down_) { + ares_socket_t socks[ARES_GETSOCK_MAXNUM]; + int socks_bitmask = ares_getsock(channel_, socks, ARES_GETSOCK_MAXNUM); + for (size_t i = 0; i < ARES_GETSOCK_MAXNUM; i++) { + if (ARES_GETSOCK_READABLE(socks_bitmask, i) || + ARES_GETSOCK_WRITABLE(socks_bitmask, i)) { + auto iter = std::find_if( + fd_node_list_.begin(), fd_node_list_.end(), + [sock = socks[i]](const auto& node) { return node->as == sock; }); + if (iter == fd_node_list_.end()) { + new_list.push_back(std::make_unique( + socks[i], polled_fd_factory_->NewGrpcPolledFdLocked(socks[i]))); + GRPC_ARES_RESOLVER_TRACE_LOG("request:%p new fd: %d", this, socks[i]); + } else { + new_list.splice(new_list.end(), fd_node_list_, iter); + } + FdNode* fd_node = new_list.back().get(); + if (ARES_GETSOCK_READABLE(socks_bitmask, i) && + !fd_node->readable_registered) { + fd_node->readable_registered = true; + if (fd_node->polled_fd->IsFdStillReadableLocked()) { + // If c-ares is interested to read and the socket already has data + // available for read, schedules OnReadable directly here. This is + // to cope with the edge-triggered poller not getting an event if no + // new data arrives and c-ares hasn't read all the data in the + // previous ares_process_fd. + GRPC_ARES_RESOLVER_TRACE_LOG( + "request:%p schedule read directly on: %d", this, fd_node->as); + event_engine_->Run( + [self = Ref(DEBUG_LOCATION, "CheckSocketsLocked"), + fd_node]() mutable { + self->OnReadable(fd_node, absl::OkStatus()); + }); + } else { + // Otherwise register with the poller for readable event. + GRPC_ARES_RESOLVER_TRACE_LOG("request:%p notify read on: %d", this, + fd_node->as); + fd_node->polled_fd->RegisterForOnReadableLocked( + [self = Ref(DEBUG_LOCATION, "CheckSocketsLocked"), + fd_node](absl::Status status) mutable { + self->OnReadable(fd_node, status); + }); + } + } + // Register write_closure if the socket is writable and write_closure + // has not been registered with this socket. + if (ARES_GETSOCK_WRITABLE(socks_bitmask, i) && + !fd_node->writable_registered) { + GRPC_ARES_RESOLVER_TRACE_LOG("request:%p notify write on: %d", this, + fd_node->as); + fd_node->writable_registered = true; + fd_node->polled_fd->RegisterForOnWriteableLocked( + [self = Ref(DEBUG_LOCATION, "CheckSocketsLocked"), + fd_node](absl::Status status) mutable { + self->OnWritable(fd_node, status); + }); + } + } + } + } + // Any remaining fds in fd_node_list_ were not returned by ares_getsock() + // and are therefore no longer in use, so they can be shut down and removed + // from the list. + while (!fd_node_list_.empty()) { + FdNode* fd_node = fd_node_list_.front().get(); + if (!fd_node->already_shutdown) { + GRPC_ARES_RESOLVER_TRACE_LOG("request: %p shutdown fd: %s", this, + fd_node->polled_fd->GetName()); + fd_node->polled_fd->ShutdownLocked(absl::OkStatus()); + fd_node->already_shutdown = true; + } + if (!fd_node->readable_registered && !fd_node->writable_registered) { + GRPC_ARES_RESOLVER_TRACE_LOG("request: %p delete fd: %s", this, + fd_node->polled_fd->GetName()); + fd_node_list_.pop_front(); + } else { + new_list.splice(new_list.end(), fd_node_list_, fd_node_list_.begin()); + } + } + fd_node_list_ = std::move(new_list); +} + +void AresResolver::MaybeStartTimerLocked() { + if (ares_backup_poll_alarm_handle_.has_value()) { + return; + } + // Initialize the backup poll alarm + GRPC_ARES_RESOLVER_TRACE_LOG( + "request:%p MaybeStartTimerLocked next ares process poll time in %zu ms", + this, Milliseconds(kAresBackupPollAlarmDuration)); + ares_backup_poll_alarm_handle_ = event_engine_->RunAfter( + kAresBackupPollAlarmDuration, + [self = Ref(DEBUG_LOCATION, "MaybeStartTimerLocked")]() { + self->OnAresBackupPollAlarm(); + }); +} + +void AresResolver::OnReadable(FdNode* fd_node, absl::Status status) { + grpc_core::MutexLock lock(&mutex_); + GPR_ASSERT(fd_node->readable_registered); + fd_node->readable_registered = false; + GRPC_ARES_RESOLVER_TRACE_LOG("OnReadable: fd: %d; request: %p; status: %s", + fd_node->as, this, status.ToString().c_str()); + if (status.ok() && !shutting_down_) { + ares_process_fd(channel_, fd_node->as, ARES_SOCKET_BAD); + } else { + // If error is not absl::OkStatus() or the resolution was cancelled, it + // means the fd has been shutdown or timed out. The pending lookups made + // on this request will be cancelled by the following ares_cancel(). The + // remaining file descriptors in this request will be cleaned up in the + // following Work() method. + ares_cancel(channel_); + } + CheckSocketsLocked(); +} + +void AresResolver::OnWritable(FdNode* fd_node, absl::Status status) { + grpc_core::MutexLock lock(&mutex_); + GPR_ASSERT(fd_node->writable_registered); + fd_node->writable_registered = false; + GRPC_ARES_RESOLVER_TRACE_LOG("OnWritable: fd: %d; request:%p; status: %s", + fd_node->as, this, status.ToString().c_str()); + if (status.ok() && !shutting_down_) { + ares_process_fd(channel_, ARES_SOCKET_BAD, fd_node->as); + } else { + // If error is not absl::OkStatus() or the resolution was cancelled, it + // means the fd has been shutdown or timed out. The pending lookups made + // on this request will be cancelled by the following ares_cancel(). The + // remaining file descriptors in this request will be cleaned up in the + // following Work() method. + ares_cancel(channel_); + } + CheckSocketsLocked(); +} + +// In case of non-responsive DNS servers, dropped packets, etc., c-ares has +// intelligent timeout and retry logic, which we can take advantage of by +// polling ares_process_fd on time intervals. Overall, the c-ares library is +// meant to be called into and given a chance to proceed name resolution: +// a) when fd events happen +// b) when some time has passed without fd events having happened +// For the latter, we use this backup poller. Also see +// https://github.com/grpc/grpc/pull/17688 description for more details. +void AresResolver::OnAresBackupPollAlarm() { + grpc_core::MutexLock lock(&mutex_); + ares_backup_poll_alarm_handle_.reset(); + GRPC_ARES_RESOLVER_TRACE_LOG( + "request:%p OnAresBackupPollAlarm shutting_down=%d.", this, + shutting_down_); + if (!shutting_down_) { + for (const auto& fd_node : fd_node_list_) { + if (!fd_node->already_shutdown) { + GRPC_ARES_RESOLVER_TRACE_LOG( + "request:%p OnAresBackupPollAlarm; ares_process_fd. fd=%s", this, + fd_node->polled_fd->GetName()); + ares_socket_t as = fd_node->polled_fd->GetWrappedAresSocketLocked(); + ares_process_fd(channel_, as, as); + } + } + MaybeStartTimerLocked(); + CheckSocketsLocked(); + } +} + +void AresResolver::OnHostbynameDoneLocked(void* arg, int status, + int /*timeouts*/, + struct hostent* hostent) { + std::unique_ptr hostname_qa( + static_cast(arg)); + auto* ares_resolver = hostname_qa->ares_resolver; + auto nh = ares_resolver->callback_map_.extract(hostname_qa->callback_map_id); + GPR_ASSERT(!nh.empty()); + GPR_ASSERT( + absl::holds_alternative( + nh.mapped())); + auto callback = absl::get( + std::move(nh.mapped())); + if (status != ARES_SUCCESS) { + std::string error_msg = + absl::StrFormat("address lookup failed for %s: %s", + hostname_qa->query_name, ares_strerror(status)); + GRPC_ARES_RESOLVER_TRACE_LOG("resolver:%p OnHostbynameDoneLocked: %s", + ares_resolver, error_msg.c_str()); + ares_resolver->event_engine_->Run( + [callback = std::move(callback), + status = AresStatusToAbslStatus(status, error_msg)]() mutable { + callback(status); + }); + return; + } + GRPC_ARES_RESOLVER_TRACE_LOG( + "resolver:%p OnHostbynameDoneLocked name=%s ARES_SUCCESS", ares_resolver, + hostname_qa->query_name.c_str()); + std::vector result; + for (size_t i = 0; hostent->h_addr_list[i] != nullptr; i++) { + switch (hostent->h_addrtype) { + case AF_INET6: { + size_t addr_len = sizeof(struct sockaddr_in6); + struct sockaddr_in6 addr; + memset(&addr, 0, addr_len); + memcpy(&addr.sin6_addr, hostent->h_addr_list[i], + sizeof(struct in6_addr)); + addr.sin6_family = static_cast(hostent->h_addrtype); + addr.sin6_port = htons(hostname_qa->port); + result.emplace_back(reinterpret_cast(&addr), addr_len); + char output[INET6_ADDRSTRLEN]; + ares_inet_ntop(AF_INET6, &addr.sin6_addr, output, INET6_ADDRSTRLEN); + GRPC_ARES_RESOLVER_TRACE_LOG( + "resolver:%p c-ares resolver gets a AF_INET6 result: \n" + " addr: %s\n port: %d\n sin6_scope_id: %d\n", + ares_resolver, output, hostname_qa->port, addr.sin6_scope_id); + break; + } + case AF_INET: { + size_t addr_len = sizeof(struct sockaddr_in); + struct sockaddr_in addr; + memset(&addr, 0, addr_len); + memcpy(&addr.sin_addr, hostent->h_addr_list[i], sizeof(struct in_addr)); + addr.sin_family = static_cast(hostent->h_addrtype); + addr.sin_port = htons(hostname_qa->port); + result.emplace_back(reinterpret_cast(&addr), addr_len); + char output[INET_ADDRSTRLEN]; + ares_inet_ntop(AF_INET, &addr.sin_addr, output, INET_ADDRSTRLEN); + GRPC_ARES_RESOLVER_TRACE_LOG( + "resolver:%p c-ares resolver gets a AF_INET result: \n" + " addr: %s\n port: %d\n", + ares_resolver, output, hostname_qa->port); + break; + } + } + } + ares_resolver->event_engine_->Run( + [callback = std::move(callback), result = std::move(result)]() mutable { + callback(std::move(result)); + }); +} + +void AresResolver::OnSRVQueryDoneLocked(void* arg, int status, int /*timeouts*/, + unsigned char* abuf, int alen) { + std::unique_ptr qa(static_cast(arg)); + auto* ares_resolver = qa->ares_resolver; + auto nh = ares_resolver->callback_map_.extract(qa->callback_map_id); + GPR_ASSERT(!nh.empty()); + GPR_ASSERT( + absl::holds_alternative( + nh.mapped())); + auto callback = absl::get( + std::move(nh.mapped())); + auto fail = [&](absl::string_view prefix) { + std::string error_message = absl::StrFormat( + "%s for %s: %s", prefix, qa->query_name, ares_strerror(status)); + GRPC_ARES_RESOLVER_TRACE_LOG("OnSRVQueryDoneLocked: %s", + error_message.c_str()); + ares_resolver->event_engine_->Run( + [callback = std::move(callback), + status = AresStatusToAbslStatus(status, error_message)]() mutable { + callback(status); + }); + }; + if (status != ARES_SUCCESS) { + fail("SRV lookup failed"); + return; + } + GRPC_ARES_RESOLVER_TRACE_LOG( + "resolver:%p OnSRVQueryDoneLocked name=%s ARES_SUCCESS", ares_resolver, + qa->query_name.c_str()); + struct ares_srv_reply* reply = nullptr; + status = ares_parse_srv_reply(abuf, alen, &reply); + GRPC_ARES_RESOLVER_TRACE_LOG("resolver:%p ares_parse_srv_reply: %d", + ares_resolver, status); + if (status != ARES_SUCCESS) { + fail("Failed to parse SRV reply"); + return; + } + std::vector result; + for (struct ares_srv_reply* srv_it = reply; srv_it != nullptr; + srv_it = srv_it->next) { + EventEngine::DNSResolver::SRVRecord record; + record.host = srv_it->host; + record.port = srv_it->port; + record.priority = srv_it->priority; + record.weight = srv_it->weight; + result.push_back(std::move(record)); + } + if (reply != nullptr) { + ares_free_data(reply); + } + ares_resolver->event_engine_->Run( + [callback = std::move(callback), result = std::move(result)]() mutable { + callback(std::move(result)); + }); +} + +void AresResolver::OnTXTDoneLocked(void* arg, int status, int /*timeouts*/, + unsigned char* buf, int len) { + std::unique_ptr qa(static_cast(arg)); + auto* ares_resolver = qa->ares_resolver; + auto nh = ares_resolver->callback_map_.extract(qa->callback_map_id); + GPR_ASSERT(!nh.empty()); + GPR_ASSERT( + absl::holds_alternative( + nh.mapped())); + auto callback = absl::get( + std::move(nh.mapped())); + auto fail = [&](absl::string_view prefix) { + std::string error_message = absl::StrFormat( + "%s for %s: %s", prefix, qa->query_name, ares_strerror(status)); + GRPC_ARES_RESOLVER_TRACE_LOG("resolver:%p OnTXTDoneLocked: %s", + ares_resolver, error_message.c_str()); + ares_resolver->event_engine_->Run( + [callback = std::move(callback), + status = AresStatusToAbslStatus(status, error_message)]() mutable { + callback(status); + }); + }; + if (status != ARES_SUCCESS) { + fail("TXT lookup failed"); + return; + } + GRPC_ARES_RESOLVER_TRACE_LOG( + "resolver:%p OnTXTDoneLocked name=%s ARES_SUCCESS", ares_resolver, + qa->query_name.c_str()); + struct ares_txt_ext* reply = nullptr; + status = ares_parse_txt_reply_ext(buf, len, &reply); + if (status != ARES_SUCCESS) { + fail("Failed to parse TXT result"); + return; + } + std::vector result; + for (struct ares_txt_ext* part = reply; part != nullptr; part = part->next) { + if (part->record_start) { + result.emplace_back(reinterpret_cast(part->txt), part->length); + } else { + absl::StrAppend( + &result.back(), + std::string(reinterpret_cast(part->txt), part->length)); + } + } + GRPC_ARES_RESOLVER_TRACE_LOG("resolver:%p Got %zu TXT records", ares_resolver, + result.size()); + if (GRPC_TRACE_FLAG_ENABLED(grpc_trace_ares_resolver)) { + for (const auto& record : result) { + gpr_log(GPR_INFO, "%s", record.c_str()); + } + } + // Clean up. + ares_free_data(reply); + ares_resolver->event_engine_->Run( + [callback = std::move(callback), result = std::move(result)]() mutable { + callback(std::move(result)); + }); +} + +} // namespace experimental +} // namespace grpc_event_engine + +void noop_inject_channel_config(ares_channel* /*channel*/) {} + +void (*event_engine_grpc_ares_test_only_inject_config)(ares_channel* channel) = + noop_inject_channel_config; + +#endif // GRPC_ARES == 1 diff --git a/src/core/lib/event_engine/ares_resolver.h b/src/core/lib/event_engine/ares_resolver.h new file mode 100644 index 00000000000..dbae35cff8f --- /dev/null +++ b/src/core/lib/event_engine/ares_resolver.h @@ -0,0 +1,147 @@ +// Copyright 2023 The 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_LIB_EVENT_ENGINE_ARES_RESOLVER_H +#define GRPC_SRC_CORE_LIB_EVENT_ENGINE_ARES_RESOLVER_H + +#include + +#include "src/core/lib/debug/trace.h" + +#if GRPC_ARES == 1 + +#include +#include + +#include + +#include "absl/base/thread_annotations.h" +#include "absl/container/flat_hash_map.h" +#include "absl/status/status.h" +#include "absl/status/statusor.h" +#include "absl/strings/string_view.h" +#include "absl/types/optional.h" +#include "absl/types/variant.h" + +#include +#include + +#include "src/core/lib/event_engine/grpc_polled_fd.h" +#include "src/core/lib/gprpp/orphanable.h" +#include "src/core/lib/gprpp/sync.h" + +namespace grpc_event_engine { +namespace experimental { + +extern grpc_core::TraceFlag grpc_trace_ares_resolver; + +#define GRPC_ARES_RESOLVER_TRACE_LOG(format, ...) \ + do { \ + if (GRPC_TRACE_FLAG_ENABLED(grpc_trace_ares_resolver)) { \ + gpr_log(GPR_INFO, "(EventEngine c-ares resolver) " format, __VA_ARGS__); \ + } \ + } while (0) + +class AresResolver : public grpc_core::InternallyRefCounted { + public: + static absl::StatusOr> + CreateAresResolver(absl::string_view dns_server, + std::unique_ptr polled_fd_factory, + std::shared_ptr event_engine); + + // Do not instantiate directly -- use CreateAresResolver() instead. + AresResolver(std::unique_ptr polled_fd_factory, + std::shared_ptr event_engine, ares_channel channel); + ~AresResolver() override; + void Orphan() override ABSL_LOCKS_EXCLUDED(mutex_); + + void LookupHostname(absl::string_view name, absl::string_view default_port, + EventEngine::DNSResolver::LookupHostnameCallback callback) + ABSL_LOCKS_EXCLUDED(mutex_); + void LookupSRV(absl::string_view name, + EventEngine::DNSResolver::LookupSRVCallback callback) + ABSL_LOCKS_EXCLUDED(mutex_); + void LookupTXT(absl::string_view name, + EventEngine::DNSResolver::LookupTXTCallback callback) + ABSL_LOCKS_EXCLUDED(mutex_); + + private: + // A FdNode saves (not owns) a live socket/fd which c-ares creates, and owns a + // GrpcPolledFd object which has a platform-agnostic interface to interact + // with the poller. The liveness of the socket means that c-ares needs us to + // monitor r/w events on this socket and notifies c-ares when such events have + // happened which we achieve through the GrpcPolledFd object. FdNode also + // handles the shutdown (maybe due to socket no longer used, finished request, + // cancel or timeout) and the destruction of the poller handle. Note that + // FdNode does not own the socket and it's the c-ares' responsibility to + // close the socket (possibly through ares_destroy). + struct FdNode { + FdNode() = default; + FdNode(ares_socket_t as, GrpcPolledFd* polled_fd) + : as(as), polled_fd(polled_fd) {} + ares_socket_t as; + std::unique_ptr polled_fd; + // true if the readable closure has been registered + bool readable_registered = false; + // true if the writable closure has been registered + bool writable_registered = false; + bool already_shutdown = false; + }; + using FdNodeList = std::list>; + + using CallbackType = + absl::variant; + + void CheckSocketsLocked() ABSL_EXCLUSIVE_LOCKS_REQUIRED(mutex_); + void MaybeStartTimerLocked() ABSL_EXCLUSIVE_LOCKS_REQUIRED(mutex_); + void OnReadable(FdNode* fd_node, absl::Status status) + ABSL_LOCKS_EXCLUDED(mutex_); + void OnWritable(FdNode* fd_node, absl::Status status) + ABSL_LOCKS_EXCLUDED(mutex_); + void OnAresBackupPollAlarm() ABSL_LOCKS_EXCLUDED(mutex_); + + // These callbacks are invoked from the c-ares library, so disable thread + // safety analysis. We are guaranteed to be holding mutex_. + static void OnHostbynameDoneLocked(void* arg, int status, int /*timeouts*/, + struct hostent* hostent) + ABSL_NO_THREAD_SAFETY_ANALYSIS; + static void OnSRVQueryDoneLocked(void* arg, int status, int /*timeouts*/, + unsigned char* abuf, + int alen) ABSL_NO_THREAD_SAFETY_ANALYSIS; + static void OnTXTDoneLocked(void* arg, int status, int /*timeouts*/, + unsigned char* buf, + int len) ABSL_NO_THREAD_SAFETY_ANALYSIS; + + grpc_core::Mutex mutex_; + bool shutting_down_ ABSL_GUARDED_BY(mutex_) = false; + ares_channel channel_ ABSL_GUARDED_BY(mutex_); + FdNodeList fd_node_list_ ABSL_GUARDED_BY(mutex_); + int id_ ABSL_GUARDED_BY(mutex_) = 0; + absl::flat_hash_map callback_map_ ABSL_GUARDED_BY(mutex_); + absl::optional ares_backup_poll_alarm_handle_ + ABSL_GUARDED_BY(mutex_); + std::unique_ptr polled_fd_factory_; + std::shared_ptr event_engine_; +}; + +} // namespace experimental +} // namespace grpc_event_engine + +// Exposed in this header for C-core tests only +extern void (*event_engine_grpc_ares_test_only_inject_config)( + ares_channel* channel); + +#endif // GRPC_ARES == 1 +#endif // GRPC_SRC_CORE_LIB_EVENT_ENGINE_ARES_RESOLVER_H diff --git a/src/core/lib/event_engine/grpc_polled_fd.h b/src/core/lib/event_engine/grpc_polled_fd.h new file mode 100644 index 00000000000..bb66089d5ad --- /dev/null +++ b/src/core/lib/event_engine/grpc_polled_fd.h @@ -0,0 +1,73 @@ +// Copyright 2023 The 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_LIB_EVENT_ENGINE_GRPC_POLLED_FD_H +#define GRPC_SRC_CORE_LIB_EVENT_ENGINE_GRPC_POLLED_FD_H + +#include + +#if GRPC_ARES == 1 + +#include + +#include "absl/functional/any_invocable.h" +#include "absl/status/status.h" + +#include "src/core/lib/iomgr/error.h" + +namespace grpc_event_engine { +namespace experimental { + +// A wrapped fd that integrates with the EventEngine poller of the current +// platform. A GrpcPolledFd knows how to create grpc platform-specific poller +// handle from "ares_socket_t" sockets, and then sign up for +// readability/writeability with that poller handle, and do shutdown and +// destruction. +class GrpcPolledFd { + public: + virtual ~GrpcPolledFd() {} + // Called when c-ares library is interested and there's no pending callback + virtual void RegisterForOnReadableLocked( + absl::AnyInvocable read_closure) = 0; + // Called when c-ares library is interested and there's no pending callback + virtual void RegisterForOnWriteableLocked( + absl::AnyInvocable write_closure) = 0; + // Indicates if there is data left even after just being read from + virtual bool IsFdStillReadableLocked() = 0; + // Called once and only once. Must cause cancellation of any pending + // read/write callbacks. + virtual void ShutdownLocked(grpc_error_handle error) = 0; + // Get the underlying ares_socket_t that this was created from + virtual ares_socket_t GetWrappedAresSocketLocked() = 0; + // A unique name, for logging + virtual const char* GetName() const = 0; +}; + +// A GrpcPolledFdFactory is 1-to-1 with and owned by a GrpcAresRequest. It knows +// how to create GrpcPolledFd's for the current platform, and the +// GrpcAresRequest uses it for all of its fd's. +class GrpcPolledFdFactory { + public: + virtual ~GrpcPolledFdFactory() {} + // Creates a new wrapped fd for the current platform + virtual GrpcPolledFd* NewGrpcPolledFdLocked(ares_socket_t as) = 0; + // Optionally configures the ares channel after creation + virtual void ConfigureAresChannelLocked(ares_channel channel) = 0; +}; + +} // namespace experimental +} // namespace grpc_event_engine + +#endif // GRPC_ARES == 1 +#endif // GRPC_SRC_CORE_LIB_EVENT_ENGINE_GRPC_POLLED_FD_H diff --git a/src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h b/src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h new file mode 100644 index 00000000000..da54dbd3517 --- /dev/null +++ b/src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h @@ -0,0 +1,112 @@ +// Copyright 2023 The 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_LIB_EVENT_ENGINE_POSIX_ENGINE_GRPC_POLLED_FD_POSIX_H +#define GRPC_SRC_CORE_LIB_EVENT_ENGINE_POSIX_ENGINE_GRPC_POLLED_FD_POSIX_H + +#include + +#include "src/core/lib/iomgr/port.h" + +#if GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_ARES_EV_DRIVER) + +#include +#include + +#include +#include + +#include + +#include "absl/functional/any_invocable.h" +#include "absl/status/status.h" +#include "absl/strings/str_cat.h" + +#include "src/core/lib/event_engine/grpc_polled_fd.h" +#include "src/core/lib/event_engine/posix_engine/event_poller.h" +#include "src/core/lib/event_engine/posix_engine/posix_engine_closure.h" +#include "src/core/lib/iomgr/error.h" + +namespace grpc_event_engine { +namespace experimental { + +class GrpcPolledFdPosix : public GrpcPolledFd { + public: + GrpcPolledFdPosix(ares_socket_t as, EventHandle* handle) + : name_(absl::StrCat("c-ares fd: ", static_cast(as))), + as_(as), + handle_(handle) {} + + ~GrpcPolledFdPosix() override { + // c-ares library will close the fd. This fd may be picked up immediately by + // another thread and should not be closed by the following OrphanHandle. + int phony_release_fd; + handle_->OrphanHandle(/*on_done=*/nullptr, &phony_release_fd, + "c-ares query finished"); + } + + void RegisterForOnReadableLocked( + absl::AnyInvocable read_closure) override { + handle_->NotifyOnRead(new PosixEngineClosure(std::move(read_closure), + /*is_permanent=*/false)); + } + + void RegisterForOnWriteableLocked( + absl::AnyInvocable write_closure) override { + handle_->NotifyOnWrite(new PosixEngineClosure(std::move(write_closure), + /*is_permanent=*/false)); + } + + bool IsFdStillReadableLocked() override { + size_t bytes_available = 0; + return ioctl(handle_->WrappedFd(), FIONREAD, &bytes_available) == 0 && + bytes_available > 0; + } + + void ShutdownLocked(grpc_error_handle error) override { + handle_->ShutdownHandle(error); + } + + ares_socket_t GetWrappedAresSocketLocked() override { return as_; } + + const char* GetName() const override { return name_.c_str(); } + + private: + const std::string name_; + const ares_socket_t as_; + EventHandle* handle_; +}; + +class GrpcPolledFdFactoryPosix : public GrpcPolledFdFactory { + public: + explicit GrpcPolledFdFactoryPosix(PosixEventPoller* poller) + : poller_(poller) {} + + GrpcPolledFd* NewGrpcPolledFdLocked(ares_socket_t as) override { + return new GrpcPolledFdPosix( + as, + poller_->CreateHandle(as, "c-ares socket", poller_->CanTrackErrors())); + } + + void ConfigureAresChannelLocked(ares_channel /*channel*/) override {} + + private: + PosixEventPoller* poller_; +}; + +} // namespace experimental +} // namespace grpc_event_engine + +#endif // GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_ARES_EV_DRIVER) +#endif // GRPC_SRC_CORE_LIB_EVENT_ENGINE_POSIX_ENGINE_GRPC_POLLED_FD_POSIX_H diff --git a/src/core/lib/event_engine/posix_engine/posix_engine.cc b/src/core/lib/event_engine/posix_engine/posix_engine.cc index a1d217cd3bd..8e71e435f94 100644 --- a/src/core/lib/event_engine/posix_engine/posix_engine.cc +++ b/src/core/lib/event_engine/posix_engine/posix_engine.cc @@ -18,6 +18,7 @@ #include #include #include +#include #include #include #include @@ -37,8 +38,10 @@ #include #include "src/core/lib/debug/trace.h" +#include "src/core/lib/event_engine/grpc_polled_fd.h" #include "src/core/lib/event_engine/poller.h" #include "src/core/lib/event_engine/posix.h" +#include "src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h" #include "src/core/lib/event_engine/posix_engine/tcp_socket_utils.h" #include "src/core/lib/event_engine/posix_engine/timer.h" #include "src/core/lib/event_engine/tcp_socket_utils.h" @@ -487,10 +490,49 @@ EventEngine::TaskHandle PosixEventEngine::RunAfterInternal( return handle; } +#if GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP) + +PosixEventEngine::PosixDNSResolver::PosixDNSResolver( + grpc_core::OrphanablePtr ares_resolver) + : ares_resolver_(std::move(ares_resolver)) {} + +void PosixEventEngine::PosixDNSResolver::LookupHostname( + LookupHostnameCallback on_resolve, absl::string_view name, + absl::string_view default_port) { + ares_resolver_->LookupHostname(name, default_port, std::move(on_resolve)); +} + +void PosixEventEngine::PosixDNSResolver::LookupSRV(LookupSRVCallback on_resolve, + absl::string_view name) { + ares_resolver_->LookupSRV(name, std::move(on_resolve)); +} + +void PosixEventEngine::PosixDNSResolver::LookupTXT(LookupTXTCallback on_resolve, + absl::string_view name) { + ares_resolver_->LookupTXT(name, std::move(on_resolve)); +} + +#endif // GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP) + absl::StatusOr> PosixEventEngine::GetDNSResolver( - EventEngine::DNSResolver::ResolverOptions const& /*options*/) { + const EventEngine::DNSResolver::ResolverOptions& options) { +#if GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP) + auto ares_resolver = AresResolver::CreateAresResolver( + options.dns_server, + std::make_unique(poller_manager_->Poller()), + shared_from_this()); + if (!ares_resolver.ok()) { + return ares_resolver.status(); + } + return std::make_unique( + std::move(*ares_resolver)); +#else // GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP) + // TODO(yijiem): Implement a basic A/AAAA-only native resolver in + // PosixEventEngine. + (void)options; grpc_core::Crash("unimplemented"); +#endif // GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP) } bool PosixEventEngine::IsWorkerThread() { grpc_core::Crash("unimplemented"); } diff --git a/src/core/lib/event_engine/posix_engine/posix_engine.h b/src/core/lib/event_engine/posix_engine/posix_engine.h index 7fc40a1d33a..bd28bde9509 100644 --- a/src/core/lib/event_engine/posix_engine/posix_engine.h +++ b/src/core/lib/event_engine/posix_engine/posix_engine.h @@ -34,11 +34,13 @@ #include #include +#include "src/core/lib/event_engine/ares_resolver.h" #include "src/core/lib/event_engine/handle_containers.h" #include "src/core/lib/event_engine/posix.h" #include "src/core/lib/event_engine/posix_engine/event_poller.h" #include "src/core/lib/event_engine/posix_engine/timer_manager.h" #include "src/core/lib/event_engine/thread_pool/thread_pool.h" +#include "src/core/lib/gprpp/orphanable.h" #include "src/core/lib/gprpp/sync.h" #include "src/core/lib/iomgr/port.h" #include "src/core/lib/surface/init_internally.h" @@ -138,7 +140,11 @@ class PosixEventEngine final : public PosixEventEngineWithFdSupport, public: class PosixDNSResolver : public EventEngine::DNSResolver { public: - ~PosixDNSResolver() override; + PosixDNSResolver() = delete; +#if GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP) + explicit PosixDNSResolver( + grpc_core::OrphanablePtr ares_resolver); +#endif // GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP) void LookupHostname(LookupHostnameCallback on_resolve, absl::string_view name, absl::string_view default_port) override; @@ -146,6 +152,11 @@ class PosixEventEngine final : public PosixEventEngineWithFdSupport, absl::string_view name) override; void LookupTXT(LookupTXTCallback on_resolve, absl::string_view name) override; + +#if GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP) + private: + grpc_core::OrphanablePtr ares_resolver_; +#endif // GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP) }; #ifdef GRPC_POSIX_SOCKET_TCP diff --git a/src/core/lib/experiments/experiments.yaml b/src/core/lib/experiments/experiments.yaml index 9c5a6121f90..f34ffc7c4f0 100644 --- a/src/core/lib/experiments/experiments.yaml +++ b/src/core/lib/experiments/experiments.yaml @@ -119,7 +119,7 @@ If set, use EventEngine DNSResolver for client channel resolution expiry: 2023/10/01 owner: yijiem@google.com - test_tags: [] + test_tags: ["cancel_ares_query_test", "resolver_component_tests_runner_invoker"] allow_in_fuzzing_config: false - name: work_stealing description: diff --git a/src/python/grpcio/grpc_core_dependencies.py b/src/python/grpcio/grpc_core_dependencies.py index eadfa0e6e41..468e401a36a 100644 --- a/src/python/grpcio/grpc_core_dependencies.py +++ b/src/python/grpcio/grpc_core_dependencies.py @@ -491,6 +491,7 @@ CORE_SOURCE_FILES = [ 'src/core/lib/debug/stats.cc', 'src/core/lib/debug/stats_data.cc', 'src/core/lib/debug/trace.cc', + 'src/core/lib/event_engine/ares_resolver.cc', 'src/core/lib/event_engine/cf_engine/cf_engine.cc', 'src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc', 'src/core/lib/event_engine/channel_args_endpoint_config.cc', diff --git a/templates/CMakeLists.txt.template b/templates/CMakeLists.txt.template index ea308132378..4f058f20ece 100644 --- a/templates/CMakeLists.txt.template +++ b/templates/CMakeLists.txt.template @@ -82,7 +82,11 @@ deps.append("${_gRPC_ADDRESS_SORTING_LIBRARIES}") deps.append("${_gRPC_RE2_LIBRARIES}") deps.append("${_gRPC_UPB_LIBRARIES}") + # TODO(yijiem): These targets depend on grpc_base instead of grpc. Since we don't populate grpc_base as a cmake target, the sources all get collapsed into these targets. This workaround adds c-ares and/or re2 dependencies to these targets. We should clean this up. + if target_dict['name'] in ['frame_test']: + deps.append("${_gRPC_CARES_LIBRARIES}") if target_dict['name'] in ['grpc_authorization_provider']: + deps.append("${_gRPC_CARES_LIBRARIES}") deps.append("${_gRPC_RE2_LIBRARIES}") deps.append("${_gRPC_ALLTARGETS_LIBRARIES}") for d in target_dict.get('deps', []): diff --git a/templates/test/cpp/naming/resolver_component_tests_defs.include b/templates/test/cpp/naming/resolver_component_tests_defs.include index 94d403d52dd..6877f6a5b09 100644 --- a/templates/test/cpp/naming/resolver_component_tests_defs.include +++ b/templates/test/cpp/naming/resolver_component_tests_defs.include @@ -1,4 +1,4 @@ -<%def name="resolver_component_tests(tests)">#!/usr/bin/env python +<%def name="resolver_component_tests(tests)">#!/usr/bin/env python3 # Copyright 2015 gRPC authors. # # Licensed under the Apache License, Version 2.0 (the "License"); @@ -58,6 +58,9 @@ if cur_resolver and cur_resolver != 'ares': test_runner_log('Exit 1 without running tests.') sys.exit(1) os.environ.update({'GRPC_TRACE': 'cares_resolver,cares_address_sorting'}) +experiments = os.environ.get('GRPC_EXPERIMENTS') +if experiments is not None and 'event_engine_dns' in experiments: + os.environ.update({'GRPC_TRACE': 'event_engine_client_channel_resolver,cares_resolver'}) def wait_until_dns_server_is_up(args, dns_server_subprocess, diff --git a/test/core/event_engine/test_suite/BUILD b/test/core/event_engine/test_suite/BUILD index d3c3f6ac82c..cc485a4d9a2 100644 --- a/test/core/event_engine/test_suite/BUILD +++ b/test/core/event_engine/test_suite/BUILD @@ -50,6 +50,7 @@ grpc_cc_test( "//test/core/event_engine:event_engine_test_utils", "//test/core/event_engine/test_suite/posix:oracle_event_engine_posix", "//test/core/event_engine/test_suite/tests:client", + "//test/core/event_engine/test_suite/tests:dns", "//test/core/event_engine/test_suite/tests:server", "//test/core/event_engine/test_suite/tests:timer", ], diff --git a/test/core/event_engine/test_suite/tests/BUILD b/test/core/event_engine/test_suite/tests/BUILD index ece7abf8cfb..ef297860a97 100644 --- a/test/core/event_engine/test_suite/tests/BUILD +++ b/test/core/event_engine/test_suite/tests/BUILD @@ -54,9 +54,27 @@ grpc_cc_library( testonly = True, srcs = ["dns_test.cc"], hdrs = ["dns_test.h"], + data = [ + "dns_test_record_groups.yaml", + "//test/cpp/naming/utils:dns_resolver", + "//test/cpp/naming/utils:dns_server", + "//test/cpp/naming/utils:health_check", + "//test/cpp/naming/utils:tcp_connect", + ], + external_deps = [ + "absl/status:statusor", + "absl/strings", + "absl/strings:str_format", + "address_sorting", + ], deps = [ + "//src/core:env", "//test/core/event_engine:event_engine_test_utils", "//test/core/event_engine/test_suite:event_engine_test_framework", + "//test/core/util:fake_udp_and_tcp_server", + "//test/core/util:grpc_test_util_base", + "//test/cpp/util:get_grpc_test_runfile_dir", + "//test/cpp/util:test_util", ], alwayslink = 1, ) diff --git a/test/core/event_engine/test_suite/tests/dns_test.cc b/test/core/event_engine/test_suite/tests/dns_test.cc index 798a1871efb..cbd48e26bc7 100644 --- a/test/core/event_engine/test_suite/tests/dns_test.cc +++ b/test/core/event_engine/test_suite/tests/dns_test.cc @@ -12,10 +12,36 @@ // See the License for the specific language governing permissions and // limitations under the License. -#include +// IWYU pragma: no_include +// IWYU pragma: no_include -#include "src/core/lib/iomgr/exec_ctx.h" +#include +#include +#include +#include +#include +#include +#include + +#include "absl/status/status.h" +#include "absl/status/statusor.h" +#include "absl/strings/str_format.h" +#include "absl/strings/str_join.h" +#include "absl/strings/string_view.h" +#include "absl/types/optional.h" +#include "gmock/gmock.h" +#include "gtest/gtest.h" + +#include + +#include "src/core/lib/event_engine/tcp_socket_utils.h" +#include "src/core/lib/gprpp/notification.h" +#include "src/core/lib/iomgr/sockaddr.h" #include "test/core/event_engine/test_suite/event_engine_test_framework.h" +#include "test/core/util/fake_udp_and_tcp_server.h" +#include "test/core/util/port.h" +#include "test/cpp/util/get_grpc_test_runfile_dir.h" +#include "test/cpp/util/subprocess.h" namespace grpc_event_engine { namespace experimental { @@ -25,7 +51,493 @@ void InitDNSTests() {} } // namespace experimental } // namespace grpc_event_engine +#ifdef GPR_WINDOWS class EventEngineDNSTest : public EventEngineTest {}; -// TODO(hork): establish meaningful tests -TEST_F(EventEngineDNSTest, TODO) { grpc_core::ExecCtx exec_ctx; } +// TODO(yijiem): make the test run on Windows +TEST_F(EventEngineDNSTest, TODO) {} +#else + +namespace { + +using grpc_event_engine::experimental::EventEngine; +using grpc_event_engine::experimental::URIToResolvedAddress; +using SRVRecord = EventEngine::DNSResolver::SRVRecord; +using testing::ElementsAre; +using testing::Pointwise; +using testing::SizeIs; +using testing::UnorderedPointwise; + +// TODO(yijiem): make this portable for Windows +constexpr char kDNSTestRecordGroupsYamlPath[] = + "test/core/event_engine/test_suite/tests/dns_test_record_groups.yaml"; +// Invoke bazel's executable links to the .sh and .py scripts (don't use +// the .sh and .py suffixes) to make sure that we're using bazel's test +// environment. +constexpr char kDNSServerRelPath[] = "test/cpp/naming/utils/dns_server"; +constexpr char kDNSResolverRelPath[] = "test/cpp/naming/utils/dns_resolver"; +constexpr char kTCPConnectRelPath[] = "test/cpp/naming/utils/tcp_connect"; +constexpr char kHealthCheckRelPath[] = "test/cpp/naming/utils/health_check"; + +MATCHER(ResolvedAddressEq, "") { + const auto& addr0 = std::get<0>(arg); + const auto& addr1 = std::get<1>(arg); + return addr0.size() == addr1.size() && + memcmp(addr0.address(), addr1.address(), addr0.size()) == 0; +} + +MATCHER(SRVRecordEq, "") { + const auto& arg0 = std::get<0>(arg); + const auto& arg1 = std::get<1>(arg); + return arg0.host == arg1.host && arg0.port == arg1.port && + arg0.priority == arg1.priority && arg0.weight == arg1.weight; +} + +MATCHER(StatusCodeEq, "") { + return std::get<0>(arg).code() == std::get<1>(arg); +} + +} // namespace + +class EventEngineDNSTest : public EventEngineTest { + protected: + static void SetUpTestSuite() { + std::string test_records_path = kDNSTestRecordGroupsYamlPath; + std::string dns_server_path = kDNSServerRelPath; + std::string dns_resolver_path = kDNSResolverRelPath; + std::string tcp_connect_path = kTCPConnectRelPath; + std::string health_check_path = kHealthCheckRelPath; + absl::optional runfile_dir = grpc::GetGrpcTestRunFileDir(); + if (runfile_dir.has_value()) { + // We sure need a portable filesystem lib for this to work on Windows. + test_records_path = absl::StrJoin({*runfile_dir, test_records_path}, "/"); + dns_server_path = absl::StrJoin({*runfile_dir, dns_server_path}, "/"); + dns_resolver_path = absl::StrJoin({*runfile_dir, dns_resolver_path}, "/"); + tcp_connect_path = absl::StrJoin({*runfile_dir, tcp_connect_path}, "/"); + health_check_path = absl::StrJoin({*runfile_dir, health_check_path}, "/"); + } else { + // Invoke the .py scripts directly where they are in source code if we are + // not running with bazel. + dns_server_path += ".py"; + dns_resolver_path += ".py"; + tcp_connect_path += ".py"; + health_check_path += ".py"; + } + // 1. launch dns_server + int port = grpc_pick_unused_port_or_die(); + // -p -r + dns_server_.server_process = new grpc::SubProcess( + {dns_server_path, "-p", std::to_string(port), "-r", test_records_path}); + dns_server_.port = port; + + // 2. wait until dns_server is up (health check) + grpc::SubProcess health_check({ + health_check_path, + "-p", + std::to_string(port), + "--dns_resolver_bin_path", + dns_resolver_path, + "--tcp_connect_bin_path", + tcp_connect_path, + }); + int status = health_check.Join(); + // TODO(yijiem): make this portable for Windows + ASSERT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0); + } + + static void TearDownTestSuite() { + dns_server_.server_process->Interrupt(); + dns_server_.server_process->Join(); + delete dns_server_.server_process; + } + + std::unique_ptr CreateDefaultDNSResolver() { + std::shared_ptr test_ee(this->NewEventEngine()); + EventEngine::DNSResolver::ResolverOptions options; + options.dns_server = dns_server_.address(); + return *test_ee->GetDNSResolver(options); + } + + std::unique_ptr + CreateDNSResolverWithNonResponsiveServer() { + using FakeUdpAndTcpServer = grpc_core::testing::FakeUdpAndTcpServer; + // Start up fake non responsive DNS server + fake_dns_server_ = std::make_unique( + FakeUdpAndTcpServer::AcceptMode::kWaitForClientToSendFirstBytes, + FakeUdpAndTcpServer::CloseSocketUponCloseFromPeer); + const std::string dns_server = + absl::StrFormat("[::1]:%d", fake_dns_server_->port()); + std::shared_ptr test_ee(this->NewEventEngine()); + EventEngine::DNSResolver::ResolverOptions options; + options.dns_server = dns_server; + return *test_ee->GetDNSResolver(options); + } + + std::unique_ptr + CreateDNSResolverWithoutSpecifyingServer() { + std::shared_ptr test_ee(this->NewEventEngine()); + EventEngine::DNSResolver::ResolverOptions options; + return *test_ee->GetDNSResolver(options); + } + + struct DNSServer { + std::string address() { return "127.0.0.1:" + std::to_string(port); } + int port; + grpc::SubProcess* server_process; + }; + grpc_core::Notification dns_resolver_signal_; + + private: + static DNSServer dns_server_; + std::unique_ptr fake_dns_server_; +}; + +EventEngineDNSTest::DNSServer EventEngineDNSTest::dns_server_; + +TEST_F(EventEngineDNSTest, QueryNXHostname) { + auto dns_resolver = CreateDefaultDNSResolver(); + dns_resolver->LookupHostname( + [this](auto result) { + ASSERT_FALSE(result.ok()); + EXPECT_EQ(result.status(), + absl::NotFoundError("address lookup failed for " + "nonexisting-target.dns-test.event-" + "engine.: Domain name not found")); + dns_resolver_signal_.Notify(); + }, + "nonexisting-target.dns-test.event-engine.", /*default_port=*/"443"); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, QueryWithIPLiteral) { + auto dns_resolver = CreateDefaultDNSResolver(); + dns_resolver->LookupHostname( + [this](auto result) { + ASSERT_TRUE(result.ok()); + EXPECT_THAT(*result, + Pointwise(ResolvedAddressEq(), + {*URIToResolvedAddress("ipv4:4.3.2.1:1234")})); + dns_resolver_signal_.Notify(); + }, + "4.3.2.1:1234", + /*default_port=*/""); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, QueryARecord) { + auto dns_resolver = CreateDefaultDNSResolver(); + dns_resolver->LookupHostname( + [this](auto result) { + ASSERT_TRUE(result.ok()); + EXPECT_THAT(*result, UnorderedPointwise( + ResolvedAddressEq(), + {*URIToResolvedAddress("ipv4:1.2.3.4:443"), + *URIToResolvedAddress("ipv4:1.2.3.5:443"), + *URIToResolvedAddress("ipv4:1.2.3.6:443")})); + dns_resolver_signal_.Notify(); + }, + "ipv4-only-multi-target.dns-test.event-engine.", + /*default_port=*/"443"); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, QueryAAAARecord) { + auto dns_resolver = CreateDefaultDNSResolver(); + dns_resolver->LookupHostname( + [this](auto result) { + ASSERT_TRUE(result.ok()); + EXPECT_THAT( + *result, + UnorderedPointwise( + ResolvedAddressEq(), + {*URIToResolvedAddress("ipv6:[2607:f8b0:400a:801::1002]:443"), + *URIToResolvedAddress("ipv6:[2607:f8b0:400a:801::1003]:443"), + *URIToResolvedAddress( + "ipv6:[2607:f8b0:400a:801::1004]:443")})); + dns_resolver_signal_.Notify(); + }, + "ipv6-only-multi-target.dns-test.event-engine.:443", + /*default_port=*/""); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, TestAddressSorting) { + auto dns_resolver = CreateDefaultDNSResolver(); + dns_resolver->LookupHostname( + [this](auto result) { + ASSERT_TRUE(result.ok()); + EXPECT_THAT( + *result, + Pointwise(ResolvedAddressEq(), + {*URIToResolvedAddress("ipv6:[::1]:1234"), + *URIToResolvedAddress("ipv6:[2002::1111]:1234")})); + dns_resolver_signal_.Notify(); + }, + "ipv6-loopback-preferred-target.dns-test.event-engine.:1234", + /*default_port=*/""); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, QuerySRVRecord) { + const SRVRecord kExpectedRecords[] = { + {/*host=*/"ipv4-only-multi-target.dns-test.event-engine", /*port=*/1234, + /*priority=*/0, /*weight=*/0}, + {"ipv6-only-multi-target.dns-test.event-engine", 1234, 0, 0}, + }; + + auto dns_resolver = CreateDefaultDNSResolver(); + dns_resolver->LookupSRV( + [&kExpectedRecords, this](auto result) { + ASSERT_TRUE(result.ok()); + EXPECT_THAT(*result, Pointwise(SRVRecordEq(), kExpectedRecords)); + dns_resolver_signal_.Notify(); + }, + "_grpclb._tcp.srv-multi-target.dns-test.event-engine."); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, QuerySRVRecordWithLocalhost) { + auto dns_resolver = CreateDefaultDNSResolver(); + dns_resolver->LookupSRV( + [this](auto result) { + ASSERT_TRUE(result.ok()); + EXPECT_THAT(*result, SizeIs(0)); + dns_resolver_signal_.Notify(); + }, + "localhost:1000"); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, QueryTXTRecord) { + // clang-format off + const std::string kExpectedRecord = + "grpc_config=[{" + "\"serviceConfig\":{" + "\"loadBalancingPolicy\":\"round_robin\"," + "\"methodConfig\":[{" + "\"name\":[{" + "\"method\":\"Foo\"," + "\"service\":\"SimpleService\"" + "}]," + "\"waitForReady\":true" + "}]" + "}" + "}]"; + // clang-format on + + auto dns_resolver = CreateDefaultDNSResolver(); + dns_resolver->LookupTXT( + [&kExpectedRecord, this](auto result) { + ASSERT_TRUE(result.ok()); + EXPECT_THAT(*result, + ElementsAre(kExpectedRecord, "other_config=other config")); + dns_resolver_signal_.Notify(); + }, + "_grpc_config.simple-service.dns-test.event-engine."); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, QueryTXTRecordWithLocalhost) { + auto dns_resolver = CreateDefaultDNSResolver(); + dns_resolver->LookupTXT( + [this](auto result) { + ASSERT_TRUE(result.ok()); + EXPECT_THAT(*result, SizeIs(0)); + dns_resolver_signal_.Notify(); + }, + "localhost:1000"); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, TestCancelActiveDNSQuery) { + const std::string name = "dont-care-since-wont-be-resolved.test.com:1234"; + auto dns_resolver = CreateDNSResolverWithNonResponsiveServer(); + dns_resolver->LookupHostname( + [this](auto result) { + ASSERT_FALSE(result.ok()); + EXPECT_EQ(result.status(), + absl::CancelledError("address lookup failed for " + "dont-care-since-wont-be-resolved.test." + "com:1234: DNS query cancelled")); + dns_resolver_signal_.Notify(); + }, + name, "1234"); + dns_resolver.reset(); + dns_resolver_signal_.WaitForNotification(); +} + +#define EXPECT_SUCCESS() \ + do { \ + EXPECT_TRUE(result.ok()); \ + EXPECT_FALSE(result->empty()); \ + } while (0) + +// The following tests are almost 1-to-1 ported from +// test/core/iomgr/resolve_address_test.cc (except tests for the native DNS +// resolver and tests that would not make sense using the +// EventEngine::DNSResolver API). + +// START +TEST_F(EventEngineDNSTest, LocalHost) { + auto dns_resolver = CreateDNSResolverWithoutSpecifyingServer(); + dns_resolver->LookupHostname( + [this](auto result) { + EXPECT_SUCCESS(); + dns_resolver_signal_.Notify(); + }, + "localhost:1", ""); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, DefaultPort) { + auto dns_resolver = CreateDNSResolverWithoutSpecifyingServer(); + dns_resolver->LookupHostname( + [this](auto result) { + EXPECT_SUCCESS(); + dns_resolver_signal_.Notify(); + }, + "localhost", "1"); + dns_resolver_signal_.WaitForNotification(); +} + +// This test assumes the environment has an ipv6 loopback +TEST_F(EventEngineDNSTest, LocalhostResultHasIPv6First) { + auto dns_resolver = CreateDNSResolverWithoutSpecifyingServer(); + dns_resolver->LookupHostname( + [this](auto result) { + EXPECT_TRUE(result.ok()); + EXPECT_TRUE(!result->empty() && + (*result)[0].address()->sa_family == AF_INET6); + dns_resolver_signal_.Notify(); + }, + "localhost:1", ""); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, NonNumericDefaultPort) { + auto dns_resolver = CreateDNSResolverWithoutSpecifyingServer(); + dns_resolver->LookupHostname( + [this](auto result) { + EXPECT_SUCCESS(); + dns_resolver_signal_.Notify(); + }, + "localhost", "http"); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, MissingDefaultPort) { + auto dns_resolver = CreateDNSResolverWithoutSpecifyingServer(); + dns_resolver->LookupHostname( + [this](auto result) { + EXPECT_FALSE(result.ok()); + dns_resolver_signal_.Notify(); + }, + "localhost", ""); + dns_resolver_signal_.WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, IPv6WithPort) { + auto dns_resolver = CreateDNSResolverWithoutSpecifyingServer(); + dns_resolver->LookupHostname( + [this](auto result) { + EXPECT_SUCCESS(); + dns_resolver_signal_.Notify(); + }, + "[2001:db8::1]:1", ""); + dns_resolver_signal_.WaitForNotification(); +} + +void TestIPv6WithoutPort(std::unique_ptr dns_resolver, + grpc_core::Notification* barrier, + absl::string_view target) { + dns_resolver->LookupHostname( + [barrier](auto result) { + EXPECT_TRUE(result.ok()); + EXPECT_FALSE(result->empty()); + barrier->Notify(); + }, + target, "80"); + barrier->WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, IPv6WithoutPortNoBrackets) { + TestIPv6WithoutPort(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "2001:db8::1"); +} + +TEST_F(EventEngineDNSTest, IPv6WithoutPortWithBrackets) { + TestIPv6WithoutPort(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "[2001:db8::1]"); +} + +TEST_F(EventEngineDNSTest, IPv6WithoutPortV4MappedV6) { + TestIPv6WithoutPort(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "2001:db8::1.2.3.4"); +} + +void TestInvalidIPAddress( + std::unique_ptr dns_resolver, + grpc_core::Notification* barrier, absl::string_view target) { + dns_resolver->LookupHostname( + [barrier](auto result) { + EXPECT_FALSE(result.ok()); + barrier->Notify(); + }, + target, ""); + barrier->WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, InvalidIPv4Addresses) { + TestInvalidIPAddress(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "293.283.1238.3:1"); +} + +TEST_F(EventEngineDNSTest, InvalidIPv6Addresses) { + TestInvalidIPAddress(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "[2001:db8::11111]:1"); +} + +void TestUnparseableHostPort( + std::unique_ptr dns_resolver, + grpc_core::Notification* barrier, absl::string_view target) { + dns_resolver->LookupHostname( + [barrier](auto result) { + EXPECT_FALSE(result.ok()); + barrier->Notify(); + }, + target, "1"); + barrier->WaitForNotification(); +} + +TEST_F(EventEngineDNSTest, UnparseableHostPortsOnlyBracket) { + TestUnparseableHostPort(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "["); +} + +TEST_F(EventEngineDNSTest, UnparseableHostPortsMissingRightBracket) { + TestUnparseableHostPort(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "[::1"); +} + +TEST_F(EventEngineDNSTest, UnparseableHostPortsBadPort) { + TestUnparseableHostPort(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "[::1]bad"); +} + +TEST_F(EventEngineDNSTest, UnparseableHostPortsBadIPv6) { + TestUnparseableHostPort(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "[1.2.3.4]"); +} + +TEST_F(EventEngineDNSTest, UnparseableHostPortsBadLocalhost) { + TestUnparseableHostPort(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "[localhost]"); +} + +TEST_F(EventEngineDNSTest, UnparseableHostPortsBadLocalhostWithPort) { + TestUnparseableHostPort(CreateDNSResolverWithoutSpecifyingServer(), + &dns_resolver_signal_, "[localhost]:1"); +} +// END + +#endif // GPR_WINDOWS diff --git a/test/core/event_engine/test_suite/tests/dns_test_record_groups.yaml b/test/core/event_engine/test_suite/tests/dns_test_record_groups.yaml new file mode 100644 index 00000000000..834962ca5e0 --- /dev/null +++ b/test/core/event_engine/test_suite/tests/dns_test_record_groups.yaml @@ -0,0 +1,21 @@ +resolver_tests_common_zone_name: dns-test.event-engine. +resolver_component_tests: +- records: + ipv4-only-multi-target: + - {TTL: '2100', data: 1.2.3.4, type: A} + - {TTL: '2100', data: 1.2.3.5, type: A} + - {TTL: '2100', data: 1.2.3.6, type: A} + ipv6-only-multi-target: + - {TTL: '2100', data: '2607:f8b0:400a:801::1002', type: AAAA} + - {TTL: '2100', data: '2607:f8b0:400a:801::1003', type: AAAA} + - {TTL: '2100', data: '2607:f8b0:400a:801::1004', type: AAAA} + ipv6-loopback-preferred-target: + - {TTL: '2100', data: '2002::1111', type: AAAA} + - {TTL: '2100', data: '::1', type: AAAA} + _grpclb._tcp.srv-multi-target: + - {TTL: '2100', data: 0 0 1234 ipv4-only-multi-target, type: SRV} + - {TTL: '2100', data: 0 0 1234 ipv6-only-multi-target, type: SRV} + _grpc_config.simple-service: + - {TTL: '2100', data: 'grpc_config=[{"serviceConfig":{"loadBalancingPolicy":"round_robin","methodConfig":[{"name":[{"method":"Foo","service":"SimpleService"}],"waitForReady":true}]}}]', + type: TXT} + - {TTL: '2100', data: 'other_config=other config', type: TXT} diff --git a/test/core/http/httpcli_test.cc b/test/core/http/httpcli_test.cc index 4e20e925c79..519e7fd3aaf 100644 --- a/test/core/http/httpcli_test.cc +++ b/test/core/http/httpcli_test.cc @@ -241,7 +241,7 @@ TEST_F(HttpRequestTest, Post) { int g_fake_non_responsive_dns_server_port; -void InjectNonResponsiveDNSServer(ares_channel channel) { +void InjectNonResponsiveDNSServer(ares_channel* channel) { gpr_log(GPR_DEBUG, "Injecting broken nameserver list. Bad server address:|[::1]:%d|.", g_fake_non_responsive_dns_server_port); @@ -253,7 +253,8 @@ void InjectNonResponsiveDNSServer(ares_channel channel) { dns_server_addrs[0].tcp_port = g_fake_non_responsive_dns_server_port; dns_server_addrs[0].udp_port = g_fake_non_responsive_dns_server_port; dns_server_addrs[0].next = nullptr; - GPR_ASSERT(ares_set_servers_ports(channel, dns_server_addrs) == ARES_SUCCESS); + GPR_ASSERT(ares_set_servers_ports(*channel, dns_server_addrs) == + ARES_SUCCESS); } TEST_F(HttpRequestTest, CancelGetDuringDNSResolution) { @@ -263,7 +264,7 @@ TEST_F(HttpRequestTest, CancelGetDuringDNSResolution) { kWaitForClientToSendFirstBytes, grpc_core::testing::FakeUdpAndTcpServer::CloseSocketUponCloseFromPeer); g_fake_non_responsive_dns_server_port = fake_dns_server.port(); - void (*prev_test_only_inject_config)(ares_channel channel) = + void (*prev_test_only_inject_config)(ares_channel * channel) = grpc_ares_test_only_inject_config; grpc_ares_test_only_inject_config = InjectNonResponsiveDNSServer; // Run the same test on several threads in parallel to try to trigger races diff --git a/test/core/iomgr/resolve_address_test.cc b/test/core/iomgr/resolve_address_test.cc index 2e32d903bd4..4dd7e04f08b 100644 --- a/test/core/iomgr/resolve_address_test.cc +++ b/test/core/iomgr/resolve_address_test.cc @@ -176,7 +176,7 @@ class ResolveAddressTest : public ::testing::Test { grpc_pollset_set* pollset_set_; // the default value of grpc_ares_test_only_inject_config, which might // be modified during a test - void (*default_inject_config_)(ares_channel channel) = nullptr; + void (*default_inject_config_)(ares_channel* channel) = nullptr; }; } // namespace @@ -390,7 +390,7 @@ namespace { int g_fake_non_responsive_dns_server_port; -void InjectNonResponsiveDNSServer(ares_channel channel) { +void InjectNonResponsiveDNSServer(ares_channel* channel) { gpr_log(GPR_DEBUG, "Injecting broken nameserver list. Bad server address:|[::1]:%d|.", g_fake_non_responsive_dns_server_port); @@ -403,7 +403,7 @@ void InjectNonResponsiveDNSServer(ares_channel channel) { dns_server_addrs[0].tcp_port = g_fake_non_responsive_dns_server_port; dns_server_addrs[0].udp_port = g_fake_non_responsive_dns_server_port; dns_server_addrs[0].next = nullptr; - ASSERT_EQ(ares_set_servers_ports(channel, dns_server_addrs), ARES_SUCCESS); + ASSERT_EQ(ares_set_servers_ports(*channel, dns_server_addrs), ARES_SUCCESS); } } // namespace diff --git a/test/cpp/naming/BUILD b/test/cpp/naming/BUILD index 32d1b462c9b..462c0c74e33 100644 --- a/test/cpp/naming/BUILD +++ b/test/cpp/naming/BUILD @@ -38,6 +38,7 @@ grpc_cc_test( name = "cancel_ares_query_test", srcs = ["cancel_ares_query_test.cc"], external_deps = ["gtest"], + tags = ["cancel_ares_query_test"], deps = [ "//:gpr", "//:grpc", diff --git a/test/cpp/naming/cancel_ares_query_test.cc b/test/cpp/naming/cancel_ares_query_test.cc index c2da75a13d1..1e2f6287941 100644 --- a/test/cpp/naming/cancel_ares_query_test.cc +++ b/test/cpp/naming/cancel_ares_query_test.cc @@ -39,6 +39,7 @@ #include "src/core/lib/debug/stats.h" #include "src/core/lib/debug/stats_data.h" #include "src/core/lib/event_engine/default_event_engine.h" +#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/gpr/string.h" #include "src/core/lib/gprpp/crash.h" #include "src/core/lib/gprpp/orphanable.h" @@ -54,6 +55,7 @@ #include "test/core/util/fake_udp_and_tcp_server.h" #include "test/core/util/port.h" #include "test/core/util/test_config.h" +#include "test/cpp/util/test_config.h" #ifdef GPR_WINDOWS #include "src/core/lib/iomgr/sockaddr_windows.h" @@ -310,8 +312,13 @@ void TestCancelDuringActiveQuery( // The DNS resolution timeout should fire well before the // RPC's deadline expires. expected_status_code = GRPC_STATUS_UNAVAILABLE; - expected_error_message_substring = - absl::StrCat("DNS resolution failed for ", name); + if (grpc_core::IsEventEngineDnsEnabled()) { + expected_error_message_substring = + absl::StrCat("errors resolving ", name); + } else { + expected_error_message_substring = + absl::StrCat("DNS resolution failed for ", name); + } grpc_arg arg; arg.type = GRPC_ARG_INTEGER; arg.key = const_cast(GRPC_ARG_DNS_ARES_QUERY_TIMEOUT_MS); @@ -424,8 +431,9 @@ TEST_F( } // namespace int main(int argc, char** argv) { - grpc::testing::TestEnvironment env(&argc, argv); ::testing::InitGoogleTest(&argc, argv); + grpc::testing::InitTest(&argc, &argv, true); + grpc::testing::TestEnvironment env(&argc, argv); auto result = RUN_ALL_TESTS(); return result; } diff --git a/test/cpp/naming/generate_resolver_component_tests.bzl b/test/cpp/naming/generate_resolver_component_tests.bzl index 626d9c0f46f..edd9c4e40e3 100755 --- a/test/cpp/naming/generate_resolver_component_tests.bzl +++ b/test/cpp/naming/generate_resolver_component_tests.bzl @@ -21,6 +21,10 @@ load("//bazel:grpc_build_system.bzl", "grpc_cc_binary", "grpc_cc_test") # buildifier: disable=unnamed-macro def generate_resolver_component_tests(): + """Generate address_sorting_test and resolver_component_test suite with different configurations. + + Note that the resolver_component_test suite's configuration is 2 dimensional: security and whether to enable the event_engine_dns experiment. + """ for unsecure_build_config_suffix in ["_unsecure", ""]: grpc_cc_test( name = "address_sorting_test%s" % unsecure_build_config_suffix, @@ -58,6 +62,7 @@ def generate_resolver_component_tests(): "//:grpc++%s" % unsecure_build_config_suffix, "//:grpc%s" % unsecure_build_config_suffix, "//:gpr", + "//src/core:ares_resolver", "//test/cpp/util:test_config", ], tags = ["no_windows"], @@ -84,7 +89,7 @@ def generate_resolver_component_tests(): "//test/cpp/naming/utils:dns_server", "//test/cpp/naming/utils:dns_resolver", "//test/cpp/naming/utils:tcp_connect", - "resolver_test_record_groups.yaml", # include the transitive dependency so that the dns sever py binary can locate this + "resolver_test_record_groups.yaml", # include the transitive dependency so that the dns server py binary can locate this ], args = [ "--test_bin_name=resolver_component_test%s" % unsecure_build_config_suffix, @@ -93,5 +98,5 @@ def generate_resolver_component_tests(): # The test is highly flaky on AWS workers that we use for running ARM64 tests. # The "no_arm64" tag can be used to skip it. # (see https://github.com/grpc/grpc/issues/25289). - tags = ["no_windows", "no_mac", "no_arm64"], + tags = ["no_windows", "no_mac", "no_arm64", "resolver_component_tests_runner_invoker"], ) diff --git a/test/cpp/naming/resolver_component_test.cc b/test/cpp/naming/resolver_component_test.cc index 1316143246d..7fe90b704e5 100644 --- a/test/cpp/naming/resolver_component_test.cc +++ b/test/cpp/naming/resolver_component_test.cc @@ -47,7 +47,9 @@ #include "src/core/lib/address_utils/sockaddr_utils.h" #include "src/core/lib/channel/channel_args.h" #include "src/core/lib/config/core_configuration.h" +#include "src/core/lib/event_engine/ares_resolver.h" #include "src/core/lib/event_engine/default_event_engine.h" +#include "src/core/lib/experiments/experiments.h" #include "src/core/lib/gpr/string.h" #include "src/core/lib/gprpp/crash.h" #include "src/core/lib/gprpp/host_port.h" @@ -255,11 +257,20 @@ void PollPollsetUntilRequestDone(ArgsStruct* args) { GPR_ASSERT(gpr_time_cmp(time_left, gpr_time_0(GPR_TIMESPAN)) >= 0); grpc_pollset_worker* worker = nullptr; grpc_core::ExecCtx exec_ctx; - GRPC_LOG_IF_ERROR( - "pollset_work", - grpc_pollset_work( - args->pollset, &worker, - grpc_core::Timestamp::FromTimespecRoundUp(NSecondDeadline(1)))); + if (grpc_core::IsEventEngineDnsEnabled()) { + // This essentially becomes a condition variable. + GRPC_LOG_IF_ERROR( + "pollset_work", + grpc_pollset_work( + args->pollset, &worker, + grpc_core::Timestamp::FromTimespecRoundUp(deadline))); + } else { + GRPC_LOG_IF_ERROR( + "pollset_work", + grpc_pollset_work( + args->pollset, &worker, + grpc_core::Timestamp::FromTimespecRoundUp(NSecondDeadline(1)))); + } } gpr_event_set(&args->ev, reinterpret_cast(1)); } @@ -529,7 +540,7 @@ int g_fake_non_responsive_dns_server_port = -1; // resolver. This is useful to effectively mock /etc/resolv.conf settings // (and equivalent on Windows), which unit tests don't have write permissions. // -void InjectBrokenNameServerList(ares_channel channel) { +void InjectBrokenNameServerList(ares_channel* channel) { struct ares_addr_port_node dns_server_addrs[2]; memset(dns_server_addrs, 0, sizeof(dns_server_addrs)); std::string unused_host; @@ -558,7 +569,8 @@ void InjectBrokenNameServerList(ares_channel channel) { dns_server_addrs[1].tcp_port = atoi(local_dns_server_port.c_str()); dns_server_addrs[1].udp_port = atoi(local_dns_server_port.c_str()); dns_server_addrs[1].next = nullptr; - GPR_ASSERT(ares_set_servers_ports(channel, dns_server_addrs) == ARES_SUCCESS); + GPR_ASSERT(ares_set_servers_ports(*channel, dns_server_addrs) == + ARES_SUCCESS); } void StartResolvingLocked(grpc_core::Resolver* r) { r->StartLocked(); } @@ -591,7 +603,12 @@ void RunResolvesRelevantRecordsTest( grpc_core::testing::FakeUdpAndTcpServer::CloseSocketUponCloseFromPeer); g_fake_non_responsive_dns_server_port = fake_non_responsive_dns_server->port(); - grpc_ares_test_only_inject_config = InjectBrokenNameServerList; + if (grpc_core::IsEventEngineDnsEnabled()) { + event_engine_grpc_ares_test_only_inject_config = + InjectBrokenNameServerList; + } else { + grpc_ares_test_only_inject_config = InjectBrokenNameServerList; + } whole_uri = absl::StrCat("dns:///", absl::GetFlag(FLAGS_target_name)); } else if (absl::GetFlag(FLAGS_inject_broken_nameserver_list) == "False") { gpr_log(GPR_INFO, "Specifying authority in uris to: %s", @@ -673,13 +690,15 @@ TEST(ResolverComponentTest, TestDoesntCrashOrHangWith1MsTimeout) { } // namespace int main(int argc, char** argv) { - grpc_init(); - grpc::testing::TestEnvironment env(&argc, argv); ::testing::InitGoogleTest(&argc, argv); + // Need before TestEnvironment construct for --grpc_experiments flag at + // least. grpc::testing::InitTest(&argc, &argv, true); + grpc::testing::TestEnvironment env(&argc, argv); if (absl::GetFlag(FLAGS_target_name).empty()) { grpc_core::Crash("Missing target_name param."); } + grpc_init(); auto result = RUN_ALL_TESTS(); grpc_shutdown(); return result; diff --git a/test/cpp/naming/resolver_component_tests_runner.py b/test/cpp/naming/resolver_component_tests_runner.py index 33118701d1a..cf648ac1614 100755 --- a/test/cpp/naming/resolver_component_tests_runner.py +++ b/test/cpp/naming/resolver_component_tests_runner.py @@ -1,4 +1,4 @@ -#!/usr/bin/env python +#!/usr/bin/env python3 # Copyright 2015 gRPC authors. # # Licensed under the Apache License, Version 2.0 (the "License"); @@ -58,6 +58,9 @@ if cur_resolver and cur_resolver != 'ares': test_runner_log('Exit 1 without running tests.') sys.exit(1) os.environ.update({'GRPC_TRACE': 'cares_resolver,cares_address_sorting'}) +experiments = os.environ.get('GRPC_EXPERIMENTS') +if experiments is not None and 'event_engine_dns' in experiments: + os.environ.update({'GRPC_TRACE': 'event_engine_client_channel_resolver,cares_resolver'}) def wait_until_dns_server_is_up(args, dns_server_subprocess, diff --git a/test/cpp/naming/utils/BUILD b/test/cpp/naming/utils/BUILD index 02b82b03392..6927a29bd1b 100644 --- a/test/cpp/naming/utils/BUILD +++ b/test/cpp/naming/utils/BUILD @@ -48,3 +48,9 @@ grpc_py_binary( testonly = True, srcs = ["tcp_connect.py"], ) + +grpc_py_binary( + name = "health_check", + testonly = True, + srcs = ["health_check.py"], +) diff --git a/test/cpp/naming/utils/dns_resolver.py b/test/cpp/naming/utils/dns_resolver.py index 91b01e856d1..923aec43e5f 100755 --- a/test/cpp/naming/utils/dns_resolver.py +++ b/test/cpp/naming/utils/dns_resolver.py @@ -1,4 +1,4 @@ -#!/usr/bin/env python2.7 +#!/usr/bin/env python3 # Copyright 2015 gRPC authors. # # Licensed under the Apache License, Version 2.0 (the "License"); diff --git a/test/cpp/naming/utils/dns_server.py b/test/cpp/naming/utils/dns_server.py index 02e1541875b..f9c118df7a1 100755 --- a/test/cpp/naming/utils/dns_server.py +++ b/test/cpp/naming/utils/dns_server.py @@ -1,4 +1,4 @@ -#!/usr/bin/env python2.7 +#!/usr/bin/env python3 # Copyright 2015 gRPC authors. # # Licensed under the Apache License, Version 2.0 (the "License"); diff --git a/test/cpp/naming/utils/health_check.py b/test/cpp/naming/utils/health_check.py new file mode 100644 index 00000000000..a7950c65f6b --- /dev/null +++ b/test/cpp/naming/utils/health_check.py @@ -0,0 +1,121 @@ +#!/usr/bin/env python3 +# Copyright 2023 The 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. + +import argparse +import platform +import subprocess +import sys +import time + + +def test_runner_log(msg): + sys.stderr.write("\n%s: %s\n" % (__file__, msg)) + + +def python_args(arg_list): + if platform.system() == "Windows": + return [sys.executable] + arg_list + return arg_list + + +def wait_until_dns_server_is_up(args): + for i in range(0, 30): + test_runner_log( + "Health check: attempt to connect to DNS server over TCP." + ) + tcp_connect_subprocess = subprocess.Popen( + python_args( + [ + args.tcp_connect_bin_path, + "--server_host", + "127.0.0.1", + "--server_port", + str(args.dns_server_port), + "--timeout", + str(1), + ] + ) + ) + tcp_connect_subprocess.communicate() + if tcp_connect_subprocess.returncode == 0: + test_runner_log( + ( + "Health check: attempt to make an A-record " + "query to DNS server." + ) + ) + dns_resolver_subprocess = subprocess.Popen( + python_args( + [ + args.dns_resolver_bin_path, + "--qname", + "health-check-local-dns-server-is-alive.resolver-tests.grpctestingexp", + "--server_host", + "127.0.0.1", + "--server_port", + str(args.dns_server_port), + ] + ), + stdout=subprocess.PIPE, + ) + dns_resolver_stdout, _ = dns_resolver_subprocess.communicate( + str.encode("ascii") + ) + if dns_resolver_subprocess.returncode == 0: + if "123.123.123.123".encode("ascii") in dns_resolver_stdout: + test_runner_log( + ( + "DNS server is up! " + "Successfully reached it over UDP and TCP." + ) + ) + return + time.sleep(1) + test_runner_log( + ( + "Failed to reach DNS server over TCP and/or UDP. " + "Exitting without running tests." + ) + ) + sys.exit(1) + + +def main(): + argp = argparse.ArgumentParser(description="Make DNS queries for A records") + argp.add_argument( + "-p", + "--dns_server_port", + default=None, + type=int, + help=("Port that local DNS server is listening on."), + ) + argp.add_argument( + "--dns_resolver_bin_path", + default=None, + type=str, + help=("Path to the DNS health check utility."), + ) + argp.add_argument( + "--tcp_connect_bin_path", + default=None, + type=str, + help=("Path to the TCP health check utility."), + ) + args = argp.parse_args() + wait_until_dns_server_is_up(args) + + +if __name__ == "__main__": + main() diff --git a/test/cpp/naming/utils/run_dns_server_for_lb_interop_tests.py b/test/cpp/naming/utils/run_dns_server_for_lb_interop_tests.py index 403b23dbde3..ccc06eaab88 100755 --- a/test/cpp/naming/utils/run_dns_server_for_lb_interop_tests.py +++ b/test/cpp/naming/utils/run_dns_server_for_lb_interop_tests.py @@ -1,4 +1,4 @@ -#!/usr/bin/env python2.7 +#!/usr/bin/env python3 # Copyright 2015 gRPC authors. # # Licensed under the Apache License, Version 2.0 (the "License"); diff --git a/test/cpp/naming/utils/tcp_connect.py b/test/cpp/naming/utils/tcp_connect.py index f06a5e4655e..41ce399febe 100755 --- a/test/cpp/naming/utils/tcp_connect.py +++ b/test/cpp/naming/utils/tcp_connect.py @@ -1,4 +1,4 @@ -#!/usr/bin/env python2.7 +#!/usr/bin/env python3 # Copyright 2015 gRPC authors. # # Licensed under the Apache License, Version 2.0 (the "License"); diff --git a/test/cpp/util/BUILD b/test/cpp/util/BUILD index 8438cf60794..d8aed93aa50 100644 --- a/test/cpp/util/BUILD +++ b/test/cpp/util/BUILD @@ -423,3 +423,19 @@ grpc_cc_test( ":test_util", ], ) + +grpc_cc_library( + name = "get_grpc_test_runfile_dir", + srcs = [ + "get_grpc_test_runfile_dir.cc", + ], + hdrs = [ + "get_grpc_test_runfile_dir.h", + ], + external_deps = [ + "absl/types:optional", + ], + deps = [ + "//src/core:env", + ], +) diff --git a/test/cpp/util/get_grpc_test_runfile_dir.cc b/test/cpp/util/get_grpc_test_runfile_dir.cc new file mode 100644 index 00000000000..4e0ce11ca5a --- /dev/null +++ b/test/cpp/util/get_grpc_test_runfile_dir.cc @@ -0,0 +1,29 @@ +// Copyright 2023 The 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. + +#include "test/cpp/util/get_grpc_test_runfile_dir.h" + +#include "src/core/lib/gprpp/env.h" + +namespace grpc { + +absl::optional GetGrpcTestRunFileDir() { + absl::optional test_srcdir = grpc_core::GetEnv("TEST_SRCDIR"); + if (!test_srcdir.has_value()) { + return absl::nullopt; + } + return *test_srcdir + "/com_github_grpc_grpc"; +} + +} // namespace grpc diff --git a/test/cpp/util/get_grpc_test_runfile_dir.h b/test/cpp/util/get_grpc_test_runfile_dir.h new file mode 100644 index 00000000000..885a44e0518 --- /dev/null +++ b/test/cpp/util/get_grpc_test_runfile_dir.h @@ -0,0 +1,32 @@ +// Copyright 2023 The 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_TEST_CPP_UTIL_GET_GRPC_TEST_RUNFILE_DIR_H +#define GRPC_TEST_CPP_UTIL_GET_GRPC_TEST_RUNFILE_DIR_H + +#include + +#include "absl/types/optional.h" + +namespace grpc { + +// Gets the absolute path of the runfile directory (a bazel/blaze concept) for a +// gRPC test. The path to the data files can be referred by joining the runfile +// directory with the workspace-relative path (e.g. +// "test/cpp/util/get_grpc_test_runfile_dir.h"). +absl::optional GetGrpcTestRunFileDir(); + +} // namespace grpc + +#endif // GRPC_TEST_CPP_UTIL_GET_GRPC_TEST_RUNFILE_DIR_H diff --git a/tools/doxygen/Doxyfile.c++.internal b/tools/doxygen/Doxyfile.c++.internal index 92d52f2675e..a9f3726fbfd 100644 --- a/tools/doxygen/Doxyfile.c++.internal +++ b/tools/doxygen/Doxyfile.c++.internal @@ -2060,6 +2060,8 @@ src/core/lib/debug/stats_data.cc \ src/core/lib/debug/stats_data.h \ src/core/lib/debug/trace.cc \ src/core/lib/debug/trace.h \ +src/core/lib/event_engine/ares_resolver.cc \ +src/core/lib/event_engine/ares_resolver.h \ src/core/lib/event_engine/cf_engine/cf_engine.cc \ src/core/lib/event_engine/cf_engine/cf_engine.h \ src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc \ @@ -2075,6 +2077,7 @@ src/core/lib/event_engine/default_event_engine_factory.h \ src/core/lib/event_engine/event_engine.cc \ src/core/lib/event_engine/forkable.cc \ src/core/lib/event_engine/forkable.h \ +src/core/lib/event_engine/grpc_polled_fd.h \ src/core/lib/event_engine/handle_containers.h \ src/core/lib/event_engine/memory_allocator.cc \ src/core/lib/event_engine/memory_allocator_factory.h \ @@ -2087,6 +2090,7 @@ src/core/lib/event_engine/posix_engine/ev_poll_posix.h \ src/core/lib/event_engine/posix_engine/event_poller.h \ src/core/lib/event_engine/posix_engine/event_poller_posix_default.cc \ src/core/lib/event_engine/posix_engine/event_poller_posix_default.h \ +src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h \ src/core/lib/event_engine/posix_engine/internal_errqueue.cc \ src/core/lib/event_engine/posix_engine/internal_errqueue.h \ src/core/lib/event_engine/posix_engine/lockfree_event.cc \ diff --git a/tools/doxygen/Doxyfile.core.internal b/tools/doxygen/Doxyfile.core.internal index 2eae3f09ed2..e1778392a62 100644 --- a/tools/doxygen/Doxyfile.core.internal +++ b/tools/doxygen/Doxyfile.core.internal @@ -1838,6 +1838,8 @@ src/core/lib/debug/stats_data.cc \ src/core/lib/debug/stats_data.h \ src/core/lib/debug/trace.cc \ src/core/lib/debug/trace.h \ +src/core/lib/event_engine/ares_resolver.cc \ +src/core/lib/event_engine/ares_resolver.h \ src/core/lib/event_engine/cf_engine/cf_engine.cc \ src/core/lib/event_engine/cf_engine/cf_engine.h \ src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc \ @@ -1853,6 +1855,7 @@ src/core/lib/event_engine/default_event_engine_factory.h \ src/core/lib/event_engine/event_engine.cc \ src/core/lib/event_engine/forkable.cc \ src/core/lib/event_engine/forkable.h \ +src/core/lib/event_engine/grpc_polled_fd.h \ src/core/lib/event_engine/handle_containers.h \ src/core/lib/event_engine/memory_allocator.cc \ src/core/lib/event_engine/memory_allocator_factory.h \ @@ -1865,6 +1868,7 @@ src/core/lib/event_engine/posix_engine/ev_poll_posix.h \ src/core/lib/event_engine/posix_engine/event_poller.h \ src/core/lib/event_engine/posix_engine/event_poller_posix_default.cc \ src/core/lib/event_engine/posix_engine/event_poller_posix_default.h \ +src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h \ src/core/lib/event_engine/posix_engine/internal_errqueue.cc \ src/core/lib/event_engine/posix_engine/internal_errqueue.h \ src/core/lib/event_engine/posix_engine/lockfree_event.cc \