diff --git a/CMakeLists.txt b/CMakeLists.txt index 04f58d2b97551..b450360640312 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -209,6 +209,7 @@ set(gRPC_ABSL_USED_TARGETS absl_nullability absl_numeric_representation absl_optional + absl_overload absl_prefetch absl_random_bit_gen_ref absl_random_distributions @@ -5065,7 +5066,9 @@ target_link_libraries(grpc++ ${_gRPC_ALLTARGETS_LIBRARIES} absl::dynamic_annotations absl::prefetch + absl::fixed_array absl::layout + absl::overload absl::absl_check absl::absl_log absl::charset @@ -5813,7 +5816,9 @@ target_link_libraries(grpc++_unsecure ${_gRPC_ALLTARGETS_LIBRARIES} absl::dynamic_annotations absl::prefetch + absl::fixed_array absl::layout + absl::overload absl::absl_check absl::absl_log absl::charset @@ -12076,8 +12081,6 @@ if(_gRPC_PLATFORM_LINUX OR _gRPC_PLATFORM_MAC OR _gRPC_PLATFORM_POSIX) target_link_libraries(channelz_tool_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload grpcpp_channelz ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc_test_util @@ -12489,8 +12492,6 @@ target_include_directories(cli_call_test target_link_libraries(cli_call_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc++_test_util ) @@ -19836,8 +19837,6 @@ target_include_directories(grpc_cli target_link_libraries(grpc_cli ${_gRPC_ALLTARGETS_LIBRARIES} - absl::fixed_array - absl::overload grpc++ ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc++_test_config @@ -20622,8 +20621,6 @@ if(_gRPC_PLATFORM_LINUX OR _gRPC_PLATFORM_POSIX) target_link_libraries(grpc_tool_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload grpc++_reflection ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc++_test_config @@ -23176,11 +23173,15 @@ add_executable(interop_client ${_gRPC_PROTO_GENS_DIR}/src/proto/grpc/testing/test.grpc.pb.cc ${_gRPC_PROTO_GENS_DIR}/src/proto/grpc/testing/test.pb.h ${_gRPC_PROTO_GENS_DIR}/src/proto/grpc/testing/test.grpc.pb.h + src/cpp/ext/otel/otel_client_call_tracer.cc + src/cpp/ext/otel/otel_plugin.cc + src/cpp/ext/otel/otel_server_call_tracer.cc test/core/credentials/call/oauth2/oauth2_utils.cc test/cpp/interop/backend_metrics_lb_policy.cc test/cpp/interop/client.cc test/cpp/interop/client_helper.cc test/cpp/interop/interop_client.cc + test/cpp/interop/otel_helper.cc ) if(WIN32 AND MSVC) if(BUILD_SHARED_LIBS) @@ -23210,6 +23211,8 @@ target_include_directories(interop_client target_link_libraries(interop_client ${_gRPC_ALLTARGETS_LIBRARIES} + opentelemetry-cpp::api + opentelemetry-cpp::trace grpc++_test_config grpc++_test_util ) @@ -23231,9 +23234,13 @@ add_executable(interop_server ${_gRPC_PROTO_GENS_DIR}/src/proto/grpc/testing/test.grpc.pb.cc ${_gRPC_PROTO_GENS_DIR}/src/proto/grpc/testing/test.pb.h ${_gRPC_PROTO_GENS_DIR}/src/proto/grpc/testing/test.grpc.pb.h + src/cpp/ext/otel/otel_client_call_tracer.cc + src/cpp/ext/otel/otel_plugin.cc + src/cpp/ext/otel/otel_server_call_tracer.cc src/cpp/server/orca/orca_service.cc test/cpp/interop/interop_server.cc test/cpp/interop/interop_server_bootstrap.cc + test/cpp/interop/otel_helper.cc test/cpp/interop/server_helper.cc ) if(WIN32 AND MSVC) @@ -23264,6 +23271,8 @@ target_include_directories(interop_server target_link_libraries(interop_server ${_gRPC_ALLTARGETS_LIBRARIES} + opentelemetry-cpp::api + opentelemetry-cpp::trace grpc++_test_config grpc++_test_util ) @@ -24184,8 +24193,6 @@ if(_gRPC_PLATFORM_LINUX OR _gRPC_PLATFORM_MAC OR _gRPC_PLATFORM_POSIX) target_link_libraries(latent_see_tool_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload grpc++ ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc_test_util @@ -34140,7 +34147,9 @@ target_link_libraries(status_test gtest absl::dynamic_annotations absl::prefetch + absl::fixed_array absl::layout + absl::overload absl::absl_check absl::absl_log absl::charset @@ -39745,8 +39754,6 @@ target_include_directories(xds_audit_logger_registry_test target_link_libraries(xds_audit_logger_registry_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload grpc++ ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc_test_util @@ -40469,7 +40476,9 @@ target_link_libraries(xds_client_test gtest absl::dynamic_annotations absl::prefetch + absl::fixed_array absl::layout + absl::overload absl::absl_check absl::absl_log absl::charset @@ -41553,7 +41562,9 @@ target_link_libraries(xds_cluster_resource_type_test gtest absl::dynamic_annotations absl::prefetch + absl::fixed_array absl::layout + absl::overload absl::absl_check absl::absl_log absl::charset @@ -42568,8 +42579,6 @@ target_include_directories(xds_common_types_test target_link_libraries(xds_common_types_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload grpc++ ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc_test_util @@ -45299,7 +45308,9 @@ target_link_libraries(xds_endpoint_resource_type_test gtest absl::dynamic_annotations absl::prefetch + absl::fixed_array absl::layout + absl::overload absl::absl_check absl::absl_log absl::charset @@ -48017,8 +48028,6 @@ target_include_directories(xds_http_filters_test target_link_libraries(xds_http_filters_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload grpc++ ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc_test_util @@ -48469,8 +48478,6 @@ target_include_directories(xds_lb_policy_registry_test target_link_libraries(xds_lb_policy_registry_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload grpc++ ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc_test_util @@ -48989,8 +48996,6 @@ target_include_directories(xds_listener_resource_type_test target_link_libraries(xds_listener_resource_type_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload grpc++ ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc_test_util @@ -53047,8 +53052,6 @@ target_include_directories(xds_route_config_resource_type_test target_link_libraries(xds_route_config_resource_type_test ${_gRPC_ALLTARGETS_LIBRARIES} gtest - absl::fixed_array - absl::overload grpc++ ${_gRPC_PROTOBUF_PROTOC_LIBRARIES} grpc_test_util @@ -56157,7 +56160,7 @@ generate_pkgconfig( "gRPC++" "C++ wrapper for gRPC" "${gRPC_CPP_VERSION}" - "absl_absl_check absl_absl_log absl_algorithm_container absl_any_invocable absl_base absl_bind_front absl_bits absl_btree absl_charset absl_check absl_cleanup absl_config absl_cord absl_core_headers absl_dynamic_annotations absl_flags absl_flags_marshalling absl_flat_hash_map absl_flat_hash_set absl_function_ref absl_hash absl_inlined_vector absl_layout absl_log absl_log_globals absl_log_severity absl_memory absl_no_destructor absl_node_hash_map absl_optional absl_prefetch absl_random_bit_gen_ref absl_random_distributions absl_random_random absl_span absl_status absl_statusor absl_str_format absl_string_view absl_strings absl_strings_internal absl_synchronization absl_time absl_type_traits absl_utility gpr grpc" + "absl_absl_check absl_absl_log absl_algorithm_container absl_any_invocable absl_base absl_bind_front absl_bits absl_btree absl_charset absl_check absl_cleanup absl_config absl_cord absl_core_headers absl_dynamic_annotations absl_fixed_array absl_flags absl_flags_marshalling absl_flat_hash_map absl_flat_hash_set absl_function_ref absl_hash absl_inlined_vector absl_layout absl_log absl_log_globals absl_log_severity absl_memory absl_no_destructor absl_node_hash_map absl_optional absl_overload absl_prefetch absl_random_bit_gen_ref absl_random_distributions absl_random_random absl_span absl_status absl_statusor absl_str_format absl_string_view absl_strings absl_strings_internal absl_synchronization absl_time absl_type_traits absl_utility gpr grpc" "libcares openssl re2 zlib" "-lgrpc++" "-laddress_sorting -lupb_textformat_lib -lupb_json_lib -lupb_reflection_lib -lupb_wire_lib -lupb_message_lib -lutf8_range_lib -lupb_descriptor_lib -lupb_mini_descriptor_lib -lupb_mini_table_lib -lupb_hash_lib -lupb_mem_lib -lupb_base_lib -lupb_lex_lib" @@ -56168,7 +56171,7 @@ generate_pkgconfig( "gRPC++ unsecure" "C++ wrapper for gRPC without SSL" "${gRPC_CPP_VERSION}" - "absl_absl_check absl_absl_log absl_algorithm_container absl_any_invocable absl_base absl_bind_front absl_bits absl_btree absl_charset absl_check absl_cleanup absl_config absl_cord absl_core_headers absl_dynamic_annotations absl_flags absl_flags_marshalling absl_flat_hash_map absl_flat_hash_set absl_function_ref absl_hash absl_inlined_vector absl_layout absl_log absl_log_globals absl_log_severity absl_memory absl_no_destructor absl_node_hash_map absl_optional absl_prefetch absl_random_bit_gen_ref absl_random_distributions absl_random_random absl_span absl_status absl_statusor absl_str_format absl_string_view absl_strings absl_strings_internal absl_synchronization absl_time absl_type_traits absl_utility gpr grpc_unsecure" + "absl_absl_check absl_absl_log absl_algorithm_container absl_any_invocable absl_base absl_bind_front absl_bits absl_btree absl_charset absl_check absl_cleanup absl_config absl_cord absl_core_headers absl_dynamic_annotations absl_fixed_array absl_flags absl_flags_marshalling absl_flat_hash_map absl_flat_hash_set absl_function_ref absl_hash absl_inlined_vector absl_layout absl_log absl_log_globals absl_log_severity absl_memory absl_no_destructor absl_node_hash_map absl_optional absl_overload absl_prefetch absl_random_bit_gen_ref absl_random_distributions absl_random_random absl_span absl_status absl_statusor absl_str_format absl_string_view absl_strings absl_strings_internal absl_synchronization absl_time absl_type_traits absl_utility gpr grpc_unsecure" "libcares zlib" "-lgrpc++_unsecure" "-laddress_sorting -lupb_textformat_lib -lupb_reflection_lib -lupb_wire_lib -lupb_message_lib -lutf8_range_lib -lupb_descriptor_lib -lupb_mini_descriptor_lib -lupb_mini_table_lib -lupb_hash_lib -lupb_mem_lib -lupb_base_lib -lupb_lex_lib" @@ -56180,7 +56183,7 @@ if(gRPC_BUILD_GRPCPP_OTEL_PLUGIN) "gRPC++ OpenTelemetry Plugin" "OpenTelemetry Plugin for gRPC C++" "${gRPC_CPP_VERSION}" - "absl_absl_check absl_absl_log absl_algorithm_container absl_any_invocable absl_base absl_bind_front absl_bits absl_btree absl_charset absl_check absl_cleanup absl_config absl_cord absl_core_headers absl_dynamic_annotations absl_flags absl_flags_marshalling absl_flat_hash_map absl_flat_hash_set absl_function_ref absl_hash absl_inlined_vector absl_layout absl_log absl_log_globals absl_log_severity absl_memory absl_no_destructor absl_node_hash_map absl_optional absl_prefetch absl_random_bit_gen_ref absl_random_distributions absl_random_random absl_span absl_status absl_statusor absl_str_format absl_string_view absl_strings absl_strings_internal absl_synchronization absl_time absl_type_traits absl_utility gpr grpc grpc++ opentelemetry_api" + "absl_absl_check absl_absl_log absl_algorithm_container absl_any_invocable absl_base absl_bind_front absl_bits absl_btree absl_charset absl_check absl_cleanup absl_config absl_cord absl_core_headers absl_dynamic_annotations absl_fixed_array absl_flags absl_flags_marshalling absl_flat_hash_map absl_flat_hash_set absl_function_ref absl_hash absl_inlined_vector absl_layout absl_log absl_log_globals absl_log_severity absl_memory absl_no_destructor absl_node_hash_map absl_optional absl_overload absl_prefetch absl_random_bit_gen_ref absl_random_distributions absl_random_random absl_span absl_status absl_statusor absl_str_format absl_string_view absl_strings absl_strings_internal absl_synchronization absl_time absl_type_traits absl_utility gpr grpc grpc++ opentelemetry_api" "libcares openssl re2 zlib" "-lgrpcpp_otel_plugin" "-laddress_sorting -lupb_textformat_lib -lupb_json_lib -lupb_reflection_lib -lupb_wire_lib -lupb_message_lib -lutf8_range_lib -lupb_descriptor_lib -lupb_mini_descriptor_lib -lupb_mini_table_lib -lupb_hash_lib -lupb_mem_lib -lupb_base_lib -lupb_lex_lib" diff --git a/MODULE.bazel b/MODULE.bazel index bbb0219f4cac3..73154141a9807 100644 --- a/MODULE.bazel +++ b/MODULE.bazel @@ -51,6 +51,7 @@ bazel_dep(name = "googletest", version = "1.17.0", repo_name = "com_google_googl bazel_dep(name = "opencensus-cpp", version = "0.0.0-20230502-50eb5de.bcr.2", repo_name = "io_opencensus_cpp") bazel_dep(name = "openssl", version = "3.3.1.bcr.1") bazel_dep(name = "opentelemetry-cpp", version = "1.19.0", repo_name = "io_opentelemetry_cpp") +bazel_dep(name = "opentelemetry-proto", version = "1.8.0", repo_name = "opentelemetry_proto") # --- Protobuf related packages. bazel_dep(name = "protobuf", version = "35.1", repo_name = "com_google_protobuf") diff --git a/bazel/grpc_build_system.bzl b/bazel/grpc_build_system.bzl index b1df59866f298..1ae3b00c9090d 100644 --- a/bazel/grpc_build_system.bzl +++ b/bazel/grpc_build_system.bzl @@ -82,7 +82,7 @@ def _get_external_deps(external_deps): elif dep == "protobuf_clib": ret.extend(["@com_google_protobuf//src/google/protobuf/compiler:code_generator", "@com_google_protobuf//src/google/protobuf/compiler:importer"]) elif dep == "protobuf_headers": - ret.extend(["@com_google_protobuf//:protobuf_headers", "@com_google_protobuf//src/google/protobuf/io", "@com_google_protobuf//src/google/protobuf/io:printer", "@com_google_protobuf//src/google/protobuf/io:tokenizer"]) + ret.extend(["@com_google_protobuf//:protobuf_headers", "@com_google_protobuf//src/google/protobuf/io", "@com_google_protobuf//src/google/protobuf/io:printer", "@com_google_protobuf//src/google/protobuf/io:tokenizer", "@com_google_protobuf//src/google/protobuf/json:json"]) elif dep.startswith("absl/"): ret.append("@com_google_absl//" + dep) elif dep.startswith("google/"): diff --git a/bazel/grpc_deps.bzl b/bazel/grpc_deps.bzl index 6a6939494dcc8..ee49d9f20e148 100644 --- a/bazel/grpc_deps.bzl +++ b/bazel/grpc_deps.bzl @@ -381,6 +381,17 @@ def grpc_deps(): ], ) + if "opentelemetry_proto" not in native.existing_rules(): + http_archive( + name = "opentelemetry_proto", + sha256 = "07e26be7cc16630f5cc8910076a086f6858288fb9868ef648d08cb5cf876d756", + strip_prefix = "opentelemetry-proto-1.8.0", + urls = [ + "https://storage.googleapis.com/grpc-bazel-mirror/github.com/open-telemetry/opentelemetry-proto/archive/refs/tags/v1.8.0.tar.gz", + "https://github.com/open-telemetry/opentelemetry-proto/archive/refs/tags/v1.8.0.tar.gz", + ], + ) + if "grpc_proto" not in native.existing_rules(): http_archive( name = "grpc_proto", diff --git a/build_autogenerated.yaml b/build_autogenerated.yaml index bd93da8d1b8f4..6940e2691c620 100644 --- a/build_autogenerated.yaml +++ b/build_autogenerated.yaml @@ -4250,6 +4250,7 @@ libs: - src/cpp/server/secure_server_credentials.h - src/cpp/server/thread_pool_interface.h - src/cpp/thread_manager/thread_manager.h + - third_party/utf8_range/utf8_validity.h src: - src/core/client_channel/virtual_channel.cc - src/cpp/client/call_context_registry.cc @@ -4306,7 +4307,9 @@ libs: deps: - absl/base:dynamic_annotations - absl/base:prefetch + - absl/container:fixed_array - absl/container:layout + - absl/functional:overload - absl/log:absl_check - absl/log:absl_log - absl/strings:charset @@ -4633,6 +4636,7 @@ libs: - src/cpp/server/health/default_health_check_service.h - src/cpp/server/thread_pool_interface.h - src/cpp/thread_manager/thread_manager.h + - third_party/utf8_range/utf8_validity.h src: - src/core/client_channel/virtual_channel.cc - src/cpp/client/call_context_registry.cc @@ -4679,7 +4683,9 @@ libs: deps: - absl/base:dynamic_annotations - absl/base:prefetch + - absl/container:fixed_array - absl/container:layout + - absl/functional:overload - absl/log:absl_check - absl/log:absl_log - absl/strings:charset @@ -9014,7 +9020,6 @@ targets: - test/cpp/util/proto_file_parser.h - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - - third_party/utf8_range/utf8_validity.h src: - src/core/ext/transport/chaotic_good/chaotic_good_frame.proto - src/proto/grpc/channelz/v2/latent_see.proto @@ -9070,8 +9075,6 @@ targets: - third_party/googletest/googlemock/src/gmock_main.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - grpcpp_channelz - protoc - grpc_test_util @@ -9365,7 +9368,6 @@ targets: - test/cpp/util/proto_file_parser.h - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - - third_party/utf8_range/utf8_validity.h src: - src/proto/grpc/reflection/v1alpha/reflection.proto - src/proto/grpc/testing/echo.proto @@ -9384,8 +9386,6 @@ targets: - test/cpp/util/service_describer.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - protoc - grpc++_test_util - name: client_auth_filter_test @@ -15031,7 +15031,6 @@ targets: - test/cpp/util/proto_file_parser.h - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - - third_party/utf8_range/utf8_validity.h src: - src/proto/grpc/reflection/v1alpha/reflection.proto - test/cpp/util/cli_call.cc @@ -15042,8 +15041,6 @@ targets: - test/cpp/util/proto_reflection_descriptor_database.cc - test/cpp/util/service_describer.cc deps: - - absl/container:fixed_array - - absl/functional:overload - grpc++ - protoc - grpc++_test_config @@ -15235,7 +15232,6 @@ targets: - test/cpp/util/proto_file_parser.h - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - - third_party/utf8_range/utf8_validity.h src: - src/proto/grpc/testing/echo.proto - src/proto/grpc/testing/echo_messages.proto @@ -15253,8 +15249,6 @@ targets: - test/cpp/util/service_describer.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - grpc++_reflection - protoc - grpc++_test_config @@ -17201,20 +17195,31 @@ targets: run: false language: c++ headers: + - src/cpp/ext/otel/key_value_iterable.h + - src/cpp/ext/otel/otel_client_call_tracer.h + - src/cpp/ext/otel/otel_plugin.h + - src/cpp/ext/otel/otel_server_call_tracer.h - test/core/credentials/call/oauth2/oauth2_utils.h - test/cpp/interop/backend_metrics_lb_policy.h - test/cpp/interop/client_helper.h - test/cpp/interop/interop_client.h + - test/cpp/interop/otel_helper.h src: - src/proto/grpc/testing/empty.proto - src/proto/grpc/testing/messages.proto - src/proto/grpc/testing/test.proto + - src/cpp/ext/otel/otel_client_call_tracer.cc + - src/cpp/ext/otel/otel_plugin.cc + - src/cpp/ext/otel/otel_server_call_tracer.cc - test/core/credentials/call/oauth2/oauth2_utils.cc - test/cpp/interop/backend_metrics_lb_policy.cc - test/cpp/interop/client.cc - test/cpp/interop/client_helper.cc - test/cpp/interop/interop_client.cc + - test/cpp/interop/otel_helper.cc deps: + - opentelemetry-cpp::api + - opentelemetry-cpp::trace - grpc++_test_config - grpc++_test_util - name: interop_server @@ -17222,17 +17227,28 @@ targets: run: false language: c++ headers: + - src/cpp/ext/otel/key_value_iterable.h + - src/cpp/ext/otel/otel_client_call_tracer.h + - src/cpp/ext/otel/otel_plugin.h + - src/cpp/ext/otel/otel_server_call_tracer.h - src/cpp/server/orca/orca_service.h + - test/cpp/interop/otel_helper.h - test/cpp/interop/server_helper.h src: - src/proto/grpc/testing/empty.proto - src/proto/grpc/testing/messages.proto - src/proto/grpc/testing/test.proto + - src/cpp/ext/otel/otel_client_call_tracer.cc + - src/cpp/ext/otel/otel_plugin.cc + - src/cpp/ext/otel/otel_server_call_tracer.cc - src/cpp/server/orca/orca_service.cc - test/cpp/interop/interop_server.cc - test/cpp/interop/interop_server_bootstrap.cc + - test/cpp/interop/otel_helper.cc - test/cpp/interop/server_helper.cc deps: + - opentelemetry-cpp::api + - opentelemetry-cpp::trace - grpc++_test_config - grpc++_test_util - name: invalid_call_argument_test @@ -18055,7 +18071,6 @@ targets: - test/cpp/util/proto_file_parser.h - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - - third_party/utf8_range/utf8_validity.h src: - src/core/ext/transport/chaotic_good/chaotic_good_frame.proto - src/proto/grpc/channelz/v2/channelz.proto @@ -18108,8 +18123,6 @@ targets: - third_party/googletest/googlemock/src/gmock_main.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - grpc++ - protoc - grpc_test_util @@ -24265,7 +24278,8 @@ targets: gtest: true build: test language: c++ - headers: [] + headers: + - third_party/utf8_range/utf8_validity.h src: - third_party/googleapis/google/rpc/status.proto - src/cpp/client/global_callback_hook.cc @@ -24275,7 +24289,9 @@ targets: - gtest - absl/base:dynamic_annotations - absl/base:prefetch + - absl/container:fixed_array - absl/container:layout + - absl/functional:overload - absl/log:absl_check - absl/log:absl_log - absl/strings:charset @@ -28339,7 +28355,6 @@ targets: - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - third_party/protoc-gen-validate/validate/validate.h - - third_party/utf8_range/utf8_validity.h src: - src/proto/grpc/reflection/v1alpha/reflection.proto - third_party/cel-spec/proto/cel/expr/checked.proto @@ -28439,8 +28454,6 @@ targets: - test/cpp/util/service_describer.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - grpc++ - protoc - grpc_test_util @@ -28487,6 +28500,7 @@ targets: - test/core/xds/xds_client_test_peer.h - test/core/xds/xds_transport_fake.h - third_party/protoc-gen-validate/validate/validate.h + - third_party/utf8_range/utf8_validity.h src: - test/core/event_engine/fuzzing_event_engine/fuzzing_event_engine.proto - third_party/cel-spec/proto/cel/expr/checked.proto @@ -28634,7 +28648,9 @@ targets: - gtest - absl/base:dynamic_annotations - absl/base:prefetch + - absl/container:fixed_array - absl/container:layout + - absl/functional:overload - absl/log:absl_check - absl/log:absl_log - absl/strings:charset @@ -28803,6 +28819,7 @@ targets: headers: - test/core/test_util/scoped_env_var.h - third_party/protoc-gen-validate/validate/validate.h + - third_party/utf8_range/utf8_validity.h src: - third_party/cel-spec/proto/cel/expr/checked.proto - third_party/cel-spec/proto/cel/expr/syntax.proto @@ -28932,7 +28949,9 @@ targets: - gtest - absl/base:dynamic_annotations - absl/base:prefetch + - absl/container:fixed_array - absl/container:layout + - absl/functional:overload - absl/log:absl_check - absl/log:absl_log - absl/strings:charset @@ -29261,7 +29280,6 @@ targets: - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - third_party/protoc-gen-validate/validate/validate.h - - third_party/utf8_range/utf8_validity.h src: - src/proto/grpc/reflection/v1alpha/reflection.proto - third_party/cel-spec/proto/cel/expr/checked.proto @@ -29432,8 +29450,6 @@ targets: - test/cpp/util/service_describer.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - grpc++ - protoc - grpc_test_util @@ -30109,6 +30125,7 @@ targets: headers: - test/core/test_util/scoped_env_var.h - third_party/protoc-gen-validate/validate/validate.h + - third_party/utf8_range/utf8_validity.h src: - third_party/envoy-api/envoy/annotations/deprecation.proto - third_party/envoy-api/envoy/annotations/resource.proto @@ -30185,7 +30202,9 @@ targets: - gtest - absl/base:dynamic_annotations - absl/base:prefetch + - absl/container:fixed_array - absl/container:layout + - absl/functional:overload - absl/log:absl_check - absl/log:absl_log - absl/strings:charset @@ -30815,7 +30834,6 @@ targets: - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - third_party/protoc-gen-validate/validate/validate.h - - third_party/utf8_range/utf8_validity.h src: - src/proto/grpc/reflection/v1alpha/reflection.proto - test/core/event_engine/fuzzing_event_engine/fuzzing_event_engine.proto @@ -30943,8 +30961,6 @@ targets: - test/cpp/util/service_describer.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - grpc++ - protoc - grpc_test_util @@ -30961,7 +30977,6 @@ targets: - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - third_party/protoc-gen-validate/validate/validate.h - - third_party/utf8_range/utf8_validity.h src: - src/proto/grpc/reflection/v1alpha/reflection.proto - third_party/cel-spec/proto/cel/expr/checked.proto @@ -31071,8 +31086,6 @@ targets: - test/cpp/util/service_describer.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - grpc++ - protoc - grpc_test_util @@ -31088,7 +31101,6 @@ targets: - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - third_party/protoc-gen-validate/validate/validate.h - - third_party/utf8_range/utf8_validity.h src: - src/proto/grpc/reflection/v1alpha/reflection.proto - third_party/cel-spec/proto/cel/expr/checked.proto @@ -31215,8 +31227,6 @@ targets: - test/cpp/util/service_describer.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - grpc++ - protoc - grpc_test_util @@ -32330,7 +32340,6 @@ targets: - test/cpp/util/proto_reflection_descriptor_database.h - test/cpp/util/service_describer.h - third_party/protoc-gen-validate/validate/validate.h - - third_party/utf8_range/utf8_validity.h src: - src/proto/grpc/lookup/v1/rls_config.proto - src/proto/grpc/reflection/v1alpha/reflection.proto @@ -32432,8 +32441,6 @@ targets: - test/cpp/util/service_describer.cc deps: - gtest - - absl/container:fixed_array - - absl/functional:overload - grpc++ - protoc - grpc_test_util diff --git a/gRPC-C++.podspec b/gRPC-C++.podspec index 7a43530841b04..5fa37e957ebc2 100644 --- a/gRPC-C++.podspec +++ b/gRPC-C++.podspec @@ -245,6 +245,7 @@ Pod::Spec.new do |s| ss.dependency 'abseil/base/prefetch', abseil_version ss.dependency 'abseil/cleanup/cleanup', abseil_version ss.dependency 'abseil/container/btree', abseil_version + ss.dependency 'abseil/container/fixed_array', abseil_version ss.dependency 'abseil/container/flat_hash_map', abseil_version ss.dependency 'abseil/container/flat_hash_set', abseil_version ss.dependency 'abseil/container/inlined_vector', abseil_version @@ -255,6 +256,7 @@ Pod::Spec.new do |s| ss.dependency 'abseil/functional/any_invocable', abseil_version ss.dependency 'abseil/functional/bind_front', abseil_version ss.dependency 'abseil/functional/function_ref', abseil_version + ss.dependency 'abseil/functional/overload', abseil_version ss.dependency 'abseil/hash/hash', abseil_version ss.dependency 'abseil/log/absl_check', abseil_version ss.dependency 'abseil/log/absl_log', abseil_version @@ -1715,6 +1717,7 @@ Pod::Spec.new do |s| 'third_party/utf8_range/utf8_range.h', 'third_party/utf8_range/utf8_range_neon.inc', 'third_party/utf8_range/utf8_range_sse.inc', + 'third_party/utf8_range/utf8_validity.h', 'third_party/xxhash/xxhash.h', 'third_party/zlib/crc32.h', 'third_party/zlib/deflate.h', @@ -3110,6 +3113,7 @@ Pod::Spec.new do |s| 'third_party/utf8_range/utf8_range.h', 'third_party/utf8_range/utf8_range_neon.inc', 'third_party/utf8_range/utf8_range_sse.inc', + 'third_party/utf8_range/utf8_validity.h', 'third_party/xxhash/xxhash.h', 'third_party/zlib/crc32.h', 'third_party/zlib/deflate.h', diff --git a/src/android/test/interop/app/CMakeLists.txt b/src/android/test/interop/app/CMakeLists.txt index 285d9019d0662..8bfdebd5fd8da 100644 --- a/src/android/test/interop/app/CMakeLists.txt +++ b/src/android/test/interop/app/CMakeLists.txt @@ -114,6 +114,8 @@ add_library(grpc-interop ${GRPC_SRC_DIR}/test/cpp/interop/backend_metrics_lb_policy.cc ${GRPC_SRC_DIR}/test/cpp/interop/interop_client.h ${GRPC_SRC_DIR}/test/cpp/interop/interop_client.cc + ${GRPC_SRC_DIR}/test/cpp/interop/otel_helper.h + ${GRPC_SRC_DIR}/test/cpp/interop/otel_helper.cc ${GRPC_SRC_DIR}/test/core/test_util/histogram.h ${GRPC_SRC_DIR}/test/core/test_util/histogram.cc ${GRPC_SRC_DIR}/test/core/test_util/test_lb_policies.h diff --git a/src/core/call/metadata_batch.h b/src/core/call/metadata_batch.h index eb988ccf0bd38..db828f798c18b 100644 --- a/src/core/call/metadata_batch.h +++ b/src/core/call/metadata_batch.h @@ -317,7 +317,7 @@ struct GrpcServerStatsBinMetadata : public SimpleSliceBasedMetadata { // grpc-trace-bin metadata trait. struct GrpcTraceBinMetadata : public SimpleSliceBasedMetadata { - static constexpr bool kPublishToApp = false; + static constexpr bool kPublishToApp = true; static constexpr bool kRepeatable = false; static constexpr bool kTransferOnTrailersOnly = false; using CompressionTraits = FrequentKeyWithNoValueCompressionCompressor; @@ -326,7 +326,7 @@ struct GrpcTraceBinMetadata : public SimpleSliceBasedMetadata { // grpc-tags-bin metadata trait. struct GrpcTagsBinMetadata : public SimpleSliceBasedMetadata { - static constexpr bool kPublishToApp = false; + static constexpr bool kPublishToApp = true; static constexpr bool kRepeatable = false; static constexpr bool kTransferOnTrailersOnly = false; using CompressionTraits = FrequentKeyWithNoValueCompressionCompressor; diff --git a/src/core/lib/surface/call_utils.h b/src/core/lib/surface/call_utils.h index 234ab18ce1693..020b8d0463cb5 100644 --- a/src/core/lib/surface/call_utils.h +++ b/src/core/lib/surface/call_utils.h @@ -112,6 +112,12 @@ class PublishToAppEncoder { if constexpr (std::is_same::value) { Append(Which::key(), value); } + if constexpr (std::is_same::value) { + Append(Which::key(), value); + } + if constexpr (std::is_same::value) { + Append(Which::key(), value); + } } } diff --git a/src/python/grpcio_tests/tests/bazel_namespace_package_hack.py b/src/python/grpcio_tests/tests/bazel_namespace_package_hack.py index 994a8e1e8004b..068bc5757fe34 100644 --- a/src/python/grpcio_tests/tests/bazel_namespace_package_hack.py +++ b/src/python/grpcio_tests/tests/bazel_namespace_package_hack.py @@ -38,3 +38,19 @@ def sys_path_to_site_dir_hack(): items.append(item) for item in items: site.addsitedir(item) + try: + import pkgutil + + import google + + google.__path__ = pkgutil.extend_path(google.__path__, google.__name__) + except Exception: + pass + try: + import pkgutil + + import src + + src.__path__ = pkgutil.extend_path(src.__path__, src.__name__) + except Exception: + pass diff --git a/src/python/grpcio_tests/tests/interop/BUILD.bazel b/src/python/grpcio_tests/tests/interop/BUILD.bazel index fd6327691f477..c172d29f90331 100644 --- a/src/python/grpcio_tests/tests/interop/BUILD.bazel +++ b/src/python/grpcio_tests/tests/interop/BUILD.bazel @@ -25,17 +25,41 @@ py_library( ], ) +py_library( + name = "otel_interop_helper", + srcs = ["otel_interop_helper.py"], + deps = [ + "//src/proto/grpc/testing:py_test_proto", + "//src/python/grpcio/grpc:grpcio", + "@opentelemetry_proto//:common_proto_py", + "@opentelemetry_proto//:trace_proto_py", + "@opentelemetry_proto//:trace_service_grpc_py", + "@opentelemetry_proto//:trace_service_proto_py", + requirement("opentelemetry-api"), + requirement("opentelemetry-sdk"), + ], +) + py_library( name = "client_lib", srcs = ["client.py"], - imports = ["../../"], + imports = [ + "../../", + "../../../../", + ], deps = [ ":methods", + ":otel_interop_helper", ":resources", "//src/proto/grpc/testing:py_test_proto", + "//src/proto/grpc/testing:test_py_pb2_grpc", "//src/python/grpcio/grpc:grpcio", + "//src/python/grpcio_observability/grpc_observability:pyobservability", + "//src/python/grpcio_tests/tests:bazel_namespace_package_hack", requirement("absl-py"), requirement("google-auth"), + requirement("opentelemetry-api"), + requirement("opentelemetry-sdk"), ], ) @@ -49,7 +73,10 @@ py_binary( py_library( name = "methods", srcs = ["methods.py"], - imports = ["../../"], + imports = [ + "../../", + "../../../../", + ], deps = [ "//src/proto/grpc/testing:empty_py_pb2", "//src/proto/grpc/testing:py_messages_proto", @@ -80,7 +107,10 @@ py_library( py_library( name = "service", srcs = ["service.py"], - imports = ["../../"], + imports = [ + "../../", + "../../../../", + ], deps = [ "//src/proto/grpc/testing:empty_py_pb2", "//src/proto/grpc/testing:py_messages_proto", @@ -92,14 +122,23 @@ py_library( py_library( name = "server", srcs = ["server.py"], - imports = ["../../"], + imports = [ + "../../", + "../../../../", + ], deps = [ + ":otel_interop_helper", ":resources", ":service", "//src/proto/grpc/testing:py_test_proto", + "//src/proto/grpc/testing:test_py_pb2_grpc", "//src/python/grpcio/grpc:grpcio", + "//src/python/grpcio_observability/grpc_observability:pyobservability", + "//src/python/grpcio_tests/tests:bazel_namespace_package_hack", "//src/python/grpcio_tests/tests/unit:test_common", requirement("absl-py"), + requirement("opentelemetry-api"), + requirement("opentelemetry-sdk"), ], ) diff --git a/src/python/grpcio_tests/tests/interop/_insecure_intraop_test.py b/src/python/grpcio_tests/tests/interop/_insecure_intraop_test.py index 16f7bdfe110aa..6f34ac0d43e78 100644 --- a/src/python/grpcio_tests/tests/interop/_insecure_intraop_test.py +++ b/src/python/grpcio_tests/tests/interop/_insecure_intraop_test.py @@ -13,8 +13,17 @@ # limitations under the License. """Insecure client-server interoperability as a unit test.""" +import os import unittest +os.environ["GRPC_BAZEL_RUNTIME"] = "1" +try: + from tests import bazel_namespace_package_hack + + bazel_namespace_package_hack.sys_path_to_site_dir_hack() +except ImportError: + pass + import grpc from src.proto.grpc.testing import test_pb2_grpc diff --git a/src/python/grpcio_tests/tests/interop/_secure_intraop_test.py b/src/python/grpcio_tests/tests/interop/_secure_intraop_test.py index b5568290c6c6c..89818f709a296 100644 --- a/src/python/grpcio_tests/tests/interop/_secure_intraop_test.py +++ b/src/python/grpcio_tests/tests/interop/_secure_intraop_test.py @@ -13,8 +13,17 @@ # limitations under the License. """Secure client-server interoperability as a unit test.""" +import os import unittest +os.environ["GRPC_BAZEL_RUNTIME"] = "1" +try: + from tests import bazel_namespace_package_hack + + bazel_namespace_package_hack.sys_path_to_site_dir_hack() +except ImportError: + pass + import grpc import grpc.experimental @@ -65,7 +74,6 @@ def tearDown(self): class SecureInteropWithSyncPrivateKeyOffloadingTest( _intraop_test_case.IntraopTestCase, unittest.TestCase ): - def setUp(self): self.server = test_common.test_server() test_pb2_grpc.add_TestServiceServicer_to_server( @@ -110,7 +118,6 @@ def tearDown(self): class SecureInteropWithAsyncPrivateKeyOffloadingTest( _intraop_test_case.IntraopTestCase, unittest.TestCase ): - def setUp(self): self.server = test_common.test_server() test_pb2_grpc.add_TestServiceServicer_to_server( diff --git a/src/python/grpcio_tests/tests/interop/client.py b/src/python/grpcio_tests/tests/interop/client.py index 6075ac5a5882c..f44cc830b2058 100644 --- a/src/python/grpcio_tests/tests/interop/client.py +++ b/src/python/grpcio_tests/tests/interop/client.py @@ -15,16 +15,25 @@ import os +os.environ["GRPC_BAZEL_RUNTIME"] = "1" +try: + from tests import bazel_namespace_package_hack + + bazel_namespace_package_hack.sys_path_to_site_dir_hack() +except ImportError: + pass + +# pylint: disable=wrong-import-position from absl import app from absl.flags import argparse_flags -from google import auth as google_auth -from google.auth import jwt as google_auth_jwt import grpc from src.proto.grpc.testing import test_pb2_grpc from tests.interop import methods from tests.interop import resources +# pylint: enable=wrong-import-position + def parse_interop_client_args(argv): parser = argparse_flags.ArgumentParser() @@ -93,18 +102,27 @@ def parse_interop_client_args(argv): " round_robin " + "or pick_first)." ), ) + parser.add_argument( + "--enable_opentelemetry", + default=False, + type=resources.parse_bool, + help="enable OpenTelemetry tracing/observability", + ) return parser.parse_args(argv[1:]) def _create_call_credentials(args): + from google import auth as google_auth + from google.auth import jwt as google_auth_jwt + if args.test_case == "oauth2_auth_token": - google_credentials, unused_project_id = google_auth.default( + google_credentials, _unused_project_id = google_auth.default( scopes=[args.oauth_scope] ) google_credentials.refresh(google_auth.transport.requests.Request()) return grpc.access_token_call_credentials(google_credentials.token) elif args.test_case == "compute_engine_creds": - google_credentials, unused_project_id = google_auth.default( + google_credentials, _unused_project_id = google_auth.default( scopes=[args.oauth_scope] ) return grpc.metadata_call_credentials( @@ -129,6 +147,8 @@ def _create_call_credentials(args): def get_secure_channel_parameters(args): + from google import auth as google_auth + call_credentials = _create_call_credentials(args) channel_opts = () @@ -143,7 +163,7 @@ def get_secure_channel_parameters(args): if args.custom_credentials_type is not None: if args.custom_credentials_type == "compute_engine_channel_creds": assert call_credentials is None - google_credentials, unused_project_id = google_auth.default( + google_credentials, _unused_project_id = google_auth.default( scopes=[args.oauth_scope] ) call_creds = grpc.metadata_call_credentials( @@ -186,6 +206,59 @@ def get_secure_channel_parameters(args): return channel_credentials, channel_opts +class _OTelClientInterceptor(grpc.UnaryUnaryClientInterceptor): + def __init__(self, tracer): + self._tracer = tracer + + def intercept_unary_unary(self, continuation, client_call_details, request): + from grpc._interceptor import _ClientCallDetails + from opentelemetry import trace + + from tests.interop import otel_interop_helper + + method = client_call_details.method + full_method = method.lstrip("/") + sent_span_name = f"Sent.{full_method}" + attempt_span_name = f"Attempt.{full_method}" + + sent_span = self._tracer.start_span( + sent_span_name, kind=trace.SpanKind.CLIENT + ) + sent_ctx = trace.set_span_in_context(sent_span) + attempt_span = self._tracer.start_span( + attempt_span_name, kind=trace.SpanKind.CLIENT, context=sent_ctx + ) + attempt_span.set_attribute("previous-rpc-attempts", 0) + attempt_span.set_attribute("transparent-retry", False) + attempt_span.add_event("Outbound message") + + trace_bin_bytes = otel_interop_helper.pack_grpc_trace_bin( + attempt_span.get_span_context().trace_id, + attempt_span.get_span_context().span_id, + ) + + metadata = list(client_call_details.metadata or []) + metadata.append(("grpc-trace-bin", trace_bin_bytes)) + + new_details = _ClientCallDetails( + method=client_call_details.method, + timeout=client_call_details.timeout, + metadata=metadata, + credentials=client_call_details.credentials, + wait_for_ready=client_call_details.wait_for_ready, + compression=client_call_details.compression, + ) + + try: + response = continuation(new_details, request) + attempt_span.add_event("Inbound message") + return response + finally: + attempt_span.end() + sent_span.end() + otel_interop_helper.flush_tracer_provider() + + def _create_channel(args): target = "{}:{}".format(args.server_host, args.server_port) @@ -195,9 +268,19 @@ def _create_channel(args): or args.custom_credentials_type is not None ): channel_credentials, options = get_secure_channel_parameters(args) - return grpc.secure_channel(target, channel_credentials, options) + channel = grpc.secure_channel(target, channel_credentials, options) else: - return grpc.insecure_channel(target) + channel = grpc.insecure_channel(target) + + if args.enable_opentelemetry: + from tests.interop import otel_interop_helper + + _, tracer = otel_interop_helper.init_tracer_provider() + channel = grpc.intercept_channel( + channel, _OTelClientInterceptor(tracer) + ) + + return channel def create_stub(channel, args): @@ -220,6 +303,10 @@ def test_interoperability(args): stub = create_stub(channel, args) test_case = _test_case_from_arg(args.test_case) test_case.test_interoperability(stub, args) + if args.enable_opentelemetry: + from tests.interop import otel_interop_helper + + otel_interop_helper.flush_tracer_provider() if __name__ == "__main__": diff --git a/src/python/grpcio_tests/tests/interop/otel_interop_helper.py b/src/python/grpcio_tests/tests/interop/otel_interop_helper.py new file mode 100644 index 0000000000000..24b378c1ba310 --- /dev/null +++ b/src/python/grpcio_tests/tests/interop/otel_interop_helper.py @@ -0,0 +1,272 @@ +# Copyright 2026 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. +"""OpenTelemetry Tracing Interop Helper for Python gRPC Interop Client/Server""" + +import os +from typing import Optional, Tuple + +import grpc +from opentelemetry import trace +from opentelemetry.proto.collector.trace.v1 import trace_service_pb2 +from opentelemetry.proto.collector.trace.v1 import trace_service_pb2_grpc +from opentelemetry.proto.common.v1 import common_pb2 +from opentelemetry.proto.trace.v1 import trace_pb2 +from opentelemetry.sdk.trace import ReadableSpan +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export import SpanExporter + + +class OTLPSpanExporter(SpanExporter): + """Exporter that sends OTLP spans to OTLP Collector over gRPC.""" + + def __init__(self, endpoint: str): + if endpoint.startswith("http://"): + endpoint = endpoint[7:] + elif endpoint.startswith("https://"): + endpoint = endpoint[8:] + self._channel = grpc.insecure_channel(endpoint) + self._stub = trace_service_pb2_grpc.TraceServiceStub(self._channel) + + def export(self, spans: Tuple[ReadableSpan, ...]) -> None: + if not spans: + return + + otlp_spans = [] + for span in spans: + ctx = span.context + parent_ctx = span.parent + + trace_id_bytes = ctx.trace_id.to_bytes(16, "big") + span_id_bytes = ctx.span_id.to_bytes(8, "big") + parent_span_id_bytes = ( + parent_ctx.span_id.to_bytes(8, "big") + if parent_ctx and parent_ctx.span_id + else b"" + ) + + proto_attributes = [] + if span.attributes: + for k, v in span.attributes.items(): + kv = common_pb2.KeyValue(key=k) + if isinstance(v, bool): + kv.value.bool_value = v + elif isinstance(v, int): + kv.value.int_value = v + elif isinstance(v, float): + kv.value.double_value = v + else: + kv.value.string_value = str(v) + proto_attributes.append(kv) + + proto_events = [] + if span.events: + for event in span.events: + e = trace_pb2.Span.Event( + name=event.name, + time_unix_nano=event.timestamp, + ) + proto_events.append(e) + + kind = ( + trace_pb2.Span.SpanKind.SPAN_KIND_CLIENT + if span.kind == trace.SpanKind.CLIENT + else trace_pb2.Span.SpanKind.SPAN_KIND_SERVER + ) + + proto_span = trace_pb2.Span( + trace_id=trace_id_bytes, + span_id=span_id_bytes, + parent_span_id=parent_span_id_bytes, + name=span.name, + kind=kind, + start_time_unix_nano=span.start_time, + end_time_unix_nano=span.end_time, + attributes=proto_attributes, + events=proto_events, + ) + otlp_spans.append(proto_span) + + scope_spans = trace_pb2.ScopeSpans(spans=otlp_spans) + resource_spans = trace_pb2.ResourceSpans(scope_spans=[scope_spans]) + request = trace_service_pb2.ExportTraceServiceRequest( + resource_spans=[resource_spans] + ) + + try: + self._stub.Export(request, timeout=5) + except Exception as e: + print(f"OTLPSpanExporter Export exception: {e}", flush=True) + + def shutdown(self) -> None: + self._channel.close() + + def force_flush(self, timeout_millis: int = 30000) -> bool: + return True + + +_GLOBAL_PROVIDER: Optional[TracerProvider] = None + + +def init_tracer_provider() -> Tuple[TracerProvider, trace.Tracer]: + global _GLOBAL_PROVIDER + if _GLOBAL_PROVIDER is None: + endpoint = os.environ.get( + "OTEL_EXPORTER_OTLP_ENDPOINT", "http://localhost:4317" + ) + exporter = OTLPSpanExporter(endpoint) + processor = SimpleSpanProcessor(exporter) + _GLOBAL_PROVIDER = TracerProvider() + _GLOBAL_PROVIDER.add_span_processor(processor) + trace.set_tracer_provider(_GLOBAL_PROVIDER) + tracer = trace.get_tracer("grpc-python-interop") + return _GLOBAL_PROVIDER, tracer + + +def flush_tracer_provider(): + if _GLOBAL_PROVIDER: + _GLOBAL_PROVIDER.force_flush() + + +def pack_grpc_trace_bin( + trace_id_int: int, span_id_int: int, is_sampled: bool = True +) -> bytes: + trace_id_bytes = trace_id_int.to_bytes(16, "big") + span_id_bytes = span_id_int.to_bytes(8, "big") + options = 1 if is_sampled else 0 + return ( + b"\x00\x00" + + trace_id_bytes + + b"\x01" + + span_id_bytes + + b"\x02" + + bytes([options]) + ) + + +def unpack_grpc_trace_bin( + header_bytes: bytes, +) -> Tuple[Optional[int], Optional[int], bool]: + if len(header_bytes) >= 29 and header_bytes[0] == 0: + trace_id_int = int.from_bytes(header_bytes[2:18], "big") + span_id_int = int.from_bytes(header_bytes[19:27], "big") + is_sampled = bool(header_bytes[28] & 1) + return trace_id_int, span_id_int, is_sampled + return None, None, False + + +def parse_traceparent( + header_str: str, +) -> Tuple[Optional[int], Optional[int], bool]: + parts = header_str.split("-") + if len(parts) >= 4 and parts[0] == "00": + try: + trace_id_int = int(parts[1], 16) + span_id_int = int(parts[2], 16) + is_sampled = (int(parts[3], 16) & 1) != 0 + return trace_id_int, span_id_int, is_sampled + except ValueError: + pass + return None, None, False + + +class OTelServerInterceptor(grpc.ServerInterceptor): + """Server interceptor to extract trace context and create Recv span.""" + + def __init__(self, tracer: trace.Tracer): + self._tracer = tracer + + def intercept_service(self, continuation, handler_call_details): + trace_bin_header = None + traceparent_header = None + for k, v in handler_call_details.invocation_metadata: + k_str = ( + k.decode("ascii", errors="ignore") + if isinstance(k, bytes) + else str(k) + ) + if k_str.lower() == "grpc-trace-bin": + trace_bin_header = v + elif k_str.lower() == "traceparent": + traceparent_header = ( + v if isinstance(v, str) else v.decode("latin1") + ) + + parent_ctx = None + trace_id, parent_span_id, is_sampled = None, None, False + + if trace_bin_header: + if isinstance(trace_bin_header, str): + trace_bin_header = trace_bin_header.encode("latin1") + ( + trace_id, + parent_span_id, + is_sampled, + ) = unpack_grpc_trace_bin(trace_bin_header) + elif traceparent_header: + ( + trace_id, + parent_span_id, + is_sampled, + ) = parse_traceparent(traceparent_header) + + if trace_id and parent_span_id: + parent_ctx = trace.SpanContext( + trace_id=trace_id, + span_id=parent_span_id, + is_remote=True, + trace_flags=trace.TraceFlags(1 if is_sampled else 0), + ) + + method = handler_call_details.method + full_method = method.lstrip("/") + span_name = f"Recv.{full_method}" + + if parent_ctx: + ctx = trace.set_span_in_context(trace.NonRecordingSpan(parent_ctx)) + server_span = self._tracer.start_span( + span_name, kind=trace.SpanKind.SERVER, context=ctx + ) + else: + server_span = self._tracer.start_span( + span_name, kind=trace.SpanKind.SERVER + ) + + server_span.add_event("Inbound message") + + handler = continuation(handler_call_details) + + if handler is None: + server_span.end() + return None + + if handler.unary_unary: + orig_func = handler.unary_unary + + def wrapper(request, context): + try: + res = orig_func(request, context) + server_span.add_event("Outbound message") + return res + finally: + server_span.end() + flush_tracer_provider() + + return grpc.unary_unary_rpc_method_handler( + wrapper, + request_deserializer=handler.request_deserializer, + response_serializer=handler.response_serializer, + ) + + return handler diff --git a/src/python/grpcio_tests/tests/interop/security_test.py b/src/python/grpcio_tests/tests/interop/security_test.py index 76001bb99ffda..f56ccbe7a94e2 100644 --- a/src/python/grpcio_tests/tests/interop/security_test.py +++ b/src/python/grpcio_tests/tests/interop/security_test.py @@ -15,6 +15,7 @@ import faulthandler from functools import partial import gc +import os import queue import sys import threading @@ -22,6 +23,14 @@ import unittest import weakref +os.environ["GRPC_BAZEL_RUNTIME"] = "1" +try: + from tests import bazel_namespace_package_hack + + bazel_namespace_package_hack.sys_path_to_site_dir_hack() +except ImportError: + pass + import grpc import grpc.experimental diff --git a/src/python/grpcio_tests/tests/interop/server.py b/src/python/grpcio_tests/tests/interop/server.py index cc705b026da9b..e9189422431b8 100644 --- a/src/python/grpcio_tests/tests/interop/server.py +++ b/src/python/grpcio_tests/tests/interop/server.py @@ -13,8 +13,20 @@ # limitations under the License. """The Python implementation of the GRPC interoperability test server.""" +import os + +os.environ["GRPC_BAZEL_RUNTIME"] = "1" +try: + from tests import bazel_namespace_package_hack + + bazel_namespace_package_hack.sys_path_to_site_dir_hack() +except ImportError: + pass + +# pylint: disable=wrong-import-position from concurrent import futures import logging +import signal from absl import app from absl.flags import argparse_flags @@ -25,6 +37,8 @@ from tests.interop import service from tests.unit import test_common +# pylint: enable=wrong-import-position + logging.basicConfig() _LOGGER = logging.getLogger(__name__) @@ -46,6 +60,12 @@ def parse_interop_server_arguments(argv): type=resources.parse_bool, help="require an ALTS connection", ) + parser.add_argument( + "--enable_opentelemetry", + default=False, + type=resources.parse_bool, + help="enable OpenTelemetry tracing/observability", + ) return parser.parse_args(argv[1:]) @@ -58,8 +78,42 @@ def get_server_credentials(use_tls): return grpc.alts_server_credentials() +def _serve_internal(server, enable_otel=False): + def _sig_handler(signum, frame): + _LOGGER.info("Received signal %d, stopping server...", signum) + if enable_otel: + from tests.interop import otel_interop_helper + + otel_interop_helper.flush_tracer_provider() + server.stop(0) + + signal.signal(signal.SIGTERM, _sig_handler) + signal.signal(signal.SIGINT, _sig_handler) + + server.start() + _LOGGER.info("Server serving.") + server.wait_for_termination() + if enable_otel: + from tests.interop import otel_interop_helper + + otel_interop_helper.flush_tracer_provider() + _LOGGER.info("Server stopped; exiting.") + + def serve(args): - server = test_common.test_server() + if args.enable_opentelemetry: + from tests.interop import otel_interop_helper + + _, tracer = otel_interop_helper.init_tracer_provider() + interceptor = otel_interop_helper.OTelServerInterceptor(tracer) + server = grpc.server( + futures.ThreadPoolExecutor(max_workers=10), + options=(("grpc.so_reuseport", 0),), + interceptors=(interceptor,), + ) + else: + server = test_common.test_server() + test_pb2_grpc.add_TestServiceServicer_to_server( service.TestService(), server ) @@ -69,10 +123,7 @@ def serve(args): else: server.add_insecure_port("[::]:{}".format(args.port)) - server.start() - _LOGGER.info("Server serving.") - server.wait_for_termination() - _LOGGER.info("Server stopped; exiting.") + _serve_internal(server, enable_otel=args.enable_opentelemetry) if __name__ == "__main__": diff --git a/templates/MODULE.bazel.inja b/templates/MODULE.bazel.inja index ccb498396367d..60f7d8fa5de37 100644 --- a/templates/MODULE.bazel.inja +++ b/templates/MODULE.bazel.inja @@ -51,6 +51,7 @@ bazel_dep(name = "googletest", version = "1.17.0", repo_name = "com_google_googl bazel_dep(name = "opencensus-cpp", version = "0.0.0-20230502-50eb5de.bcr.2", repo_name = "io_opencensus_cpp") bazel_dep(name = "openssl", version = "3.3.1.bcr.1") bazel_dep(name = "opentelemetry-cpp", version = "1.19.0", repo_name = "io_opentelemetry_cpp") +bazel_dep(name = "opentelemetry-proto", version = "1.8.0", repo_name = "opentelemetry_proto") # --- Protobuf related packages. bazel_dep(name = "protobuf", version = "35.1", repo_name = "com_google_protobuf") diff --git a/test/cpp/interop/BUILD b/test/cpp/interop/BUILD index 166dd00c6997e..da9e136e75b92 100644 --- a/test/cpp/interop/BUILD +++ b/test/cpp/interop/BUILD @@ -94,6 +94,28 @@ grpc_cc_binary( ], ) +grpc_cc_library( + name = "otel_helper_lib", + srcs = [ + "otel_helper.cc", + ], + hdrs = [ + "otel_helper.h", + ], + external_deps = [ + "absl/flags:flag", + "absl/log", + "otel/api", + "otel/exporters/otlp:otlp_grpc_exporter", + "otel/sdk/src/trace", + "otel/sdk:headers", + ], + deps = [ + "//:grpc++", + "//src/cpp/ext/otel:otel_plugin", + ], +) + grpc_cc_library( name = "interop_server_lib", srcs = [ @@ -105,6 +127,7 @@ grpc_cc_library( "absl/log:check", ], deps = [ + ":otel_helper_lib", ":server_helper_lib", "//:gpr", "//:grpc", @@ -171,6 +194,7 @@ grpc_cc_library( ], deps = [ ":client_helper_lib", + ":otel_helper_lib", "//:gpr", "//:grpc", "//:grpc++", @@ -786,3 +810,25 @@ grpc_cc_library( "//src/core:pollset_set", ], ) + +grpc_cc_binary( + name = "otlp_collector", + srcs = [ + "otlp_collector.cc", + ], + external_deps = [ + "absl/flags:flag", + "absl/flags:parse", + "absl/log", + "protobuf", + "protobuf_clib", + "protobuf_headers", + ], + tags = ["nobuilder"], + deps = [ + "//:gpr", + "//:grpc++", + "@opentelemetry_proto//:trace_service_grpc_cc", + "@opentelemetry_proto//:trace_service_proto_cc", + ], +) diff --git a/test/cpp/interop/client.cc b/test/cpp/interop/client.cc index 9cdf206447ef0..199e7cfb283e3 100644 --- a/test/cpp/interop/client.cc +++ b/test/cpp/interop/client.cc @@ -30,6 +30,7 @@ #include "test/core/test_util/test_config.h" #include "test/cpp/interop/client_helper.h" #include "test/cpp/interop/interop_client.h" +#include "test/cpp/interop/otel_helper.h" #include "test/cpp/util/test_config.h" #include "absl/flags/flag.h" #include "absl/log/log.h" @@ -204,6 +205,7 @@ ParseAdditionalMetadataFlag(const std::string& flag) { int main(int argc, char** argv) { grpc::testing::TestEnvironment env(&argc, argv); grpc::testing::InitTest(&argc, &argv, true); + grpc::testing::interop::MaybeRegisterOpenTelemetry(); LOG(INFO) << "Testing these cases: " << absl::GetFlag(FLAGS_test_case); int ret = 0; @@ -221,12 +223,14 @@ int main(int argc, char** argv) { factories; if (!additional_metadata->empty()) { factories.emplace_back( - new grpc::testing::AdditionalMetadataInterceptorFactory( + std::make_unique< + grpc::testing::AdditionalMetadataInterceptorFactory>( *additional_metadata)); } if (absl::GetFlag(FLAGS_log_metadata_and_status)) { factories.emplace_back( - new grpc::testing::MetadataAndStatusLoggerInterceptorFactory()); + std::make_unique< + grpc::testing::MetadataAndStatusLoggerInterceptorFactory>()); } std::string service_config_json = absl::GetFlag(FLAGS_service_config_json); @@ -357,5 +361,6 @@ int main(int argc, char** argv) { ret = 1; } + grpc::testing::interop::ForceFlushOpenTelemetry(); return ret; } diff --git a/test/cpp/interop/interop_server.cc b/test/cpp/interop/interop_server.cc index b2ca0eafae557..626f515d23a99 100644 --- a/test/cpp/interop/interop_server.cc +++ b/test/cpp/interop/interop_server.cc @@ -38,6 +38,7 @@ #include "src/proto/grpc/testing/empty.pb.h" #include "src/proto/grpc/testing/messages.pb.h" #include "src/proto/grpc/testing/test.grpc.pb.h" +#include "test/cpp/interop/otel_helper.h" #include "test/cpp/interop/server_helper.h" #include "test/cpp/util/test_config.h" #include "absl/flags/flag.h" @@ -419,6 +420,7 @@ void grpc::testing::interop::RunServer( ServerStartedCondition* server_started_condition, std::unique_ptr>> server_options) { + grpc::testing::interop::MaybeRegisterOpenTelemetry(); GRPC_CHECK_NE(port, 0); std::ostringstream server_address; server_address << "0.0.0.0:" << port; @@ -457,4 +459,5 @@ void grpc::testing::interop::RunServer( gpr_sleep_until(gpr_time_add(gpr_now(GPR_CLOCK_REALTIME), gpr_time_from_seconds(5, GPR_TIMESPAN))); } + grpc::testing::interop::ForceFlushOpenTelemetry(); } diff --git a/test/cpp/interop/interop_server_bootstrap.cc b/test/cpp/interop/interop_server_bootstrap.cc index 8bead419554e9..7e3dffcda288a 100644 --- a/test/cpp/interop/interop_server_bootstrap.cc +++ b/test/cpp/interop/interop_server_bootstrap.cc @@ -32,6 +32,7 @@ int main(int argc, char** argv) { grpc::testing::TestEnvironment env(&argc, argv); grpc::testing::InitTest(&argc, &argv, true); signal(SIGINT, sigint_handler); + signal(SIGTERM, sigint_handler); grpc::testing::interop::RunServer( grpc::testing::CreateInteropServerCredentials()); diff --git a/test/cpp/interop/observability_interop_server_bootstrap.cc b/test/cpp/interop/observability_interop_server_bootstrap.cc index f5b5af76e5a0a..30d1bf1b587d5 100644 --- a/test/cpp/interop/observability_interop_server_bootstrap.cc +++ b/test/cpp/interop/observability_interop_server_bootstrap.cc @@ -38,6 +38,7 @@ int main(int argc, char** argv) { grpc::testing::TestEnvironment env(&argc, argv); grpc::testing::InitTest(&argc, &argv, true); signal(SIGINT, sigint_handler); + signal(SIGTERM, sigint_handler); if (absl::GetFlag(FLAGS_enable_observability)) { // TODO(someone): remove deprecated usage diff --git a/test/cpp/interop/otel_helper.cc b/test/cpp/interop/otel_helper.cc new file mode 100644 index 0000000000000..75d3be208f146 --- /dev/null +++ b/test/cpp/interop/otel_helper.cc @@ -0,0 +1,112 @@ +// +// +// Copyright 2026 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/interop/otel_helper.h" + +#ifndef HAVE_ABSEIL +#define HAVE_ABSEIL +#endif + +#include +#include +#include +#include + +#include "absl/flags/flag.h" +#include "absl/log/log.h" + +#if __has_include() && \ + __has_include("opentelemetry/exporters/otlp/otlp_grpc_exporter_factory.h") +#define GRPC_HAS_OTEL_TRACING 1 +#endif + +#ifdef GRPC_HAS_OTEL_TRACING +#include + +#include "opentelemetry/exporters/otlp/otlp_grpc_exporter_factory.h" +#include "opentelemetry/sdk/trace/simple_processor_factory.h" +#include "opentelemetry/sdk/trace/tracer_provider.h" +#endif + +ABSL_FLAG(bool, enable_opentelemetry, false, + "Whether to enable OpenTelemetry Tracing"); + +namespace grpc { +namespace testing { +namespace interop { + +#ifdef GRPC_HAS_OTEL_TRACING +static std::shared_ptr + g_tracer_provider; +static std::once_flag g_otel_init_once; +#endif + +void MaybeRegisterOpenTelemetry() { +#ifdef GRPC_HAS_OTEL_TRACING + std::call_once(g_otel_init_once, []() { + if (!absl::GetFlag(FLAGS_enable_opentelemetry)) { + return; + } + const char* otel_traces_exporter = std::getenv("OTEL_TRACES_EXPORTER"); + if (otel_traces_exporter != nullptr && + std::string(otel_traces_exporter) == "none") { + LOG(INFO) << "OTEL_TRACES_EXPORTER is set to none. Tracing is disabled."; + return; + } + + // Create OTLP Grpc Exporter + opentelemetry::exporter::otlp::OtlpGrpcExporterOptions opts; + auto exporter = + opentelemetry::exporter::otlp::OtlpGrpcExporterFactory::Create(opts); + auto processor = + opentelemetry::sdk::trace::SimpleSpanProcessorFactory::Create( + std::move(exporter)); + auto provider = std::make_shared( + std::move(processor)); + std::atomic_store(&g_tracer_provider, provider); + + auto status = + grpc::OpenTelemetryPluginBuilder() + .SetTracerProvider(provider) + .SetTextMapPropagator(grpc::OpenTelemetryPluginBuilder:: + MakeGrpcTraceBinTextMapPropagator()) + .BuildAndRegisterGlobal(); + + if (!status.ok()) { + LOG(ERROR) << "Failed to register gRPC OpenTelemetry Plugin: " + << status.ToString(); + } else { + LOG(INFO) + << "Successfully registered gRPC OpenTelemetry Plugin for tracing."; + } + }); +#endif +} + +void ForceFlushOpenTelemetry() { +#ifdef GRPC_HAS_OTEL_TRACING + auto provider = std::atomic_load(&g_tracer_provider); + if (provider != nullptr) { + provider->ForceFlush(); + } +#endif +} + +} // namespace interop +} // namespace testing +} // namespace grpc diff --git a/test/cpp/interop/otel_helper.h b/test/cpp/interop/otel_helper.h new file mode 100644 index 0000000000000..780e896fd973c --- /dev/null +++ b/test/cpp/interop/otel_helper.h @@ -0,0 +1,33 @@ +// +// +// Copyright 2026 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_INTEROP_OTEL_HELPER_H +#define GRPC_TEST_CPP_INTEROP_OTEL_HELPER_H + +namespace grpc { +namespace testing { +namespace interop { + +void MaybeRegisterOpenTelemetry(); +void ForceFlushOpenTelemetry(); + +} // namespace interop +} // namespace testing +} // namespace grpc + +#endif // GRPC_TEST_CPP_INTEROP_OTEL_HELPER_H diff --git a/test/cpp/interop/otlp_collector.cc b/test/cpp/interop/otlp_collector.cc new file mode 100644 index 0000000000000..f7397a1272fbe --- /dev/null +++ b/test/cpp/interop/otlp_collector.cc @@ -0,0 +1,147 @@ +// +// +// Copyright 2026 gRPC authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +// +// + +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "opentelemetry/proto/collector/trace/v1/trace_service.grpc.pb.h" +#include "absl/flags/flag.h" +#include "absl/flags/parse.h" +#include "absl/log/log.h" + +ABSL_FLAG(int, port, 0, "Port to listen on"); +ABSL_FLAG(std::string, file, "", "File to write JSON spans to"); + +class TraceServiceServiceImpl final + : public opentelemetry::proto::collector::trace::v1::TraceService::Service { + public: + explicit TraceServiceServiceImpl(std::string file_path) + : file_path_(std::move(file_path)) {} + + grpc::Status Export( + grpc::ServerContext* /*context*/, + const opentelemetry::proto::collector::trace::v1:: + ExportTraceServiceRequest* request, + opentelemetry::proto::collector::trace::v1::ExportTraceServiceResponse* + /*response*/) override { + std::string json_string; + google::protobuf::json::PrintOptions options; + options.add_whitespace = true; + options.always_print_fields_with_no_presence = true; + options.preserve_proto_field_names = true; + auto status = google::protobuf::json::MessageToJsonString( + *request, &json_string, options); + if (!status.ok()) { + LOG(ERROR) << "Failed to serialize ExportTraceServiceRequest to JSON: " + << status.ToString(); + return grpc::Status(grpc::StatusCode::INTERNAL, + "Failed to serialize to JSON"); + } + + std::lock_guard lock(mu_); + requests_json_.push_back(json_string); + + // Write to a temporary file first and rename to avoid read-during-write + // data race + std::string tmp_file = file_path_ + ".tmp"; + std::ofstream out(tmp_file, std::ios::trunc); + if (!out) { + LOG(ERROR) << "Failed to open file for writing: " << tmp_file; + return grpc::Status(grpc::StatusCode::INTERNAL, "Failed to open file"); + } + out << "[\n"; + for (size_t i = 0; i < requests_json_.size(); ++i) { + out << requests_json_[i]; + if (i + 1 < requests_json_.size()) { + out << ",\n"; + } + } + out << "\n]\n"; + if (out.fail() || out.bad()) { + LOG(ERROR) << "Failed to write JSON spans to file: " << tmp_file; + return grpc::Status(grpc::StatusCode::INTERNAL, + "Failed to write trace output"); + } + out.close(); + + if (std::rename(tmp_file.c_str(), file_path_.c_str()) != 0) { + LOG(ERROR) << "Failed to rename temporary file " << tmp_file << " to " + << file_path_; + return grpc::Status(grpc::StatusCode::INTERNAL, + "Failed to update trace output file"); + } + + return grpc::Status::OK; + } + + private: + std::string file_path_; + std::vector requests_json_; + std::mutex mu_; +}; + +static std::atomic g_got_sigint{false}; + +static void sig_handler(int /*sig*/) { + g_got_sigint.store(true, std::memory_order_relaxed); +} + +int main(int argc, char** argv) { + absl::ParseCommandLine(argc, argv); + int port = absl::GetFlag(FLAGS_port); + std::string file_path = absl::GetFlag(FLAGS_file); + + if (port == 0) { + LOG(ERROR) << "--port is required"; + return 1; + } + if (file_path.empty()) { + LOG(ERROR) << "--file is required"; + return 1; + } + + TraceServiceServiceImpl service(file_path); + + grpc::ServerBuilder builder; + builder.AddListeningPort("0.0.0.0:" + std::to_string(port), + grpc::InsecureServerCredentials()); + builder.RegisterService(&service); + + signal(SIGINT, sig_handler); + signal(SIGTERM, sig_handler); + + std::unique_ptr server(builder.BuildAndStart()); + LOG(INFO) << "OTLP Collector listening on port " << port << "..."; + while (!g_got_sigint.load(std::memory_order_relaxed)) { + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + } + server->Shutdown(); + server.reset(); + + return 0; +} diff --git a/test/cpp/interop/otlp_collector.py b/test/cpp/interop/otlp_collector.py new file mode 100644 index 0000000000000..596ab24832b94 --- /dev/null +++ b/test/cpp/interop/otlp_collector.py @@ -0,0 +1,79 @@ +# +# +# Copyright 2026 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 +from concurrent import futures +import json +import os +import signal +import threading +import time + +from google.protobuf import json_format +import grpc +from opentelemetry.proto.collector.trace.v1 import trace_service_pb2 +from opentelemetry.proto.collector.trace.v1 import trace_service_pb2_grpc + + +class TraceServiceServicer(trace_service_pb2_grpc.TraceServiceServicer): + def __init__(self, file_path): + self._file_path = file_path + self._requests = [] + self._lock = threading.Lock() + + def Export(self, request, context): + req_dict = json_format.MessageToDict( + request, preserving_proto_field_name=True + ) + with self._lock: + self._requests.append(req_dict) + tmp_file = self._file_path + ".tmp" + with open(tmp_file, "w") as f: + json.dump(self._requests, f, indent=2) + f.flush() + os.replace(tmp_file, self._file_path) + + return trace_service_pb2.ExportTraceServiceResponse() + + +def serve(): + parser = argparse.ArgumentParser() + parser.add_argument("--port", type=int, required=True) + parser.add_argument("--file", type=str, required=True) + args = parser.parse_args() + + server = grpc.server(futures.ThreadPoolExecutor(max_workers=10)) + trace_service_pb2_grpc.add_TraceServiceServicer_to_server( + TraceServiceServicer(args.file), server + ) + server.add_insecure_port(f"[::]:{args.port}") + server.start() + print(f"OTLP Collector listening on port {args.port}...") + + def sig_handler(signum, frame): + print(f"OTLP Collector received signal {signum}, stopping server...") + server.stop(0) + + signal.signal(signal.SIGTERM, sig_handler) + signal.signal(signal.SIGINT, sig_handler) + + server.wait_for_termination() + + +if __name__ == "__main__": + serve() diff --git a/test/cpp/interop/run_otel_interop_test.py b/test/cpp/interop/run_otel_interop_test.py new file mode 100755 index 0000000000000..bdf2625bafffb --- /dev/null +++ b/test/cpp/interop/run_otel_interop_test.py @@ -0,0 +1,478 @@ +#!/usr/bin/env python3 +# +# Copyright 2026 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 json +import os +import socket +import subprocess +import sys +import time + + +def get_free_port(): + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + s.bind(("", 0)) + port = s.getsockname()[1] + s.close() + return port + + +def run_cmd(args, desc, env=None, cwd=None): + print(f"Executing: {' '.join(args)} ({desc})") + proc_env = os.environ.copy() + if "CC" not in proc_env and os.path.exists("/usr/bin/gcc"): + proc_env["CC"] = "/usr/bin/gcc" + if "CXX" not in proc_env and os.path.exists("/usr/bin/g++"): + proc_env["CXX"] = "/usr/bin/g++" + if env: + proc_env.update(env) + res = subprocess.run( + args, env=proc_env, cwd=cwd, capture_output=True, text=True + ) + if res.returncode != 0: + print(f"Error executing {desc}:") + print("STDOUT:", res.stdout) + print("STDERR:", res.stderr) + sys.exit(res.returncode) + return res + + +def start_proc(args, env, desc, cwd=None): + print(f"Starting in background: {' '.join(args)} ({desc})") + # Inherit system environment and merge with custom variables + proc_env = os.environ.copy() + if "CC" not in proc_env and os.path.exists("/usr/bin/gcc"): + proc_env["CC"] = "/usr/bin/gcc" + if "CXX" not in proc_env and os.path.exists("/usr/bin/g++"): + proc_env["CXX"] = "/usr/bin/g++" + proc_env.update(env) + return subprocess.Popen( + args, + env=proc_env, + cwd=cwd, + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + ) + + +def verify_spans(spans_file): + start_time = time.time() + all_spans = [] + client_span = None + attempt_span = None + server_span = None + + print("Verifying spans with polling...") + while time.time() - start_time < 5.0: + if not os.path.exists(spans_file): + time.sleep(0.5) + continue + try: + with open(spans_file, "r") as f: + requests = json.load(f) + except (json.JSONDecodeError, IOError): + time.sleep(0.5) + continue + + all_spans = [] + for req in requests: + for resource_spans in req.get("resource_spans", []): + for scope_spans in resource_spans.get("scope_spans", []): + for span in scope_spans.get("spans", []): + all_spans.append(span) + + client_span = next( + ( + s + for s in all_spans + if s.get("name") + in ( + "Sent.grpc.testing.TestService.EmptyCall", + "Sent.grpc.testing.TestService/EmptyCall", + ) + ), + None, + ) + if not client_span: + time.sleep(0.5) + continue + + trace_id = client_span.get("trace_id") + if not trace_id: + time.sleep(0.5) + continue + + attempt_span = next( + ( + s + for s in all_spans + if s.get("name") + in ( + "Attempt.grpc.testing.TestService.EmptyCall", + "Attempt.grpc.testing.TestService/EmptyCall", + ) + and s.get("trace_id") == trace_id + ), + None, + ) + server_span = next( + ( + s + for s in all_spans + if s.get("name") + in ( + "Recv.grpc.testing.TestService.EmptyCall", + "Recv.grpc.testing.TestService/EmptyCall", + ) + and s.get("trace_id") == trace_id + ), + None, + ) + + if not attempt_span or not server_span: + time.sleep(0.5) + continue + + # Found all three spans matching the Trace ID + break + else: + print("Assertion Failed: Timed out waiting for all 3 expected spans.") + return False + + print(f"Collected {len(all_spans)} spans:") + for s in all_spans: + print( + f" Span: '{s.get('name')}' (TraceID: {s.get('trace_id')}, SpanID: {s.get('span_id')}, ParentID: {s.get('parent_span_id')})" + ) + + # 1. Assert all spans share the same Trace ID + trace_id = client_span.get("trace_id") + if attempt_span.get("trace_id") != trace_id: + print( + f"Assertion Failed: Attempt span Trace ID mismatch. Expected {trace_id}, got {attempt_span.get('trace_id')}" + ) + return False + if server_span.get("trace_id") != trace_id: + print( + f"Assertion Failed: Server span Trace ID mismatch. Expected {trace_id}, got {server_span.get('trace_id')}" + ) + return False + + # 2. Assert Parent-Child hierarchy + # Attempt Span must be a child of Client Span + if attempt_span.get("parent_span_id") != client_span.get("span_id"): + print( + f"Assertion Failed: Attempt span is not child of Client span. Expected parent {client_span.get('span_id')}, got {attempt_span.get('parent_span_id')}" + ) + return False + # Server Span must be a child of Attempt Span + if server_span.get("parent_span_id") != attempt_span.get("span_id"): + print( + f"Assertion Failed: Server span is not child of Attempt span. Expected parent {attempt_span.get('span_id')}, got {server_span.get('parent_span_id')}" + ) + return False + + # 3. Verify Attempt Span Attributes + attributes = attempt_span.get("attributes", []) + prev_attempts = None + trans_retry = None + for attr in attributes: + key = attr.get("key") + val_dict = attr.get("value", {}) + if key == "previous-rpc-attempts": + prev_attempts = val_dict.get("int_value") + elif key == "transparent-retry": + trans_retry = val_dict.get("bool_value") + + if prev_attempts is None: + print( + "Assertion Failed: Attribute 'previous-rpc-attempts' not found on Attempt span." + ) + return False + if int(prev_attempts) != 0: + print( + f"Assertion Failed: Attribute 'previous-rpc-attempts' value mismatch. Expected 0, got {prev_attempts}" + ) + return False + if trans_retry is None: + print( + "Assertion Failed: Attribute 'transparent-retry' not found on Attempt span." + ) + return False + if trans_retry is not False: + print( + f"Assertion Failed: Attribute 'transparent-retry' value mismatch. Expected False, got {trans_retry}" + ) + return False + + # 4. Verify Events + attempt_events = [e.get("name") for e in attempt_span.get("events", [])] + client_events = [e.get("name") for e in client_span.get("events", [])] + if ( + "Outbound message" not in attempt_events + and "Outbound message" not in client_events + ): + print( + "Assertion Failed: Event 'Outbound message' not found on Client or Attempt span." + ) + return False + if ( + "Inbound message" not in attempt_events + and "Inbound message" not in client_events + ): + print( + "Assertion Failed: Event 'Inbound message' not found on Client or Attempt span." + ) + return False + + server_events = [e.get("name") for e in server_span.get("events", [])] + if "Inbound message" not in server_events: + print( + "Assertion Failed: Event 'Inbound message' not found on Server span." + ) + return False + if "Outbound message" not in server_events: + print( + "Assertion Failed: Event 'Outbound message' not found on Server span." + ) + return False + + print("All span assertions passed successfully!") + return True + + +def main(): + ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "../../..")) + os.chdir(ROOT) + + parser = argparse.ArgumentParser( + description="Build and run OpenTelemetry Interop tests locally" + ) + parser.add_argument( + "--skip_build", + action="store_true", + help="Skip building targets with Bazel/Gradle", + ) + parser.add_argument( + "--client", + choices=["c++", "java", "python"], + default="c++", + help="Client language (c++, java, or python)", + ) + parser.add_argument( + "--server", + choices=["c++", "java", "python"], + default="c++", + help="Server language (c++, java, or python)", + ) + args = parser.parse_args() + + if not args.skip_build: + if args.client == "c++" or args.server == "c++": + run_cmd( + [ + "./tools/bazel", + "build", + "--macos_minimum_os=11.0", + "//test/cpp/interop:interop_client", + "//test/cpp/interop:interop_server", + "//test/cpp/interop:otlp_collector", + ], + "Building C++ interop targets", + ) + if args.client == "java" or args.server == "java": + run_cmd( + [ + "./gradlew", + ":grpc-interop-testing:installDist", + "-x", + "test", + "-PskipCodegen=true", + ], + "Building Java interop targets", + cwd="../grpc-java", + ) + if args.client == "python" or args.server == "python": + run_cmd( + [ + "./tools/bazel", + "build", + "--macos_minimum_os=11.0", + "//src/python/grpcio_tests/tests/interop:client", + "//src/python/grpcio_tests/tests/interop:server_bin", + ], + "Building Python interop targets", + ) + + collector_port = get_free_port() + server_port = get_free_port() + spans_file = os.path.abspath("captured_spans.json") + if os.path.exists(spans_file): + os.remove(spans_file) + + collector_proc = None + server_proc = None + try: + # Start Collector + collector_proc = start_proc( + [ + "./bazel-bin/test/cpp/interop/otlp_collector", + f"--port={collector_port}", + f"--file={spans_file}", + ], + {}, + "OTLP Collector", + ) + time.sleep(1) # wait for collector to start listening + + # Base env for OTLP exporter + base_env = { + "OTEL_EXPORTER_OTLP_ENDPOINT": f"http://localhost:{collector_port}", + "OTEL_EXPORTER_OTLP_PROTOCOL": "grpc", + "OTEL_TRACES_EXPORTER": "otlp", + "OTEL_METRICS_EXPORTER": "none", + "OTEL_LOGS_EXPORTER": "none", + "GRPC_BAZEL_RUNTIME": "1", + } + + server_env = base_env.copy() + if args.server in ("c++", "java"): + server_env["GRPC_EXPERIMENTAL_ENABLE_OTEL_TRACING"] = "true" + + if args.server == "c++": + server_proc = start_proc( + [ + "./bazel-bin/test/cpp/interop/interop_server", + f"--port={server_port}", + "--enable_opentelemetry=true", + ], + server_env, + "C++ Interop Server", + ) + elif args.server == "java": + server_proc = start_proc( + [ + "../grpc-java/interop-testing/build/install/grpc-interop-testing/bin/test-server", + f"--port={server_port}", + "--use_tls=false", + "--enable_opentelemetry=true", + ], + server_env, + "Java Interop Server", + ) + elif args.server == "python": + server_proc = start_proc( + [ + "./bazel-bin/src/python/grpcio_tests/tests/interop/server_bin", + f"--port={server_port}", + "--use_tls=false", + "--enable_opentelemetry=true", + ], + server_env, + "Python Interop Server", + ) + time.sleep(3) # wait for server to bind and start + + # Run Client + print(f"Running {args.client.upper()} Client...") + client_env = base_env.copy() + if args.client in ("c++", "java"): + client_env["GRPC_EXPERIMENTAL_ENABLE_OTEL_TRACING"] = "true" + + if args.client == "c++": + client_res = run_cmd( + [ + "./bazel-bin/test/cpp/interop/interop_client", + "--server_host=localhost", + f"--server_port={server_port}", + "--test_case=empty_unary", + "--enable_opentelemetry=true", + ], + "Running C++ Interop Client", + env=client_env, + ) + elif args.client == "java": + client_res = run_cmd( + [ + "../grpc-java/interop-testing/build/install/grpc-interop-testing/bin/test-client", + "--server_host=localhost", + f"--server_port={server_port}", + "--test_case=empty_unary", + "--use_tls=false", + "--enable_opentelemetry=true", + ], + "Running Java Interop Client", + env=client_env, + ) + elif args.client == "python": + client_res = run_cmd( + [ + "./bazel-bin/src/python/grpcio_tests/tests/interop/client", + "--server_host=localhost", + f"--server_port={server_port}", + "--test_case=empty_unary", + "--use_tls=false", + "--enable_opentelemetry=true", + ], + "Running Python Interop Client", + env=client_env, + ) + + print("Client finished. Waiting for spans to flush...") + time.sleep(2) + + finally: + # Cleanup server and collector processes + # Cleanup server and collector processes safely + if server_proc: + print("Terminating server...") + try: + server_proc.terminate() + server_proc.wait(timeout=5) + except subprocess.TimeoutExpired: + try: + server_proc.kill() + except Exception: + pass + except Exception as e: + print(f"Error terminating server: {e}") + if collector_proc: + print("Terminating collector...") + try: + collector_proc.terminate() + collector_proc.wait(timeout=5) + except subprocess.TimeoutExpired: + try: + collector_proc.kill() + except Exception: + pass + except Exception as e: + print(f"Error terminating collector: {e}") + + # Perform Span Verifications + success = verify_spans(spans_file) + if success: + print("Test Result: PASSED") + sys.exit(0) + else: + print("Test Result: FAILED") + sys.exit(1) + + +if __name__ == "__main__": + main() diff --git a/tools/doxygen/Doxyfile.c++.internal b/tools/doxygen/Doxyfile.c++.internal index 54566efb1cc24..36b2cf00dfac8 100644 --- a/tools/doxygen/Doxyfile.c++.internal +++ b/tools/doxygen/Doxyfile.c++.internal @@ -3336,6 +3336,7 @@ src/cpp/util/status.cc \ src/cpp/util/string_ref.cc \ src/cpp/util/time_cc.cc \ third_party/upb/upb/generated_code_support.h \ +third_party/utf8_range/utf8_validity.h \ third_party/xxhash/xxhash.h # This tag can be used to specify the character encoding of the source files diff --git a/tools/run_tests/run_interop_tests.py b/tools/run_tests/run_interop_tests.py index d361ab2547ac3..49b2c355a0941 100755 --- a/tools/run_tests/run_interop_tests.py +++ b/tools/run_tests/run_interop_tests.py @@ -21,11 +21,175 @@ import multiprocessing import os import re +import socket import subprocess import sys import time import uuid +_collector_port = None + + +def get_free_port(): + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + s.bind(("", 0)) + port = s.getsockname()[1] + s.close() + return port + + +def verify_tracing_spans(spans_file): + start_time = time.time() + all_spans = [] + client_span = None + attempt_span = None + server_span = None + + print("Verifying tracing spans with polling...") + while time.time() - start_time < 5.0: + if not os.path.exists(spans_file): + time.sleep(0.5) + continue + try: + with open(spans_file, "r") as f: + requests = json.load(f) + except (json.JSONDecodeError, IOError): + time.sleep(0.5) + continue + + all_spans = [] + for req in requests: + for resource_spans in req.get("resource_spans", []): + for scope_spans in resource_spans.get("scope_spans", []): + for span in scope_spans.get("spans", []): + all_spans.append(span) + + client_span = next( + ( + s + for s in all_spans + if s.get("name") + in ( + "Sent.grpc.testing.TestService.EmptyCall", + "Sent.grpc.testing.TestService/EmptyCall", + ) + ), + None, + ) + if not client_span: + time.sleep(0.5) + continue + + trace_id = client_span.get("trace_id") + if not trace_id: + time.sleep(0.5) + continue + + attempt_span = next( + ( + s + for s in all_spans + if s.get("name") + in ( + "Attempt.grpc.testing.TestService.EmptyCall", + "Attempt.grpc.testing.TestService/EmptyCall", + ) + and s.get("trace_id") == trace_id + ), + None, + ) + server_span = next( + ( + s + for s in all_spans + if s.get("name") + in ( + "Recv.grpc.testing.TestService.EmptyCall", + "Recv.grpc.testing.TestService/EmptyCall", + ) + and s.get("trace_id") == trace_id + ), + None, + ) + + if not attempt_span or not server_span: + time.sleep(0.5) + continue + + # Found all three matching spans + break + else: + print( + "Verification Failed: Timed out waiting for all 3 expected spans." + ) + return False + + trace_id = client_span.get("trace_id") + if ( + attempt_span.get("trace_id") != trace_id + or server_span.get("trace_id") != trace_id + ): + print("Verification Failed: Trace ID mismatch.") + return False + + if attempt_span.get("parent_span_id") != client_span.get("span_id"): + print("Verification Failed: Attempt span parent mismatch.") + return False + if server_span.get("parent_span_id") != attempt_span.get("span_id"): + print("Verification Failed: Server span parent mismatch.") + return False + + attributes = attempt_span.get("attributes", []) + prev_attempts = None + trans_retry = None + for attr in attributes: + key = attr.get("key") + val_dict = attr.get("value", {}) + if key == "previous-rpc-attempts": + prev_attempts = val_dict.get("int_value") + elif key == "transparent-retry": + trans_retry = val_dict.get("bool_value") + + if prev_attempts is None or int(prev_attempts) != 0: + print(f"Verification Failed: previous-rpc-attempts attribute mismatch.") + return False + if trans_retry is None or trans_retry is not False: + print(f"Verification Failed: transparent-retry attribute mismatch.") + return False + + attempt_events = [e.get("name") for e in attempt_span.get("events", [])] + client_events = [e.get("name") for e in client_span.get("events", [])] + if ( + "Outbound message" not in attempt_events + and "Outbound message" not in client_events + ): + print( + "Verification Failed: Outbound event not found in Client or Attempt span." + ) + return False + if ( + "Inbound message" not in attempt_events + and "Inbound message" not in client_events + ): + print( + "Verification Failed: Inbound event not found in Client or Attempt span." + ) + return False + + server_events = [e.get("name") for e in server_span.get("events", [])] + if ( + "Inbound message" not in server_events + or "Outbound message" not in server_events + ): + print( + "Verification Failed: Inbound/Outbound event not found in Server span." + ) + return False + + print("OTel Tracing span verification succeeded!") + return True + + import python_utils.dockerjob as dockerjob import python_utils.jobset as jobset import python_utils.report_utils as report_utils @@ -736,6 +900,7 @@ def __str__(self): "special_status_message", "orca_per_rpc", "orca_oob", + "test_unary_rpc_tracing_export", ] _AUTH_TEST_CASES = [ @@ -1054,6 +1219,21 @@ def cloud_to_cloud_jobspec( sys.exit(1) client_test_case = test_case + if test_case == "test_unary_rpc_tracing_export": + client_test_case = "empty_unary" + interop_only_options += ["--enable_opentelemetry=true"] + add_env = add_env.copy() + add_env["GRPC_EXPERIMENTAL_ENABLE_OTEL_TRACING"] = "true" + add_env["OTEL_TRACES_EXPORTER"] = "otlp" + add_env["OTEL_METRICS_EXPORTER"] = "none" + add_env["OTEL_LOGS_EXPORTER"] = "none" + if _collector_port: + collector_host = ( + "host.docker.internal" if docker_image else "localhost" + ) + add_env["OTEL_EXPORTER_OTLP_ENDPOINT"] = ( + f"http://{collector_host}:{_collector_port}" + ) if test_case in _HTTP2_SERVER_TEST_CASES_THAT_USE_GRPC_CLIENTS: client_test_case = _GRPC_CLIENT_TEST_CASES_FOR_HTTP2_SERVER_TEST_CASES[ test_case @@ -1147,6 +1327,18 @@ def server_jobspec( "interop_server_%s" % language.safename ) server_cmd = ["--port=%s" % _DEFAULT_SERVER_PORT] + environ = language.global_env() + if _collector_port and language.safename in ["cxx", "java", "python"]: + server_cmd += ["--enable_opentelemetry=true"] + environ = environ.copy() + environ["GRPC_EXPERIMENTAL_ENABLE_OTEL_TRACING"] = "true" + environ["OTEL_TRACES_EXPORTER"] = "otlp" + environ["OTEL_METRICS_EXPORTER"] = "none" + environ["OTEL_LOGS_EXPORTER"] = "none" + collector_host = "host.docker.internal" if docker_image else "localhost" + environ["OTEL_EXPORTER_OTLP_ENDPOINT"] = ( + f"http://{collector_host}:{_collector_port}" + ) if transport_security == "tls": server_cmd += ["--use_tls=true"] elif transport_security == "alts": @@ -1164,7 +1356,6 @@ def server_jobspec( "--max_concurrent_streams_limit=%d" % max_concurrent_streams_limit ] cmdline = bash_cmdline(language.server_cmd(server_cmd)) - environ = language.global_env() docker_args = ["--name=%s" % container_name] if language.safename == "http2": # we are running the http2 interop server. Open next N ports beginning @@ -1196,6 +1387,11 @@ def server_jobspec( else: docker_args += ["-p", str(_DEFAULT_SERVER_PORT)] + if docker_image: + docker_args = docker_args + [ + "--add-host=host.docker.internal:host-gateway" + ] + docker_cmdline = docker_run_cmdline( cmdline, image=docker_image, @@ -1442,6 +1638,11 @@ def aggregate_http2_results(stdout): nargs="?", help="Upload test results to a specified BQ table.", ) +argp.add_argument( + "--test_case", + type=str, + help="Run a specific test case. If not specified, all test cases will be run.", +) args = argp.parse_args() servers = set( @@ -1567,7 +1768,29 @@ def aggregate_http2_results(stdout): # Start interop servers. server_jobs = {} server_addresses = {} +collector_proc = None try: + if not args.manual_run: + has_cxx = ("c++" in args.language or "all" in args.language) or ( + "c++" in servers + ) + if has_cxx: + _collector_port = get_free_port() + spans_file = os.path.abspath("captured_spans.json") + if os.path.exists(spans_file): + os.remove(spans_file) + collector_cmd = [ + "./bazel-bin/test/cpp/interop/otlp_collector", + f"--port={_collector_port}", + f"--file={spans_file}", + ] + print(f"Starting OTLP collector on port {_collector_port}...") + collector_proc = subprocess.Popen( + collector_cmd, + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + ) + time.sleep(1) for s in servers: lang = str(s) spec = server_jobspec( @@ -1610,6 +1833,10 @@ def aggregate_http2_results(stdout): for server_host_nickname in args.prod_servers: for language in languages: for test_case in _TEST_CASES: + if args.test_case and test_case != args.test_case: + continue + if test_case == "test_unary_rpc_tracing_export": + continue if not test_case in language.unimplemented_test_cases(): if ( not test_case @@ -1654,6 +1881,8 @@ def aggregate_http2_results(stdout): jobs.append(test_job) if args.http2_interop: for test_case in _HTTP2_TEST_CASES: + if args.test_case and test_case != args.test_case: + continue test_job = cloud_to_prod_jobspec( http2Interop, test_case, @@ -1674,6 +1903,8 @@ def aggregate_http2_results(stdout): for server_host_nickname in args.prod_servers: for language in languages: for test_case in _AUTH_TEST_CASES: + if args.test_case and test_case != args.test_case: + continue if ( not args.skip_compute_engine_creds or not compute_engine_creds_required( @@ -1726,6 +1957,15 @@ def aggregate_http2_results(stdout): skip_server = server_language.unimplemented_test_cases_server() for language in languages: for test_case in _TEST_CASES: + if args.test_case and test_case != args.test_case: + continue + if test_case == "test_unary_rpc_tracing_export": + allowed_tracing_languages = ["c++", "java", "python"] + if ( + str(language) not in allowed_tracing_languages + or server_name not in allowed_tracing_languages + ): + continue if not test_case in language.unimplemented_test_cases(): if not test_case in skip_server: test_job = cloud_to_cloud_jobspec( @@ -1742,6 +1982,8 @@ def aggregate_http2_results(stdout): if args.http2_interop: for test_case in _HTTP2_TEST_CASES: + if args.test_case and test_case != args.test_case: + continue if server_name == "go": # TODO(carl-mastrangelo): Reenable after https://github.com/grpc/grpc-go/issues/434 continue @@ -1764,6 +2006,8 @@ def aggregate_http2_results(stdout): for test_case in set(_HTTP2_SERVER_TEST_CASES) - set( _HTTP2_SERVER_TEST_CASES_THAT_USE_GRPC_CLIENTS ): + if args.test_case and test_case != args.test_case: + continue offset = sorted(_HTTP2_SERVER_TEST_CASES).index(test_case) server_port = _DEFAULT_SERVER_PORT + offset if not args.manual_run: @@ -1867,6 +2111,15 @@ def aggregate_http2_results(stdout): maxjobs=args.jobs, skip_jobs=args.manual_run, ) + if not num_failures and not args.manual_run: + has_tracing_test = any( + job.shortname.find("test_unary_rpc_tracing_export") != -1 + for job in jobs + ) + if has_tracing_test: + spans_file = os.path.abspath("captured_spans.json") + if not verify_tracing_spans(spans_file): + num_failures = 1 if args.bq_result_table and resultset: upload_interop_results_to_bq(resultset, args.bq_result_table) if num_failures: @@ -1892,6 +2145,18 @@ def aggregate_http2_results(stdout): else: sys.exit(0) finally: + if collector_proc: + print("Terminating OTLP collector...") + try: + collector_proc.terminate() + collector_proc.wait(timeout=5) + except subprocess.TimeoutExpired: + try: + collector_proc.kill() + except Exception: + pass + except Exception as e: + print(f"Error terminating OTLP collector: {e}") # Check if servers are still running. for server, job in list(server_jobs.items()): if not job.is_running(): diff --git a/tools/run_tests/sanity/check_bazel_workspace.py b/tools/run_tests/sanity/check_bazel_workspace.py index 7bfa8bb679656..07aa0589d99a7 100755 --- a/tools/run_tests/sanity/check_bazel_workspace.py +++ b/tools/run_tests/sanity/check_bazel_workspace.py @@ -52,6 +52,7 @@ "com_google_fuzztest", "io_opencensus_cpp", "io_opentelemetry_cpp", + "opentelemetry_proto", "envoy_api", _BAZEL_SKYLIB_DEP_NAME, _BAZEL_TOOLCHAINS_DEP_NAME, @@ -83,6 +84,7 @@ "com_google_absl", "com_google_fuzztest", "io_opencensus_cpp", + "opentelemetry_proto", _BAZEL_SKYLIB_DEP_NAME, _BAZEL_TOOLCHAINS_DEP_NAME, _BAZEL_COMPDB_DEP_NAME, diff --git a/tools/run_tests/sanity/check_version.py b/tools/run_tests/sanity/check_version.py index 3e483cce06989..e7c5ec57e98b7 100755 --- a/tools/run_tests/sanity/check_version.py +++ b/tools/run_tests/sanity/check_version.py @@ -73,7 +73,7 @@ for tag, value in list(settings.items()): if re.match(r"^[a-z]+_version$", tag): value = Version(value) - if tag != "core_version": + if tag not in ("core_version", "protobuf_version"): if value.major != top_version.major: errors += 1 print( @@ -86,8 +86,8 @@ "minor version mismatch on %s: %d vs %d" % (tag, value.minor, top_version.minor) ) - if not check_version(value): - errors += 1 - print((warning % (tag, value))) + if not check_version(value): + errors += 1 + print((warning % (tag, value))) sys.exit(errors)