[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.
<!--

If you know who should review your pull request, please assign it to
that
person, otherwise the pull request would get assigned randomly.

If your pull request is for a specific language, please add the
appropriate
lang label.

-->
This commit is contained in:
Yijie Ma 2023-07-21 13:24:16 -07:00 committed by GitHub
parent 04a7b80e98
commit a7bf07e86a
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
50 changed files with 2097 additions and 56 deletions

11
CMakeLists.txt generated
View File

@ -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
)

2
Makefile generated
View File

@ -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 \

4
Package.swift generated
View File

@ -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",

6
bazel/experiments.bzl generated
View File

@ -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",

View File

@ -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

1
config.m4 generated
View File

@ -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 \

1
config.w32 generated
View File

@ -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 " +

6
gRPC-C++.podspec generated
View File

@ -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',

7
gRPC-Core.podspec generated
View File

@ -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',

4
grpc.gemspec generated
View File

@ -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 )

3
grpc.gyp generated
View File

@ -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',

View File

@ -365,7 +365,6 @@ class EventEngine : public std::enable_shared_from_this<EventEngine> {
/// 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;

4
package.xml generated
View File

@ -1045,6 +1045,8 @@
<file baseinstalldir="/" name="src/core/lib/debug/stats_data.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/debug/trace.cc" role="src" />
<file baseinstalldir="/" name="src/core/lib/debug/trace.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/ares_resolver.cc" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/ares_resolver.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/cf_engine/cf_engine.cc" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/cf_engine/cf_engine.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/cf_engine/cfstream_endpoint.cc" role="src" />
@ -1060,6 +1062,7 @@
<file baseinstalldir="/" name="src/core/lib/event_engine/event_engine.cc" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/forkable.cc" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/forkable.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/grpc_polled_fd.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/handle_containers.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/memory_allocator.cc" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/memory_allocator_factory.h" role="src" />
@ -1072,6 +1075,7 @@
<file baseinstalldir="/" name="src/core/lib/event_engine/posix_engine/event_poller.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/posix_engine/event_poller_posix_default.cc" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/posix_engine/event_poller_posix_default.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/posix_engine/grpc_polled_fd_posix.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/posix_engine/internal_errqueue.cc" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/posix_engine/internal_errqueue.h" role="src" />
<file baseinstalldir="/" name="src/core/lib/event_engine/posix_engine/lockfree_event.cc" role="src" />

View File

@ -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",

View File

@ -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(

View File

@ -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

View File

@ -22,7 +22,6 @@
#include <chrono>
#include <memory>
#include <string>
#include <type_traits>
#include <utility>
#include <vector>
@ -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<std::vector<EventEngine::ResolvedAddress>> addresses) {
absl::StatusOr<std::vector<EventEngine::ResolvedAddress>>
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<std::vector<EventEngine::DNSResolver::SRVRecord>>
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<std::vector<std::string>> service_config) {
absl::StatusOr<std::vector<std::string>> 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<std::vector<EventEngine::ResolvedAddress>>
new_addresses) {
ValidationErrors::ScopedField field(&errors_, "hostname lookup");
absl::optional<Resolver::Result> 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<std::vector<EventEngine::DNSResolver::SRVRecord>>
srv_records) {
ValidationErrors::ScopedField field(&errors_, "srv lookup");
absl::optional<Resolver::Result> 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<std::vector<EventEngine::ResolvedAddress>>
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<std::vector<EventEngine::ResolvedAddress>>
new_balancer_addresses) {
ValidationErrors::ScopedField field(
&errors_, absl::StrCat("balancer lookup for ", authority));
absl::optional<Resolver::Result> 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<std::vector<std::string>> service_config) {
ValidationErrors::ScopedField field(&errors_, "txt lookup");
absl::optional<Resolver::Result> 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,

View File

@ -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) {
? "<null>"
: 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<PollingResolver> self =

View File

@ -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 <grpc/support/port_platform.h>
#include "src/core/lib/event_engine/ares_resolver.h"
#include <stdint.h>
#include <string>
#include <vector>
#include "src/core/lib/iomgr/port.h"
// IWYU pragma: no_include <arpa/inet.h>
// IWYU pragma: no_include <arpa/nameser.h>
// IWYU pragma: no_include <inttypes.h>
// IWYU pragma: no_include <netdb.h>
// IWYU pragma: no_include <netinet/in.h>
// IWYU pragma: no_include <stdlib.h>
// IWYU pragma: no_include <sys/socket.h>
// IWYU pragma: no_include <ratio>
#if GRPC_ARES == 1
#include <ares_nameser.h>
#include <string.h>
#include <algorithm>
#include <chrono>
#include <initializer_list>
#include <memory>
#include <type_traits>
#include <utility>
#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 <grpc/event_engine/event_engine.h>
#include <grpc/support/log.h>
#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<struct sockaddr_in*>(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<struct sockaddr_in6*>(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<grpc_core::OrphanablePtr<AresResolver>>
AresResolver::CreateAresResolver(
absl::string_view dns_server,
std::unique_ptr<GrpcPolledFdFactory> polled_fd_factory,
std::shared_ptr<EventEngine> 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<AresResolver>(
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<EventEngine::ResolvedAddress> result;
result.emplace_back(reinterpret_cast<sockaddr*>(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<EventEngine::DNSResolver::SRVRecord>());
});
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<std::string>());
});
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<GrpcPolledFdFactory> polled_fd_factory,
std::shared_ptr<EventEngine> event_engine, ares_channel channel)
: grpc_core::InternallyRefCounted<AresResolver>(
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<FdNode>(
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<HostnameQueryArg> hostname_qa(
static_cast<HostnameQueryArg*>(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<EventEngine::DNSResolver::LookupHostnameCallback>(
nh.mapped()));
auto callback = absl::get<EventEngine::DNSResolver::LookupHostnameCallback>(
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<EventEngine::ResolvedAddress> 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<unsigned char>(hostent->h_addrtype);
addr.sin6_port = htons(hostname_qa->port);
result.emplace_back(reinterpret_cast<const sockaddr*>(&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<unsigned char>(hostent->h_addrtype);
addr.sin_port = htons(hostname_qa->port);
result.emplace_back(reinterpret_cast<const sockaddr*>(&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<QueryArg> qa(static_cast<QueryArg*>(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<EventEngine::DNSResolver::LookupSRVCallback>(
nh.mapped()));
auto callback = absl::get<EventEngine::DNSResolver::LookupSRVCallback>(
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<EventEngine::DNSResolver::SRVRecord> 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<QueryArg> qa(static_cast<QueryArg*>(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<EventEngine::DNSResolver::LookupTXTCallback>(
nh.mapped()));
auto callback = absl::get<EventEngine::DNSResolver::LookupTXTCallback>(
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<std::string> result;
for (struct ares_txt_ext* part = reply; part != nullptr; part = part->next) {
if (part->record_start) {
result.emplace_back(reinterpret_cast<char*>(part->txt), part->length);
} else {
absl::StrAppend(
&result.back(),
std::string(reinterpret_cast<char*>(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

View File

@ -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 <grpc/support/port_platform.h>
#include "src/core/lib/debug/trace.h"
#if GRPC_ARES == 1
#include <list>
#include <memory>
#include <ares.h>
#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 <grpc/event_engine/event_engine.h>
#include <grpc/support/log.h>
#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<AresResolver> {
public:
static absl::StatusOr<grpc_core::OrphanablePtr<AresResolver>>
CreateAresResolver(absl::string_view dns_server,
std::unique_ptr<GrpcPolledFdFactory> polled_fd_factory,
std::shared_ptr<EventEngine> event_engine);
// Do not instantiate directly -- use CreateAresResolver() instead.
AresResolver(std::unique_ptr<GrpcPolledFdFactory> polled_fd_factory,
std::shared_ptr<EventEngine> 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<GrpcPolledFd> 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<std::unique_ptr<FdNode>>;
using CallbackType =
absl::variant<EventEngine::DNSResolver::LookupHostnameCallback,
EventEngine::DNSResolver::LookupSRVCallback,
EventEngine::DNSResolver::LookupTXTCallback>;
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<int, CallbackType> callback_map_ ABSL_GUARDED_BY(mutex_);
absl::optional<EventEngine::TaskHandle> ares_backup_poll_alarm_handle_
ABSL_GUARDED_BY(mutex_);
std::unique_ptr<GrpcPolledFdFactory> polled_fd_factory_;
std::shared_ptr<EventEngine> 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

View File

@ -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 <grpc/support/port_platform.h>
#if GRPC_ARES == 1
#include <ares.h>
#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<void(absl::Status)> read_closure) = 0;
// Called when c-ares library is interested and there's no pending callback
virtual void RegisterForOnWriteableLocked(
absl::AnyInvocable<void(absl::Status)> 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

View File

@ -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 <grpc/support/port_platform.h>
#include "src/core/lib/iomgr/port.h"
#if GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_ARES_EV_DRIVER)
#include <string.h>
#include <sys/ioctl.h>
#include <string>
#include <utility>
#include <ares.h>
#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<int>(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<void(absl::Status)> read_closure) override {
handle_->NotifyOnRead(new PosixEngineClosure(std::move(read_closure),
/*is_permanent=*/false));
}
void RegisterForOnWriteableLocked(
absl::AnyInvocable<void(absl::Status)> 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

View File

@ -18,6 +18,7 @@
#include <algorithm>
#include <atomic>
#include <chrono>
#include <cstdint>
#include <cstring>
#include <memory>
#include <string>
@ -37,8 +38,10 @@
#include <grpc/support/log.h>
#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<AresResolver> 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<std::unique_ptr<EventEngine::DNSResolver>>
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<GrpcPolledFdFactoryPosix>(poller_manager_->Poller()),
shared_from_this());
if (!ares_resolver.ok()) {
return ares_resolver.status();
}
return std::make_unique<PosixEventEngine::PosixDNSResolver>(
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"); }

View File

@ -34,11 +34,13 @@
#include <grpc/event_engine/event_engine.h>
#include <grpc/event_engine/memory_allocator.h>
#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<AresResolver> 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<AresResolver> ares_resolver_;
#endif // GRPC_ARES == 1 && defined(GRPC_POSIX_SOCKET_TCP)
};
#ifdef GRPC_POSIX_SOCKET_TCP

View File

@ -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:

View File

@ -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',

View File

@ -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', []):

View File

@ -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,

View File

@ -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",
],

View File

@ -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,
)

View File

@ -12,10 +12,36 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#include <gtest/gtest.h>
// IWYU pragma: no_include <ratio>
// IWYU pragma: no_include <arpa/inet.h>
#include "src/core/lib/iomgr/exec_ctx.h"
#include <cstdlib>
#include <cstring>
#include <initializer_list>
#include <memory>
#include <string>
#include <tuple>
#include <vector>
#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 <grpc/event_engine/event_engine.h>
#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<std::string> 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();
// <path to dns_server.py> -p <port> -r <path to records config>
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<EventEngine::DNSResolver> CreateDefaultDNSResolver() {
std::shared_ptr<EventEngine> test_ee(this->NewEventEngine());
EventEngine::DNSResolver::ResolverOptions options;
options.dns_server = dns_server_.address();
return *test_ee->GetDNSResolver(options);
}
std::unique_ptr<EventEngine::DNSResolver>
CreateDNSResolverWithNonResponsiveServer() {
using FakeUdpAndTcpServer = grpc_core::testing::FakeUdpAndTcpServer;
// Start up fake non responsive DNS server
fake_dns_server_ = std::make_unique<FakeUdpAndTcpServer>(
FakeUdpAndTcpServer::AcceptMode::kWaitForClientToSendFirstBytes,
FakeUdpAndTcpServer::CloseSocketUponCloseFromPeer);
const std::string dns_server =
absl::StrFormat("[::1]:%d", fake_dns_server_->port());
std::shared_ptr<EventEngine> test_ee(this->NewEventEngine());
EventEngine::DNSResolver::ResolverOptions options;
options.dns_server = dns_server;
return *test_ee->GetDNSResolver(options);
}
std::unique_ptr<EventEngine::DNSResolver>
CreateDNSResolverWithoutSpecifyingServer() {
std::shared_ptr<EventEngine> 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<grpc_core::testing::FakeUdpAndTcpServer> 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<EventEngine::DNSResolver> 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<EventEngine::DNSResolver> 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<EventEngine::DNSResolver> 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

View File

@ -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}

View File

@ -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

View File

@ -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

View File

@ -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",

View File

@ -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<char*>(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;
}

View File

@ -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"],
)

View File

@ -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<void*>(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;

View File

@ -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,

View File

@ -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"],
)

View File

@ -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");

View File

@ -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");

View File

@ -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()

View File

@ -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");

View File

@ -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");

View File

@ -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",
],
)

View File

@ -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<std::string> GetGrpcTestRunFileDir() {
absl::optional<std::string> test_srcdir = grpc_core::GetEnv("TEST_SRCDIR");
if (!test_srcdir.has_value()) {
return absl::nullopt;
}
return *test_srcdir + "/com_github_grpc_grpc";
}
} // namespace grpc

View File

@ -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 <string>
#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<std::string> GetGrpcTestRunFileDir();
} // namespace grpc
#endif // GRPC_TEST_CPP_UTIL_GET_GRPC_TEST_RUNFILE_DIR_H

View File

@ -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 \

View File

@ -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 \