From 7f2ecdb3f8af8463c3d291443b1c867cdab8547d Mon Sep 17 00:00:00 2001 From: Xuan Wang Date: Tue, 12 Dec 2023 10:40:59 -0800 Subject: [PATCH] [Python Observability] Building package and add to run_test (#34207) ### Changes in this PR * Refactor and remove some Core/C++ dependencies to simplify Python Observability package build process. * Refactored code to read config at Python layer. * Enable observability build from source. * Add observability to run_test. * Currently it's only enabled in Linux. * Add error handler in run_test loaders. * Current framework will always visit modules in test directory then decide which tests to skip. * Since we're not building Observability for MacOS and Windows this step will fail with error `No module named 'grpc_observability'`. * After the change we'll just skip those modules. * We still have `_sanity_test` to make sure all tests are loaded correctly for each platform. * Remov OC dependency as we're migrating to OTel. * Also removed trace from testing. * Note that trace propagation function was also removed because of this. ### Testing * Passed existing tests. * Tested locally, able to build observability from source using `GRPC_PYTHON_BUILD_WITH_CYTHON=1 pip install .`. Closes #34207 PiperOrigin-RevId: 590258014 --- BUILD | 23 -- requirements.bazel.txt | 3 +- src/python/grpcio_observability/.gitignore | 5 +- src/python/grpcio_observability/MANIFEST.in | 8 + .../_parallel_compile_patch.py | 77 +++++ .../grpc_observability/BUILD.bazel | 17 +- .../grpc_observability/_cyobservability.pyx | 29 +- .../grpc_observability/_gcp_observability.py | 98 +----- .../_observability_config.py | 129 ++++++++ .../_open_census_exporter.py | 11 +- .../grpc_observability/client_call_tracer.cc | 27 +- .../grpc_observability/client_call_tracer.h | 25 +- .../grpc_observability/observability_util.cc | 63 +--- .../grpc_observability/observability_util.h | 42 +-- .../python_census_context.cc | 17 +- .../python_census_context.h | 16 +- .../grpc_observability/rpc_encoding.cc | 29 ++ .../grpc_observability/rpc_encoding.h | 108 +++++++ .../grpc_observability/sampler.cc | 2 +- .../grpc_observability/server_call_tracer.cc | 23 +- .../grpc_observability/server_call_tracer.h | 6 +- .../grpcio_observability/grpc_version.py | 17 + .../make_grpcio_observability.py | 215 +++++++++++++ .../observability_lib_deps.py | 193 ++++++++++++ src/python/grpcio_observability/setup.py | 297 ++++++++++++++++++ src/python/grpcio_tests/tests/_loader.py | 15 +- .../tests/_sanity/_sanity_test.py | 12 + .../tests/observability/__init__.py | 13 + .../observability/_observability_test.py | 120 +------ src/python/grpcio_tests/tests/tests.json | 1 + .../grpc_version.py.template | 19 ++ .../run_tests/helper_scripts/build_python.sh | 15 +- 32 files changed, 1247 insertions(+), 428 deletions(-) create mode 100644 src/python/grpcio_observability/MANIFEST.in create mode 100644 src/python/grpcio_observability/_parallel_compile_patch.py create mode 100644 src/python/grpcio_observability/grpc_observability/_observability_config.py create mode 100644 src/python/grpcio_observability/grpc_observability/rpc_encoding.cc create mode 100644 src/python/grpcio_observability/grpc_observability/rpc_encoding.h create mode 100644 src/python/grpcio_observability/grpc_version.py create mode 100755 src/python/grpcio_observability/make_grpcio_observability.py create mode 100644 src/python/grpcio_observability/observability_lib_deps.py create mode 100644 src/python/grpcio_observability/setup.py create mode 100644 src/python/grpcio_tests/tests/observability/__init__.py create mode 100644 templates/src/python/grpcio_observability/grpc_version.py.template diff --git a/BUILD b/BUILD index 4c668a9e393..6d5130d4ebb 100644 --- a/BUILD +++ b/BUILD @@ -2308,29 +2308,6 @@ grpc_cc_library( ], ) -grpc_cc_library( - name = "grpc_rpc_encoding", - srcs = [ - "src/cpp/ext/filters/census/rpc_encoding.cc", - ], - hdrs = [ - "src/cpp/ext/filters/census/rpc_encoding.h", - ], - external_deps = [ - "absl/base", - "absl/base:core_headers", - "absl/base:endian", - "absl/meta:type_traits", - "absl/status", - "absl/strings", - "absl/time", - ], - language = "c++", - tags = ["nofixdeps"], - visibility = ["@grpc:grpc_python_observability"], - deps = ["gpr_platform"], -) - grpc_cc_library( name = "grpc_opencensus_plugin", srcs = [ diff --git a/requirements.bazel.txt b/requirements.bazel.txt index 532562389fd..0c8ec0cfbd7 100644 --- a/requirements.bazel.txt +++ b/requirements.bazel.txt @@ -3,6 +3,7 @@ coverage==4.5.4 cython==0.29.21 protobuf>=3.5.0.post1, < 4.0dev wheel==0.38.1 +google-auth==1.24.0 oauth2client==4.1.0 requests==2.25.1 urllib3==1.26.5 @@ -13,7 +14,5 @@ gevent==22.08.0 zope.event==4.5.0 setuptools==44.1.1 xds-protos==0.0.11 -opencensus==0.10.0 -opencensus-ext-stackdriver==0.8.0 absl-py==1.4.0 googleapis-common-protos==1.61.0 diff --git a/src/python/grpcio_observability/.gitignore b/src/python/grpcio_observability/.gitignore index 1516c11f4fb..288073fc29b 100644 --- a/src/python/grpcio_observability/.gitignore +++ b/src/python/grpcio_observability/.gitignore @@ -1,6 +1,7 @@ build/ -include/ +grpc_root/ +third_party/ +*.egg-info/ *.c *.cpp -*.egg-info *.so diff --git a/src/python/grpcio_observability/MANIFEST.in b/src/python/grpcio_observability/MANIFEST.in new file mode 100644 index 00000000000..8efdc8f6f21 --- /dev/null +++ b/src/python/grpcio_observability/MANIFEST.in @@ -0,0 +1,8 @@ +graft src/python/grpcio_observability/grpcio_observability.egg-info +graft grpc_observability +graft grpc_root +graft third_party +include _parallel_compile_patch.py +include grpc_version.py +include observability_lib_deps.py +include README.rst diff --git a/src/python/grpcio_observability/_parallel_compile_patch.py b/src/python/grpcio_observability/_parallel_compile_patch.py new file mode 100644 index 00000000000..34ba5379216 --- /dev/null +++ b/src/python/grpcio_observability/_parallel_compile_patch.py @@ -0,0 +1,77 @@ +# Copyright 2023 The gRPC Authors +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +"""Patches the compile() to allow enable parallel compilation of C/C++. + +build_ext has lots of C/C++ files and normally them one by one. +Enabling parallel build helps a lot. +""" + +import os + +try: + BUILD_EXT_COMPILER_JOBS = int( + os.environ["GRPC_PYTHON_BUILD_EXT_COMPILER_JOBS"] + ) +except KeyError: + import multiprocessing + + BUILD_EXT_COMPILER_JOBS = multiprocessing.cpu_count() + + +# monkey-patch for parallel compilation +# TODO(xuanwn): Use a template for this file. +def _parallel_compile( + self, + sources, + output_dir=None, + macros=None, + include_dirs=None, + debug=0, + extra_preargs=None, + extra_postargs=None, + depends=None, +): + # setup the same way as distutils.ccompiler.CCompiler + # https://github.com/python/cpython/blob/31368a4f0e531c19affe2a1becd25fc316bc7501/Lib/distutils/ccompiler.py#L564 + macros, objects, extra_postargs, pp_opts, build = self._setup_compile( + output_dir, macros, include_dirs, sources, depends, extra_postargs + ) + cc_args = self._get_cc_args(pp_opts, debug, extra_preargs) + + def _compile_single_file(obj): + try: + src, ext = build[obj] + except KeyError: + return + self._compile(obj, src, ext, cc_args, extra_postargs, pp_opts) + + # run compilation of individual files in parallel + import multiprocessing.pool + + multiprocessing.pool.ThreadPool(BUILD_EXT_COMPILER_JOBS).map( + _compile_single_file, objects + ) + return objects + + +def monkeypatch_compile_maybe(): + """ + Monkeypatching is dumb, but the build speed gain is worth it. + After python 3.12, we won't find distutils if SETUPTOOLS_USE_DISTUTILS=stdlib. + """ + use_distutils = os.environ.get("SETUPTOOLS_USE_DISTUTILS", "") + if BUILD_EXT_COMPILER_JOBS > 1 and use_distutils != "stdlib": + import distutils.ccompiler # pylint: disable=wrong-import-position + + distutils.ccompiler.CCompiler.compile = _parallel_compile diff --git a/src/python/grpcio_observability/grpc_observability/BUILD.bazel b/src/python/grpcio_observability/grpc_observability/BUILD.bazel index 3be7bf5f970..2fbfb01aa09 100644 --- a/src/python/grpcio_observability/grpc_observability/BUILD.bazel +++ b/src/python/grpcio_observability/grpc_observability/BUILD.bazel @@ -12,12 +12,9 @@ # See the License for the specific language governing permissions and # limitations under the License. -load("@grpc_python_dependencies//:requirements.bzl", "requirement") load("//bazel:cython_library.bzl", "pyx_library") -package(default_visibility = ["//visibility:public"]) - -# TODO(xuanwn): We also need support Python-native build +package(default_visibility = ["//visibility:private"]) cc_library( name = "observability", @@ -25,6 +22,7 @@ cc_library( "client_call_tracer.cc", "observability_util.cc", "python_census_context.cc", + "rpc_encoding.cc", "sampler.cc", "server_call_tracer.cc", ], @@ -33,15 +31,13 @@ cc_library( "constants.h", "observability_util.h", "python_census_context.h", + "rpc_encoding.h", "sampler.h", "server_call_tracer.h", ], includes = ["."], deps = [ - #TODO(xuanwn): Confirm only referenced code is inlcuded in shared object library - "//:grpc", - "//:grpc_rpc_encoding", - "//src/cpp/ext/gcp:observability_config", + "//:grpc_base", ], ) @@ -65,7 +61,7 @@ py_library( "_gcp_observability.py", "_measures.py", "_observability.py", - "_open_census_exporter.py", + "_observability_config.py", "_views.py", ], imports = [ @@ -78,8 +74,5 @@ py_library( ], deps = [ ":cyobservability", - "//src/python/grpcio/grpc:grpcio", - requirement("opencensus"), - requirement("opencensus-ext-stackdriver"), ], ) diff --git a/src/python/grpcio_observability/grpc_observability/_cyobservability.pyx b/src/python/grpcio_observability/grpc_observability/_cyobservability.pyx index 5da61cc7a41..f48674e01e2 100644 --- a/src/python/grpcio_observability/grpc_observability/_cyobservability.pyx +++ b/src/python/grpcio_observability/grpc_observability/_cyobservability.pyx @@ -22,7 +22,7 @@ import os from threading import Thread from typing import List, Mapping, Tuple, Union -import _observability +from grpc_observability import _observability # Time we wait for batch exporting census data # TODO(xuanwn): change interval to a more appropriate number @@ -72,7 +72,6 @@ class MetricsName(enum.Enum): CLIENT_RETRIES_PER_CALL = _CyMetricsName.CY_CLIENT_RETRIES_PER_CALL CLIENT_TRANSPARENT_RETRIES_PER_CALL = _CyMetricsName.CY_CLIENT_TRANSPARENT_RETRIES_PER_CALL CLIENT_RETRY_DELAY_PER_CALL = _CyMetricsName.CY_CLIENT_RETRY_DELAY_PER_CALL - CLIENT_TRANSPORT_LATENCY = _CyMetricsName.CY_CLIENT_TRANSPORT_LATENCY SERVER_SENT_MESSAGES_PER_RPC = _CyMetricsName.CY_SERVER_SENT_MESSAGES_PER_RPC SERVER_SENT_BYTES_PER_RPC = _CyMetricsName.CY_SERVER_SENT_BYTES_PER_RPC SERVER_RECEIVED_MESSAGES_PER_RPC = _CyMetricsName.CY_SERVER_RECEIVED_MESSAGES_PER_RPC @@ -101,28 +100,16 @@ def _start_exporting_thread(object exporter) -> None: GLOBAL_EXPORT_THREAD = Thread(target=_export_census_data, args=(exporter,)) GLOBAL_EXPORT_THREAD.start() +def activate_config(object py_config) -> None: + py_config: "_observability_config.GcpObservabilityConfig" -def set_gcp_observability_config(object py_config) -> bool: - py_config: _gcp_observability.GcpObservabilityPythonConfig - - py_labels = {} - sampling_rate = 0.0 - - cdef cGcpObservabilityConfig c_config = ReadAndActivateObservabilityConfig() - if not c_config.is_valid: - return False - - for label in c_config.labels: - py_labels[_decode(label.key)] = _decode(label.value) - - if PythonCensusTracingEnabled(): - sampling_rate = c_config.cloud_trace.sampling_rate + if (py_config.tracing_enabled): + EnablePythonCensusTracing(True); # Save sampling rate to global sampler. - ProbabilitySampler.Get().SetThreshold(sampling_rate) + ProbabilitySampler.Get().SetThreshold(py_config.sampling_rate) - py_config.set_configuration(_decode(c_config.project_id), sampling_rate, py_labels, - PythonCensusTracingEnabled(), PythonCensusStatsEnabled()) - return True + if (py_config.stats_enabled): + EnablePythonCensusStats(True); def create_client_call_tracer(bytes method_name, bytes trace_id, diff --git a/src/python/grpcio_observability/grpc_observability/_gcp_observability.py b/src/python/grpcio_observability/grpc_observability/_gcp_observability.py index d62653bcd58..58e9a9d47a8 100644 --- a/src/python/grpcio_observability/grpc_observability/_gcp_observability.py +++ b/src/python/grpcio_observability/grpc_observability/_gcp_observability.py @@ -13,20 +13,15 @@ # limitations under the License. from __future__ import annotations -from dataclasses import dataclass -from dataclasses import field import logging -import threading import time -from typing import Any, Mapping, Optional +from typing import Any import grpc -from grpc_observability import _cyobservability # pytype: disable=pyi-error -from grpc_observability._open_census_exporter import CENSUS_UPLOAD_INTERVAL_SECS -from grpc_observability._open_census_exporter import OpenCensusExporter -from opencensus.trace import execution_context -from opencensus.trace import span_context as span_context_module -from opencensus.trace import trace_options as trace_options_module + +# pytype: disable=pyi-error +from grpc_observability import _cyobservability +from grpc_observability import _observability_config _LOGGER = logging.getLogger(__name__) @@ -56,42 +51,6 @@ GRPC_STATUS_CODE_TO_STRING = { grpc.StatusCode.DATA_LOSS: "DATA_LOSS", } -GRPC_SPAN_CONTEXT = "grpc_span_context" - - -@dataclass -class GcpObservabilityPythonConfig: - _singleton = None - _lock: threading.RLock = threading.RLock() - project_id: str = "" - stats_enabled: bool = False - tracing_enabled: bool = False - labels: Optional[Mapping[str, str]] = field(default_factory=dict) - sampling_rate: Optional[float] = 0.0 - - @staticmethod - def get(): - with GcpObservabilityPythonConfig._lock: - if GcpObservabilityPythonConfig._singleton is None: - GcpObservabilityPythonConfig._singleton = ( - GcpObservabilityPythonConfig() - ) - return GcpObservabilityPythonConfig._singleton - - def set_configuration( - self, - project_id: str, - sampling_rate: Optional[float] = 0.0, - labels: Optional[Mapping[str, str]] = None, - tracing_enabled: bool = False, - stats_enabled: bool = False, - ) -> None: - self.project_id = project_id - self.stats_enabled = stats_enabled - self.tracing_enabled = tracing_enabled - self.labels = labels - self.sampling_rate = sampling_rate - # pylint: disable=no-self-use class GCPOpenCensusObservability(grpc._observability.ObservabilityPlugin): @@ -108,25 +67,22 @@ class GCPOpenCensusObservability(grpc._observability.ObservabilityPlugin): exporter: Exporter used to export data. """ - config: GcpObservabilityPythonConfig + config: _observability_config.GcpObservabilityConfig exporter: "grpc_observability.Exporter" - use_open_census_exporter: bool def __init__(self, exporter: "grpc_observability.Exporter" = None): self.exporter = None - self.config = GcpObservabilityPythonConfig.get() - self.use_open_census_exporter = False - config_valid = _cyobservability.set_gcp_observability_config( - self.config - ) - if not config_valid: - raise ValueError("Invalid configuration") + self.config = None + try: + self.config = _observability_config.read_config() + _cyobservability.activate_config(self.config) + except Exception as e: # pylint: disable=broad-except + raise ValueError(f"Reading configuration failed with: {e}") if exporter: self.exporter = exporter else: - self.exporter = OpenCensusExporter(self.config) - self.use_open_census_exporter = True + raise ValueError(f"Please provide an exporter!") if self.config.tracing_enabled: self.set_tracing(True) @@ -156,9 +112,6 @@ class GCPOpenCensusObservability(grpc._observability.ObservabilityPlugin): # TODO(xuanwn): explicit synchronization # https://github.com/grpc/grpc/issues/33262 time.sleep(_cyobservability.CENSUS_EXPORT_BATCH_INTERVAL_SECS) - if self.use_open_census_exporter: - # Sleep so StackDriver can upload data to GCP. - time.sleep(CENSUS_UPLOAD_INTERVAL_SECS) self.set_tracing(False) self.set_stats(False) _cyobservability.observability_deinit() @@ -167,20 +120,10 @@ class GCPOpenCensusObservability(grpc._observability.ObservabilityPlugin): def create_client_call_tracer( self, method_name: bytes ) -> ClientCallTracerCapsule: - grpc_span_context = execution_context.get_opencensus_attr( - GRPC_SPAN_CONTEXT + trace_id = b"TRACE_ID" + capsule = _cyobservability.create_client_call_tracer( + method_name, trace_id ) - if grpc_span_context: - trace_id = grpc_span_context.trace_id.encode("utf8") - parent_span_id = grpc_span_context.span_id.encode("utf8") - capsule = _cyobservability.create_client_call_tracer( - method_name, trace_id, parent_span_id - ) - else: - trace_id = span_context_module.generate_trace_id().encode("utf8") - capsule = _cyobservability.create_client_call_tracer( - method_name, trace_id - ) return capsule def create_server_call_tracer_factory( @@ -197,14 +140,7 @@ class GCPOpenCensusObservability(grpc._observability.ObservabilityPlugin): def save_trace_context( self, trace_id: str, span_id: str, is_sampled: bool ) -> None: - trace_options = trace_options_module.TraceOptions(0) - trace_options.set_enabled(is_sampled) - span_context = span_context_module.SpanContext( - trace_id=trace_id, - span_id=span_id, - trace_options=trace_options, - ) - execution_context.set_opencensus_attr(GRPC_SPAN_CONTEXT, span_context) + pass def record_rpc_latency( self, method: str, rpc_latency: float, status_code: grpc.StatusCode diff --git a/src/python/grpcio_observability/grpc_observability/_observability_config.py b/src/python/grpcio_observability/grpc_observability/_observability_config.py new file mode 100644 index 00000000000..f0105b7d33b --- /dev/null +++ b/src/python/grpcio_observability/grpc_observability/_observability_config.py @@ -0,0 +1,129 @@ +# Copyright 2023 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. +"""Helper to read observability config.""" + +from dataclasses import dataclass +from dataclasses import field +import json +import os +from typing import Mapping, Optional + +GRPC_GCP_OBSERVABILITY_CONFIG_FILE_ENV = "GRPC_GCP_OBSERVABILITY_CONFIG_FILE" +GRPC_GCP_OBSERVABILITY_CONFIG_ENV = "GRPC_GCP_OBSERVABILITY_CONFIG" + + +@dataclass +class GcpObservabilityConfig: + project_id: str = "" + stats_enabled: bool = False + tracing_enabled: bool = False + labels: Optional[Mapping[str, str]] = field(default_factory=dict) + sampling_rate: Optional[float] = 0.0 + + def load_from_string_content(self, config_contents: str) -> None: + """Loads the configuration from a string. + + Args: + config_contents: The configuration string. + + Raises: + ValueError: If the configuration is invalid. + """ + try: + config_json = json.loads(config_contents) + except json.decoder.JSONDecodeError: + raise ValueError("Failed to load Json configuration.") + + if config_json and not isinstance(config_json, dict): + raise ValueError("Found invalid configuration.") + + self.project_id = config_json.get("project_id", "") + self.labels = config_json.get("labels", {}) + self.stats_enabled = "cloud_monitoring" in config_json.keys() + self.tracing_enabled = "cloud_trace" in config_json.keys() + tracing_config = config_json.get("cloud_trace", {}) + self.sampling_rate = tracing_config.get("sampling_rate", 0.0) + + +def read_config() -> GcpObservabilityConfig: + """Reads the GCP observability config from the environment variables. + + Returns: + The GCP observability config. + + Raises: + ValueError: If the configuration is invalid. + """ + config_contents = _get_gcp_observability_config_contents() + config = GcpObservabilityConfig() + config.load_from_string_content(config_contents) + + if not config.project_id: + # Get project ID from GCP environment variables since project ID was not + # set it in the GCP observability config. + config.project_id = _get_gcp_project_id_from_env_var() + if not config.project_id: + # Could not find project ID from GCP environment variables either. + raise ValueError("GCP Project ID not found.") + return config + + +def _get_gcp_project_id_from_env_var() -> Optional[str]: + """Gets the project ID from the GCP environment variables. + + Returns: + The project ID, or an empty string if the project ID could not be found. + """ + + project_id = "" + project_id = os.getenv("GCP_PROJECT") + if project_id: + return project_id + + project_id = os.getenv("GCLOUD_PROJECT") + if project_id: + return project_id + + project_id = os.getenv("GOOGLE_CLOUD_PROJECT") + if project_id: + return project_id + + return project_id + + +def _get_gcp_observability_config_contents() -> str: + """Get the contents of the observability config from environment variable or file. + + Returns: + The content from environment variable. + + Raises: + ValueError: If no configuration content was found. + """ + + contents_str = "" + # First try get config from GRPC_GCP_OBSERVABILITY_CONFIG_FILE_ENV. + config_path = os.getenv(GRPC_GCP_OBSERVABILITY_CONFIG_FILE_ENV) + if config_path: + with open(config_path, "r") as f: + contents_str = f.read() + + # Next, try GRPC_GCP_OBSERVABILITY_CONFIG_ENV env var. + if not contents_str: + contents_str = os.getenv(GRPC_GCP_OBSERVABILITY_CONFIG_ENV) + + if not contents_str: + raise ValueError("Configuration content not found.") + + return contents_str diff --git a/src/python/grpcio_observability/grpc_observability/_open_census_exporter.py b/src/python/grpcio_observability/grpc_observability/_open_census_exporter.py index b3b557ca695..8dbccfdd156 100644 --- a/src/python/grpcio_observability/grpc_observability/_open_census_exporter.py +++ b/src/python/grpcio_observability/grpc_observability/_open_census_exporter.py @@ -14,10 +14,11 @@ from datetime import datetime import os -from typing import Any, List, Mapping, Optional, Tuple +from typing import List, Mapping, Optional, Tuple from google.rpc import code_pb2 from grpc_observability import _observability # pytype: disable=pyi-error +from grpc_observability import _observability_config from grpc_observability import _views from opencensus.common.transports import async_ from opencensus.ext.stackdriver import stats_exporter @@ -38,8 +39,6 @@ from opencensus.trace import time_event from opencensus.trace import trace_options from opencensus.trace import tracer -_gcp_observability = Any # grpc_observability.py imports this module. - # 60s is the default time for open census to call export. CENSUS_UPLOAD_INTERVAL_SECS = int( os.environ.get("GRPC_PYTHON_CENSUS_EXPORT_UPLOAD_INTERVAL_SECS", 20) @@ -61,16 +60,14 @@ class StackDriverAsyncTransport(async_.AsyncTransport): class OpenCensusExporter(_observability.Exporter): - config: "_gcp_observability.GcpObservabilityPythonConfig" + config: _observability_config.GcpObservabilityConfig default_labels: Optional[Mapping[str, str]] project_id: str tracer: Optional[tracer.Tracer] stats_recorder: Optional[StatsRecorder] view_manager: Optional[ViewManager] - def __init__( - self, config: "_gcp_observability.GcpObservabilityPythonConfig" - ): + def __init__(self, config: _observability_config.GcpObservabilityConfig): self.config = config.get() self.default_labels = self.config.labels self.project_id = self.config.project_id diff --git a/src/python/grpcio_observability/grpc_observability/client_call_tracer.cc b/src/python/grpcio_observability/grpc_observability/client_call_tracer.cc index 4af7d423bee..46a091a8fa2 100644 --- a/src/python/grpcio_observability/grpc_observability/client_call_tracer.cc +++ b/src/python/grpcio_observability/grpc_observability/client_call_tracer.cc @@ -12,24 +12,21 @@ // See the License for the specific language governing permissions and // limitations under the License. -#include "src/python/grpcio_observability/grpc_observability/client_call_tracer.h" +#include "client_call_tracer.h" -#include -#include -#include #include #include #include #include "absl/strings/str_cat.h" -#include "absl/strings/string_view.h" #include "absl/time/clock.h" +#include "constants.h" +#include "observability_util.h" +#include "python_census_context.h" #include -#include "src/core/lib/experiments/experiments.h" -#include "src/core/lib/gprpp/sync.h" #include "src/core/lib/slice/slice.h" namespace grpc_observability { @@ -56,7 +53,7 @@ void PythonOpenCensusCallTracer::GenerateContext() {} void PythonOpenCensusCallTracer::RecordAnnotation( absl::string_view annotation) { - if (!context_.SpanContext().IsSampled()) { + if (!context_.GetSpanContext().IsSampled()) { return; } context_.AddSpanAnnotation(annotation); @@ -64,7 +61,7 @@ void PythonOpenCensusCallTracer::RecordAnnotation( void PythonOpenCensusCallTracer::RecordAnnotation( const Annotation& annotation) { - if (!context_.SpanContext().IsSampled()) { + if (!context_.GetSpanContext().IsSampled()) { return; } @@ -93,7 +90,7 @@ PythonOpenCensusCallTracer::~PythonOpenCensusCallTracer() { if (tracing_enabled_) { context_.EndSpan(); if (IsSampled()) { - RecordSpan(context_.Span().ToCensusData()); + RecordSpan(context_.GetSpan().ToCensusData()); } } } @@ -101,7 +98,7 @@ PythonOpenCensusCallTracer::~PythonOpenCensusCallTracer() { PythonCensusContext PythonOpenCensusCallTracer::CreateCensusContextForCallAttempt() { auto context = PythonCensusContext(absl::StrCat("Attempt.", method_), - &(context_.Span()), context_.Labels()); + &(context_.GetSpan()), context_.Labels()); return context; } @@ -274,11 +271,11 @@ void PythonOpenCensusCallTracer::PythonOpenCensusCallAttemptTracer::RecordEnd( if (parent_->tracing_enabled_) { if (status_code_ != absl::StatusCode::kOk) { - context_.Span().SetStatus(StatusCodeToString(status_code_)); + context_.GetSpan().SetStatus(StatusCodeToString(status_code_)); } context_.EndSpan(); if (IsSampled()) { - RecordSpan(context_.Span().ToCensusData()); + RecordSpan(context_.GetSpan().ToCensusData()); } } @@ -289,7 +286,7 @@ void PythonOpenCensusCallTracer::PythonOpenCensusCallAttemptTracer::RecordEnd( void PythonOpenCensusCallTracer::PythonOpenCensusCallAttemptTracer:: RecordAnnotation(absl::string_view annotation) { - if (!context_.SpanContext().IsSampled()) { + if (!context_.GetSpanContext().IsSampled()) { return; } context_.AddSpanAnnotation(annotation); @@ -297,7 +294,7 @@ void PythonOpenCensusCallTracer::PythonOpenCensusCallAttemptTracer:: void PythonOpenCensusCallTracer::PythonOpenCensusCallAttemptTracer:: RecordAnnotation(const Annotation& annotation) { - if (!context_.SpanContext().IsSampled()) { + if (!context_.GetSpanContext().IsSampled()) { return; } diff --git a/src/python/grpcio_observability/grpc_observability/client_call_tracer.h b/src/python/grpcio_observability/grpc_observability/client_call_tracer.h index b94f57bbb76..47a279865af 100644 --- a/src/python/grpcio_observability/grpc_observability/client_call_tracer.h +++ b/src/python/grpcio_observability/grpc_observability/client_call_tracer.h @@ -12,8 +12,8 @@ // See the License for the specific language governing permissions and // limitations under the License. -#ifndef GRPC_PYRHON_OPENCENSUS_CLIENT_CALL_TRACER_H -#define GRPC_PYRHON_OPENCENSUS_CLIENT_CALL_TRACER_H +#ifndef GRPC_PYTHON_OPENCENSUS_CLIENT_CALL_TRACER_H +#define GRPC_PYTHON_OPENCENSUS_CLIENT_CALL_TRACER_H #include @@ -24,16 +24,11 @@ #include "absl/strings/escaping.h" #include "absl/strings/string_view.h" #include "absl/time/time.h" +#include "python_census_context.h" #include #include "src/core/lib/channel/call_tracer.h" -#include "src/core/lib/gprpp/sync.h" -#include "src/core/lib/iomgr/error.h" -#include "src/core/lib/slice/slice_buffer.h" -#include "src/core/lib/transport/metadata_batch.h" -#include "src/core/lib/transport/transport.h" -#include "src/python/grpcio_observability/grpc_observability/python_census_context.h" namespace grpc_observability { @@ -46,15 +41,15 @@ class PythonOpenCensusCallTracer : public grpc_core::ClientCallTracer { bool is_transparent_retry); std::string TraceId() override { return absl::BytesToHexString( - absl::string_view(context_.SpanContext().TraceId())); + absl::string_view(context_.GetSpanContext().TraceId())); } std::string SpanId() override { return absl::BytesToHexString( - absl::string_view(context_.SpanContext().SpanId())); + absl::string_view(context_.GetSpanContext().SpanId())); } - bool IsSampled() override { return context_.SpanContext().IsSampled(); } + bool IsSampled() override { return context_.GetSpanContext().IsSampled(); } void RecordSendInitialMetadata( grpc_metadata_batch* send_initial_metadata) override; @@ -102,15 +97,15 @@ class PythonOpenCensusCallTracer : public grpc_core::ClientCallTracer { std::string TraceId() override { return absl::BytesToHexString( - absl::string_view(context_.SpanContext().TraceId())); + absl::string_view(context_.GetSpanContext().TraceId())); } std::string SpanId() override { return absl::BytesToHexString( - absl::string_view(context_.SpanContext().SpanId())); + absl::string_view(context_.GetSpanContext().SpanId())); } - bool IsSampled() override { return context_.SpanContext().IsSampled(); } + bool IsSampled() override { return context_.GetSpanContext().IsSampled(); } void GenerateContext(); PythonOpenCensusCallAttemptTracer* StartNewAttempt( @@ -139,4 +134,4 @@ class PythonOpenCensusCallTracer : public grpc_core::ClientCallTracer { } // namespace grpc_observability -#endif // GRPC_PYRHON_OPENCENSUS_CLIENT_CALL_TRACER_H +#endif // GRPC_PYTHON_OPENCENSUS_CLIENT_CALL_TRACER_H diff --git a/src/python/grpcio_observability/grpc_observability/observability_util.cc b/src/python/grpcio_observability/grpc_observability/observability_util.cc index 3225543773d..cbbab6a8c93 100644 --- a/src/python/grpcio_observability/grpc_observability/observability_util.cc +++ b/src/python/grpcio_observability/grpc_observability/observability_util.cc @@ -12,10 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -#include "src/python/grpcio_observability/grpc_observability/observability_util.h" - -#include -#include +#include "observability_util.h" #include #include @@ -25,13 +22,13 @@ #include "absl/status/statusor.h" #include "absl/strings/string_view.h" #include "absl/types/optional.h" +#include "client_call_tracer.h" +#include "constants.h" +#include "python_census_context.h" +#include "server_call_tracer.h" #include -#include "src/cpp/ext/gcp/observability_config.h" -#include "src/python/grpcio_observability/grpc_observability/client_call_tracer.h" -#include "src/python/grpcio_observability/grpc_observability/server_call_tracer.h" - namespace grpc_observability { std::queue* g_census_data_buffer; @@ -128,56 +125,6 @@ void AddCensusDataToBuffer(const CensusData& data) { } } -GcpObservabilityConfig ReadAndActivateObservabilityConfig() { - auto config = grpc::internal::GcpObservabilityConfig::ReadFromEnv(); - if (!config.ok()) { - return GcpObservabilityConfig(); - } - - if (!config->cloud_trace.has_value() && - !config->cloud_monitoring.has_value() && - !config->cloud_logging.has_value()) { - return GcpObservabilityConfig(true); - } - - if (config->cloud_trace.has_value()) { - EnablePythonCensusTracing(true); - } - if (config->cloud_monitoring.has_value()) { - EnablePythonCensusStats(true); - } - - std::vector