From bc4ca6083267d3f91e4eb72b9e903f332e54ec7d Mon Sep 17 00:00:00 2001 From: Subham Sinha Date: Wed, 24 Jun 2026 12:03:04 +0530 Subject: [PATCH 1/6] fix(metrics): correct GFE metrics extraction and enable by default --- .../spanner_v1/metrics/metrics_interceptor.py | 365 +++++++++++++++++- .../spanner_v1/metrics/metrics_tracer.py | 79 +++- .../metrics/metrics_tracer_factory.py | 6 + .../metrics/spanner_metrics_tracer_factory.py | 7 +- .../spanner/transports/grpc_asyncio.py | 25 +- .../mockserver_tests/test_gfe_metrics.py | 206 ++++++++++ .../tests/unit/test_metrics_interceptor.py | 11 +- .../tests/unit/test_metrics_tracer.py | 33 ++ 8 files changed, 707 insertions(+), 25 deletions(-) create mode 100644 packages/google-cloud-spanner/tests/mockserver_tests/test_gfe_metrics.py diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py index 3e38c4e0191d..f60a6236d5e3 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py @@ -14,14 +14,19 @@ """Interceptor for collecting Cloud Spanner metrics.""" +import inspect +import logging import re -from typing import Dict +from typing import Any, Dict +import grpc from grpc_interceptor import ClientInterceptor from .constants import GOOGLE_CLOUD_RESOURCE_KEY, SPANNER_METHOD_PREFIX from .spanner_metrics_tracer_factory import SpannerMetricsTracerFactory +logger = logging.getLogger(__name__) + class MetricsInterceptor(ClientInterceptor): """Interceptor that collects metrics for Cloud Spanner operations.""" @@ -67,7 +72,8 @@ def _extract_resource_from_path(metadata: Dict[str, str]) -> Dict[str, str]: resources = MetricsInterceptor._parse_resource_path(path) return resources - def _set_metrics_tracer_attributes(self, resources: Dict[str, str]) -> None: + @staticmethod + def _set_metrics_tracer_attributes(resources: Dict[str, str]) -> None: """ Sets the metric tracer attributes based on the provided resources. @@ -115,17 +121,358 @@ def intercept(self, invoked_method, request_or_iterator, call_details): self._set_metrics_tracer_attributes(resources) ## Format method to be be spanner. - method_name = call_details.method.removeprefix(SPANNER_METHOD_PREFIX).replace( - "/", "." - ) + method_str = call_details.method + if isinstance(method_str, bytes): + method_str = method_str.decode("utf-8") + method_name = method_str.removeprefix(SPANNER_METHOD_PREFIX).replace("/", ".") tracer.set_method(method_name) tracer.record_attempt_start() response = invoked_method(request_or_iterator, call_details) - tracer.record_attempt_completion() - # Process and send GFE metrics if enabled - if tracer.gfe_enabled: - metadata = response.initial_metadata() + return _wrap_response(response, tracer) + + +def _wrap_response(response: Any, tracer: Any) -> Any: + """Wraps the response if it is streaming, or records metrics immediately if unary.""" + if hasattr(response, "__next__"): + return _StreamingResponseWrapper(response, tracer) + else: + # Unary call: execute completion and record metrics immediately + try: + tracer.record_attempt_completion() + metadata = [] + if hasattr(response, "initial_metadata"): + try: + metadata.extend(response.initial_metadata() or []) + except Exception as e: + logger.warning(f"Failed to retrieve initial metadata: {e}") + if hasattr(response, "trailing_metadata"): + try: + metadata.extend(response.trailing_metadata() or []) + except Exception as e: + logger.warning(f"Failed to retrieve trailing metadata: {e}") tracer.record_gfe_metrics(metadata) + except Exception as e: + logger.warning(f"Failed to record metrics: {e}") return response + + +class AsyncMetricsInterceptor( + grpc.aio.UnaryUnaryClientInterceptor, + grpc.aio.UnaryStreamClientInterceptor, + grpc.aio.StreamUnaryClientInterceptor, + grpc.aio.StreamStreamClientInterceptor, +): + """Async Interceptor that collects metrics for Cloud Spanner operations.""" + + async def intercept_unary_unary(self, continuation, client_call_details, request): + return await self._async_intercept(continuation, client_call_details, request) + + async def intercept_unary_stream(self, continuation, client_call_details, request): + return await self._async_intercept(continuation, client_call_details, request) + + async def intercept_stream_unary( + self, continuation, client_call_details, request_iterator + ): + return await self._async_intercept( + continuation, client_call_details, request_iterator + ) + + async def intercept_stream_stream( + self, continuation, client_call_details, request_iterator + ): + return await self._async_intercept( + continuation, client_call_details, request_iterator + ) + + async def _async_intercept( + self, + continuation: Any, + call_details: grpc.ClientCallDetails, + request_or_iterator: Any, + ) -> Any: + # Implementation for async interceptor + factory = SpannerMetricsTracerFactory() + tracer = SpannerMetricsTracerFactory.get_current_tracer() + if tracer is None or not factory.enabled: + return await continuation(call_details, request_or_iterator) + + if not ( + tracer.client_attributes.get("project_id") + and tracer.client_attributes.get("instance_id") + and tracer.client_attributes.get("database") + ): + resources = MetricsInterceptor._extract_resource_from_path( + call_details.metadata + ) + MetricsInterceptor._set_metrics_tracer_attributes(resources) + + method_str = call_details.method + if isinstance(method_str, bytes): + method_str = method_str.decode("utf-8") + method_name = method_str.removeprefix(SPANNER_METHOD_PREFIX).replace("/", ".") + + tracer.set_method(method_name) + tracer.record_attempt_start() + response = await continuation(call_details, request_or_iterator) + + if hasattr(response, "__anext__"): + return _AsyncStreamingResponseWrapper(response, tracer) + else: + return _AsyncUnaryResponseWrapper(response, tracer) + + +class _StreamingResponseWrapper: + """Wrapper for streaming RPC response iterators to defer metrics recording.""" + + def __init__(self, response, tracer): + self._response = response + self._tracer = tracer + self._metrics_recorded = False + self._iterator = None + + def __iter__(self): + self._iterator = iter(self._response) + return self + + def __next__(self): + if self._iterator is None: + self._iterator = iter(self._response) + try: + return next(self._iterator) + except StopIteration: + self._record_metrics() + raise + except Exception: + self._record_metrics() + raise + + def _record_metrics(self): + if self._metrics_recorded: + return + self._metrics_recorded = True + try: + self._tracer.record_attempt_completion() + metadata = [] + if hasattr(self._response, "initial_metadata"): + try: + metadata.extend(self._response.initial_metadata() or []) + except Exception as e: + logger.warning(f"Failed to retrieve initial metadata: {e}") + if hasattr(self._response, "trailing_metadata"): + try: + metadata.extend(self._response.trailing_metadata() or []) + except Exception as e: + logger.warning(f"Failed to retrieve trailing metadata: {e}") + self._tracer.record_gfe_metrics(metadata) + except Exception as e: + logger.warning(f"Failed to record metrics: {e}") + + def __del__(self): + try: + self._record_metrics() + except Exception: + pass + + def __getattr__(self, name): + return getattr(self._response, name) + + +class _AsyncUnaryResponseWrapper(grpc.aio.UnaryUnaryCall): + """Wrapper for async unary RPC response to defer metrics recording until awaited.""" + + def __init__(self, response, tracer): + self._response = response + self._tracer = tracer + self._metrics_recorded = False + + def add_done_callback(self, *args, **kwargs): + return getattr(self._response, "add_done_callback")(*args, **kwargs) + + def cancel(self, *args, **kwargs): + return getattr(self._response, "cancel")(*args, **kwargs) + + def cancelled(self, *args, **kwargs): + return getattr(self._response, "cancelled")(*args, **kwargs) + + def code(self, *args, **kwargs): + return getattr(self._response, "code")(*args, **kwargs) + + def details(self, *args, **kwargs): + return getattr(self._response, "details")(*args, **kwargs) + + def done(self, *args, **kwargs): + return getattr(self._response, "done")(*args, **kwargs) + + def initial_metadata(self, *args, **kwargs): + return getattr(self._response, "initial_metadata")(*args, **kwargs) + + def time_remaining(self, *args, **kwargs): + return getattr(self._response, "time_remaining")(*args, **kwargs) + + def trailing_metadata(self, *args, **kwargs): + return getattr(self._response, "trailing_metadata")(*args, **kwargs) + + def wait_for_connection(self, *args, **kwargs): + return getattr(self._response, "wait_for_connection")(*args, **kwargs) + + def __await__(self): + async def _wait(): + try: + return await self._response + finally: + await self._record_metrics() + + return _wait().__await__() + + async def _record_metrics(self): + if self._metrics_recorded: + return + self._metrics_recorded = True + try: + self._tracer.record_attempt_completion() + metadata = [] + if hasattr(self._response, "initial_metadata"): + try: + res = self._response.initial_metadata() + if inspect.isawaitable(res): + res = await res + metadata.extend(res or []) + except Exception as e: + logger.warning(f"Failed to retrieve initial metadata: {e}") + if hasattr(self._response, "trailing_metadata"): + try: + res = self._response.trailing_metadata() + if inspect.isawaitable(res): + res = await res + metadata.extend(res or []) + except Exception as e: + logger.warning(f"Failed to retrieve trailing metadata: {e}") + self._tracer.record_gfe_metrics(metadata) + except Exception as e: + logger.warning(f"Failed to record metrics: {e}") + + def __del__(self): + if not self._metrics_recorded: + self._metrics_recorded = True + try: + self._tracer.record_attempt_completion() + except Exception: + pass + + def __getattr__(self, name): + return getattr(self._response, name) + + +class _AsyncStreamingResponseWrapper( + grpc.aio.UnaryStreamCall, + grpc.aio.StreamUnaryCall, + grpc.aio.StreamStreamCall, +): + """Wrapper for async streaming RPC response iterators to defer metrics recording.""" + + def __init__(self, response, tracer): + self._response = response + self._tracer = tracer + self._metrics_recorded = False + self._iterator = None + + def add_done_callback(self, *args, **kwargs): + return getattr(self._response, "add_done_callback")(*args, **kwargs) + + def cancel(self, *args, **kwargs): + return getattr(self._response, "cancel")(*args, **kwargs) + + def cancelled(self, *args, **kwargs): + return getattr(self._response, "cancelled")(*args, **kwargs) + + def code(self, *args, **kwargs): + return getattr(self._response, "code")(*args, **kwargs) + + def details(self, *args, **kwargs): + return getattr(self._response, "details")(*args, **kwargs) + + def done(self, *args, **kwargs): + return getattr(self._response, "done")(*args, **kwargs) + + def initial_metadata(self, *args, **kwargs): + return getattr(self._response, "initial_metadata")(*args, **kwargs) + + def time_remaining(self, *args, **kwargs): + return getattr(self._response, "time_remaining")(*args, **kwargs) + + def trailing_metadata(self, *args, **kwargs): + return getattr(self._response, "trailing_metadata")(*args, **kwargs) + + def wait_for_connection(self, *args, **kwargs): + return getattr(self._response, "wait_for_connection")(*args, **kwargs) + + def read(self, *args, **kwargs): + return getattr(self._response, "read")(*args, **kwargs) + + def write(self, *args, **kwargs): + return getattr(self._response, "write")(*args, **kwargs) + + def done_writing(self, *args, **kwargs): + return getattr(self._response, "done_writing")(*args, **kwargs) + + def __aiter__(self): + if hasattr(self._response, "__aiter__"): + self._iterator = self._response.__aiter__() + else: + self._iterator = self._response + return self + + async def __anext__(self): + if self._iterator is None: + if hasattr(self._response, "__aiter__"): + self._iterator = self._response.__aiter__() + else: + self._iterator = self._response + try: + return await self._iterator.__anext__() + except StopAsyncIteration: + await self._record_metrics() + raise + except Exception: + await self._record_metrics() + raise + + async def _record_metrics(self): + if self._metrics_recorded: + return + self._metrics_recorded = True + try: + self._tracer.record_attempt_completion() + metadata = [] + if hasattr(self._response, "initial_metadata"): + try: + res = self._response.initial_metadata() + if inspect.isawaitable(res): + res = await res + metadata.extend(res or []) + except Exception as e: + logger.warning(f"Failed to retrieve initial metadata: {e}") + if hasattr(self._response, "trailing_metadata"): + try: + res = self._response.trailing_metadata() + if inspect.isawaitable(res): + res = await res + metadata.extend(res or []) + except Exception as e: + logger.warning(f"Failed to retrieve trailing metadata: {e}") + self._tracer.record_gfe_metrics(metadata) + except Exception as e: + logger.warning(f"Failed to record metrics: {e}") + + def __del__(self): + if not self._metrics_recorded: + self._metrics_recorded = True + try: + self._tracer.record_attempt_completion() + except Exception: + pass + + def __getattr__(self, name): + return getattr(self._response, name) diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py index f79869948f99..27d33660e736 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py @@ -19,8 +19,9 @@ while the helper classes provide additional functionality and context for the metrics being traced. """ +import re from datetime import datetime -from typing import Dict +from typing import Any, Dict, Optional from grpc import StatusCode @@ -198,6 +199,8 @@ def __init__( instrument_operation_counter: "Counter", client_attributes: Dict[str, str], gfe_enabled: bool = False, + instrument_gfe_latency: Optional["Histogram"] = None, + instrument_gfe_missing_header_count: Optional["Counter"] = None, ): """ Initialize a MetricsTracer instance with the given parameters. @@ -214,6 +217,8 @@ def __init__( instrument_operation_counter (Counter): Instrument for counting operations. client_attributes (Dict[str, str]): Dictionary of client attributes used for metrics tracing. gfe_enabled (bool, optional): Indicates if GFE metrics are enabled. Defaults to False. + instrument_gfe_latency (Histogram, optional): Instrument for measuring GFE latency. + instrument_gfe_missing_header_count (Counter, optional): Instrument for counting missing GFE headers. """ self.current_op = MetricOpTracer() self._client_attributes = client_attributes @@ -221,8 +226,10 @@ def __init__( self._instrument_attempt_counter = instrument_attempt_counter self._instrument_operation_latency = instrument_operation_latency self._instrument_operation_counter = instrument_operation_counter + self._instrument_gfe_latency = instrument_gfe_latency + self._instrument_gfe_missing_header_count = instrument_gfe_missing_header_count self.enabled = enabled - self.gfe_enabled = gfe_enabled + self.gfe_enabled = True @staticmethod def _get_ms_time_diff(start: datetime, end: datetime) -> float: @@ -399,7 +406,11 @@ def record_gfe_latency(self, latency: int) -> None: Args: latency (int): The latency duration to be recorded. """ - if not self.enabled or not HAS_OPENTELEMETRY_INSTALLED or not self.gfe_enabled: + if ( + not self.enabled + or not HAS_OPENTELEMETRY_INSTALLED + or not getattr(self, "_instrument_gfe_latency", None) + ): return self._instrument_gfe_latency.record( amount=latency, attributes=self.client_attributes @@ -409,12 +420,72 @@ def record_gfe_missing_header_count(self) -> None: """ Increments the counter for missing GFE headers. """ - if not self.enabled or not HAS_OPENTELEMETRY_INSTALLED or not self.gfe_enabled: + if ( + not self.enabled + or not HAS_OPENTELEMETRY_INSTALLED + or not getattr(self, "_instrument_gfe_missing_header_count", None) + ): return self._instrument_gfe_missing_header_count.add( amount=1, attributes=self.client_attributes ) + @staticmethod + def extract_gfe_latency(metadata: Any) -> Optional[int]: + """ + Extracts the GFE latency value (in milliseconds) from response metadata. + """ + if not metadata: + return None + + header_vals = [] + if isinstance(metadata, dict): + for key, val in metadata.items(): + if key and str(key).lower() in ("server-timing", "server_timing"): + if isinstance(val, (list, tuple)): + header_vals.extend(val) + else: + header_vals.append(val) + elif isinstance(metadata, (list, tuple)): + for item in metadata: + if isinstance(item, (list, tuple)) and len(item) == 2: + key, val = item + if key and str(key).lower() in ("server-timing", "server_timing"): + if isinstance(val, (list, tuple)): + header_vals.extend(val) + else: + header_vals.append(val) + + for header_val in header_vals: + if not header_val: + continue + if isinstance(header_val, bytes): + try: + header_val = header_val.decode("utf-8") + except Exception: + header_val = str(header_val) + elif not isinstance(header_val, str): + header_val = str(header_val) + match = re.search(r"gfet4t7;\s*dur=([0-9.]+)", header_val) + if match: + try: + return int(float(match.group(1))) + except ValueError: + pass + return None + + def record_gfe_metrics(self, metadata: Any) -> None: + """ + Extracts and records GFE metrics from the RPC response metadata. + """ + if not self.enabled or not HAS_OPENTELEMETRY_INSTALLED: + return + latency = self.extract_gfe_latency(metadata) + if latency is not None: + self.record_gfe_latency(latency) + else: + self.record_gfe_missing_header_count() + def _create_operation_otel_attributes(self) -> dict: """ Create additional attributes for operation metrics tracing. diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py index f22d285c9750..029dddfaa15a 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py @@ -85,6 +85,7 @@ def __init__(self, enabled: bool, service_name: str): project (str): The project ID for the monitored resource. """ self.enabled = enabled + self.gfe_enabled = True self._create_metric_instruments(service_name) self._client_attributes = {} @@ -268,6 +269,11 @@ def create_metrics_tracer(self) -> MetricsTracer: instrument_operation_latency=self._instrument_operation_latency, instrument_operation_counter=self._instrument_operation_counter, client_attributes=self._client_attributes.copy(), + gfe_enabled=True, + instrument_gfe_latency=getattr(self, "_instrument_gfe_latency", None), + instrument_gfe_missing_header_count=getattr( + self, "_instrument_gfe_missing_header_count", None + ), ) return metrics_tracer diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/spanner_metrics_tracer_factory.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/spanner_metrics_tracer_factory.py index 6fc5956582c1..7886e555f120 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/spanner_metrics_tracer_factory.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/spanner_metrics_tracer_factory.py @@ -51,9 +51,7 @@ class SpannerMetricsTracerFactory(MetricsTracerFactory): "current_metrics_tracer", default=None ) - def __new__( - cls, enabled: bool = True, gfe_enabled: bool = False - ) -> "SpannerMetricsTracerFactory": + def __new__(cls, enabled: bool = True) -> "SpannerMetricsTracerFactory": """ Create a new instance of SpannerMetricsTracerFactory if it doesn't already exist. @@ -63,7 +61,6 @@ def __new__( Args: enabled (bool): A flag indicating whether metrics tracing is enabled. Defaults to True. - gfe_enabled (bool): A flag indicating whether GFE metrics are enabled. Defaults to False. Returns: SpannerMetricsTracerFactory: The singleton instance of SpannerMetricsTracerFactory. @@ -83,7 +80,7 @@ def __new__( cls._generate_client_hash(client_uid) ) cls._metrics_tracer_factory.set_location(_get_cloud_region()) - cls._metrics_tracer_factory.gfe_enabled = gfe_enabled + cls._metrics_tracer_factory.gfe_enabled = True if cls._metrics_tracer_factory.enabled != enabled: cls._metrics_tracer_factory.enabled = enabled diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/services/spanner/transports/grpc_asyncio.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/services/spanner/transports/grpc_asyncio.py index c688b31eefc4..c56ab0112d23 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/services/spanner/transports/grpc_asyncio.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/services/spanner/transports/grpc_asyncio.py @@ -32,7 +32,10 @@ from google.protobuf.json_format import MessageToJson from grpc.experimental import aio # type: ignore -from google.cloud.spanner_v1.metrics.metrics_interceptor import MetricsInterceptor +from google.cloud.spanner_v1.metrics.metrics_interceptor import ( + AsyncMetricsInterceptor, + MetricsInterceptor, +) from google.cloud.spanner_v1.types import ( commit_response, location, @@ -327,6 +330,26 @@ def __init__( ], ) + if metrics_interceptor is not None: + self._metrics_interceptor = AsyncMetricsInterceptor() + # Attach interceptor directly since grpc.aio does not provide intercept_channel. + if hasattr(self._grpc_channel, "_unary_unary_interceptors"): + self._grpc_channel._unary_unary_interceptors.append( + self._metrics_interceptor + ) + if hasattr(self._grpc_channel, "_unary_stream_interceptors"): + self._grpc_channel._unary_stream_interceptors.append( + self._metrics_interceptor + ) + if hasattr(self._grpc_channel, "_stream_unary_interceptors"): + self._grpc_channel._stream_unary_interceptors.append( + self._metrics_interceptor + ) + if hasattr(self._grpc_channel, "_stream_stream_interceptors"): + self._grpc_channel._stream_stream_interceptors.append( + self._metrics_interceptor + ) + self._interceptor = _LoggingClientAIOInterceptor() self._grpc_channel._unary_unary_interceptors.append(self._interceptor) self._logged_channel = self._grpc_channel diff --git a/packages/google-cloud-spanner/tests/mockserver_tests/test_gfe_metrics.py b/packages/google-cloud-spanner/tests/mockserver_tests/test_gfe_metrics.py new file mode 100644 index 000000000000..6d4cf51bff4e --- /dev/null +++ b/packages/google-cloud-spanner/tests/mockserver_tests/test_gfe_metrics.py @@ -0,0 +1,206 @@ +# Copyright 2025 Google LLC All rights reserved. +# +# 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 os +from unittest import mock + +import grpc +from google.api_core.client_options import ClientOptions +from google.auth.credentials import AnonymousCredentials +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.metrics.export import InMemoryMetricReader + +import google.cloud.spanner_v1.client as client_mod +from google.cloud.spanner_v1 import Client +from google.cloud.spanner_v1.metrics.metrics_interceptor import MetricsInterceptor +from google.cloud.spanner_v1.metrics.spanner_metrics_tracer_factory import ( + SpannerMetricsTracerFactory, +) +from google.cloud.spanner_v1.pool import FixedSizePool +from tests.mockserver_tests.mock_server_test_base import ( + MockServerTestBase, + add_select1_result, +) + + +class TestGFEMetricsIntegration(MockServerTestBase): + def setUp(self): + super().setUp() + os.environ["SPANNER_DISABLE_BUILTIN_METRICS"] = "false" + SpannerMetricsTracerFactory._metrics_tracer_factory = None + client_mod._metrics_monitor_initialized = False + + def tearDown(self): + super().tearDown() + os.environ["SPANNER_DISABLE_BUILTIN_METRICS"] = "true" + SpannerMetricsTracerFactory._metrics_tracer_factory = None + client_mod._metrics_monitor_initialized = False + + def test_gfe_metrics_exported(self): + add_select1_result() + reader = InMemoryMetricReader() + meter_provider = MeterProvider(metric_readers=[reader]) + + orig_call = grpc._channel._UnaryStreamMultiCallable.__call__ + orig_initial_metadata = grpc._channel._MultiThreadedRendezvous.initial_metadata + orig_trailing_metadata = ( + grpc._channel._MultiThreadedRendezvous.trailing_metadata + ) + + def custom_initial_metadata(self): + mocked = getattr(self, "_is_execute_streaming_sql_mock", False) + if mocked: + return (("server-timing", "gfet4t7; dur=55"),) + return orig_initial_metadata(self) + + def custom_trailing_metadata(self): + mocked = getattr(self, "_is_execute_streaming_sql_mock", False) + if mocked: + return (("server-timing", "gfet4t7; dur=55"),) + return orig_trailing_metadata(self) + + def custom_call(self_callable, request, *args, **kwargs): + method = getattr(self_callable, "_method", b"") + method_str = method.decode("utf-8") if isinstance(method, bytes) else method + response = orig_call(self_callable, request, *args, **kwargs) + if "ExecuteStreamingSql" in method_str: + response._is_execute_streaming_sql_mock = True + return response + + try: + with ( + mock.patch( + "google.cloud.spanner_v1.metrics.metrics_tracer_factory.get_meter_provider", + return_value=meter_provider, + ), + mock.patch( + "google.cloud.spanner_v1.client.MeterProvider", + return_value=meter_provider, + ), + mock.patch( + "google.cloud.spanner_v1.client._get_spanner_emulator_host", + return_value=None, + ), + mock.patch( + "grpc._channel._UnaryStreamMultiCallable.__call__", + custom_call, + ), + mock.patch( + "grpc._channel._MultiThreadedRendezvous.initial_metadata", + custom_initial_metadata, + ), + mock.patch( + "grpc._channel._MultiThreadedRendezvous.trailing_metadata", + custom_trailing_metadata, + ), + ): + client = Client( + project="p", + credentials=AnonymousCredentials(), + client_options=ClientOptions( + api_endpoint="localhost:" + str(MockServerTestBase.port), + ), + ) + instance = client.instance("test-instance") + database = instance.database( + "test-database", + pool=FixedSizePool(size=10), + enable_interceptors_in_tests=True, + ) + database._interceptors.append(MetricsInterceptor()) + database._spanner_api = ( + None # Force recreation with the new interceptor + ) + + with database.snapshot() as snapshot: + results = snapshot.execute_sql("select 1") + # Consume the streaming results to complete the stream + list(results) + + metric_data = reader.get_metrics_data() + self.assertIsNotNone(metric_data) + metrics = { + metric.name: metric + for rm in metric_data.resource_metrics + for sm in rm.scope_metrics + for metric in sm.metrics + } + + self.assertIn("gfe_latency", metrics, f"Metrics: {list(metrics.keys())}") + gfe_metric = metrics["gfe_latency"] + point = next(iter(gfe_metric.data.data_points)) + self.assertEqual(point.sum, 55) + + finally: + pass + + def test_gfe_missing_header_count_exported(self): + add_select1_result() + reader = InMemoryMetricReader() + meter_provider = MeterProvider(metric_readers=[reader]) + + try: + with ( + mock.patch( + "google.cloud.spanner_v1.metrics.metrics_tracer_factory.get_meter_provider", + return_value=meter_provider, + ), + mock.patch( + "google.cloud.spanner_v1.client.MeterProvider", + return_value=meter_provider, + ), + mock.patch( + "google.cloud.spanner_v1.client._get_spanner_emulator_host", + return_value=None, + ), + ): + client = Client( + project="p", + credentials=AnonymousCredentials(), + client_options=ClientOptions( + api_endpoint="localhost:" + str(MockServerTestBase.port), + ), + ) + instance = client.instance("test-instance") + database = instance.database( + "test-database", + pool=FixedSizePool(size=10), + enable_interceptors_in_tests=True, + ) + database._interceptors.append(MetricsInterceptor()) + database._spanner_api = ( + None # Force recreation with the new interceptor + ) + + with database.snapshot() as snapshot: + results = snapshot.execute_sql("select 1") + list(results) + + metric_data = reader.get_metrics_data() + self.assertIsNotNone(metric_data) + metrics = { + metric.name: metric + for rm in metric_data.resource_metrics + for sm in rm.scope_metrics + for metric in sm.metrics + } + + self.assertIn( + "gfe_missing_header_count", metrics, f"Metrics: {list(metrics.keys())}" + ) + missing_metric = metrics["gfe_missing_header_count"] + point = next(iter(missing_metric.data.data_points)) + self.assertGreaterEqual(point.value, 1) + finally: + pass diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py index 6e091860b425..efa080191c9e 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py @@ -12,7 +12,7 @@ # See the License for the specific language governing permissions and # limitations under the License. -from unittest.mock import MagicMock +from unittest.mock import MagicMock, Mock import pytest @@ -41,7 +41,7 @@ def __init__(self): self.project = None self.instance = None self.database = None - self.gfe_enabled = False + self.gfe_enabled = True self.record_attempt_start = MagicMock() self.record_attempt_completion = MagicMock() self.set_method = MagicMock() @@ -99,10 +99,8 @@ def test_set_metrics_tracer_attributes(interceptor, mock_tracer_ctx): def test_intercept_with_tracer(interceptor, mock_tracer_ctx): # mock_tracer_ctx fixture sets the ContextVar - mock_tracer_ctx.gfe_enabled = False - - invoked_response = MagicMock() - invoked_response.initial_metadata.return_value = {} + invoked_response = Mock() + invoked_response.initial_metadata.return_value = [] mock_invoked_method = MagicMock(return_value=invoked_response) call_details = MagicMock( @@ -119,4 +117,5 @@ def test_intercept_with_tracer(interceptor, mock_tracer_ctx): assert response == invoked_response mock_tracer_ctx.record_attempt_start.assert_called() mock_tracer_ctx.record_attempt_completion.assert_called_once() + mock_tracer_ctx.record_gfe_metrics.assert_called_once() mock_invoked_method.assert_called_once_with("request", call_details) diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py index 90b2f2f511f9..4769974f0c8a 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py @@ -264,3 +264,36 @@ def test_record_gfe_missing_header_count(metrics_tracer): metrics_tracer.record_gfe_missing_header_count() assert mock_gfe_missing_header_count.add.call_count == 1 # Should not increment metrics_tracer.enabled = True # Reset for next test + + +def test_extract_gfe_latency(): + # Valid trailing metadata list of tuples + metadata_list = [("server-timing", "gfet4t7; dur=123")] + assert MetricsTracer.extract_gfe_latency(metadata_list) == 123 + + # Valid metadata dict + metadata_dict = {"server-timing": "gfet4t7; dur=456"} + assert MetricsTracer.extract_gfe_latency(metadata_dict) == 456 + + # Missing header + assert MetricsTracer.extract_gfe_latency([("other-header", "val")]) is None + assert MetricsTracer.extract_gfe_latency(None) is None + + +def test_record_gfe_metrics(metrics_tracer): + mock_gfe_latency = mock.create_autospec(Histogram, instance=True) + mock_gfe_missing = mock.create_autospec(Counter, instance=True) + metrics_tracer._instrument_gfe_latency = mock_gfe_latency + metrics_tracer._instrument_gfe_missing_header_count = mock_gfe_missing + metrics_tracer.gfe_enabled = True + + # With header + metrics_tracer.record_gfe_metrics([("server-timing", "gfet4t7; dur=88")]) + assert mock_gfe_latency.record.call_count == 1 + assert mock_gfe_latency.record.call_args[1]["amount"] == 88 + assert mock_gfe_missing.add.call_count == 0 + + # Without header + metrics_tracer.record_gfe_metrics([("other", "1")]) + assert mock_gfe_latency.record.call_count == 1 + assert mock_gfe_missing.add.call_count == 1 From ee696dbc4bca8d45e0b299d1638ce368c2c5130e Mon Sep 17 00:00:00 2001 From: Subham Sinha <35077434+sinhasubham@users.noreply.github.com> Date: Tue, 14 Jul 2026 12:36:22 +0530 Subject: [PATCH 2/6] feat(metrics): add AFE latency metrics and simplify metadata extraction --- .../cloud/spanner_v1/metrics/constants.py | 8 +- .../spanner_v1/metrics/metrics_interceptor.py | 30 +----- .../spanner_v1/metrics/metrics_tracer.py | 95 +++++++++++++++++++ .../metrics/metrics_tracer_factory.py | 20 ++++ ...fe_metrics.py => test_frontend_metrics.py} | 22 ++++- .../tests/unit/test_metrics_interceptor.py | 2 + .../tests/unit/test_metrics_tracer.py | 71 ++++++++++++++ 7 files changed, 216 insertions(+), 32 deletions(-) rename packages/google-cloud-spanner/tests/mockserver_tests/{test_gfe_metrics.py => test_frontend_metrics.py} (89%) diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/constants.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/constants.py index a5f709881b12..2e213dabccc2 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/constants.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/constants.py @@ -58,13 +58,19 @@ METRIC_NAME_ATTEMPT_LATENCIES = "attempt_latencies" METRIC_NAME_OPERATION_COUNT = "operation_count" METRIC_NAME_ATTEMPT_COUNT = "attempt_count" -METRIC_NAME_GFE_LATENCY = "gfe_latency" +METRIC_NAME_GFE_LATENCY = "gfe_latencies" METRIC_NAME_GFE_MISSING_HEADER_COUNT = "gfe_missing_header_count" +METRIC_NAME_AFE_LATENCY = "afe_latencies" +METRIC_NAME_AFE_MISSING_HEADER_COUNT = "afe_missing_header_count" METRIC_NAMES = [ METRIC_NAME_OPERATION_LATENCIES, METRIC_NAME_ATTEMPT_LATENCIES, METRIC_NAME_OPERATION_COUNT, METRIC_NAME_ATTEMPT_COUNT, + METRIC_NAME_GFE_LATENCY, + METRIC_NAME_GFE_MISSING_HEADER_COUNT, + METRIC_NAME_AFE_LATENCY, + METRIC_NAME_AFE_MISSING_HEADER_COUNT, ] METRIC_EXPORT_INTERVAL_MS = 60000 # 1 Minute diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py index f60a6236d5e3..6e865dd1272b 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py @@ -147,12 +147,8 @@ def _wrap_response(response: Any, tracer: Any) -> Any: metadata.extend(response.initial_metadata() or []) except Exception as e: logger.warning(f"Failed to retrieve initial metadata: {e}") - if hasattr(response, "trailing_metadata"): - try: - metadata.extend(response.trailing_metadata() or []) - except Exception as e: - logger.warning(f"Failed to retrieve trailing metadata: {e}") tracer.record_gfe_metrics(metadata) + tracer.record_afe_metrics(metadata) except Exception as e: logger.warning(f"Failed to record metrics: {e}") return response @@ -260,12 +256,8 @@ def _record_metrics(self): metadata.extend(self._response.initial_metadata() or []) except Exception as e: logger.warning(f"Failed to retrieve initial metadata: {e}") - if hasattr(self._response, "trailing_metadata"): - try: - metadata.extend(self._response.trailing_metadata() or []) - except Exception as e: - logger.warning(f"Failed to retrieve trailing metadata: {e}") self._tracer.record_gfe_metrics(metadata) + self._tracer.record_afe_metrics(metadata) except Exception as e: logger.warning(f"Failed to record metrics: {e}") @@ -341,15 +333,8 @@ async def _record_metrics(self): metadata.extend(res or []) except Exception as e: logger.warning(f"Failed to retrieve initial metadata: {e}") - if hasattr(self._response, "trailing_metadata"): - try: - res = self._response.trailing_metadata() - if inspect.isawaitable(res): - res = await res - metadata.extend(res or []) - except Exception as e: - logger.warning(f"Failed to retrieve trailing metadata: {e}") self._tracer.record_gfe_metrics(metadata) + self._tracer.record_afe_metrics(metadata) except Exception as e: logger.warning(f"Failed to record metrics: {e}") @@ -454,15 +439,8 @@ async def _record_metrics(self): metadata.extend(res or []) except Exception as e: logger.warning(f"Failed to retrieve initial metadata: {e}") - if hasattr(self._response, "trailing_metadata"): - try: - res = self._response.trailing_metadata() - if inspect.isawaitable(res): - res = await res - metadata.extend(res or []) - except Exception as e: - logger.warning(f"Failed to retrieve trailing metadata: {e}") self._tracer.record_gfe_metrics(metadata) + self._tracer.record_afe_metrics(metadata) except Exception as e: logger.warning(f"Failed to record metrics: {e}") diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py index 27d33660e736..0146b35ad14d 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py @@ -185,6 +185,8 @@ class should not have any knowledge about the observability framework used for m _instrument_operation_latency: "Histogram" _instrument_gfe_latency: "Histogram" _instrument_gfe_missing_header_count: "Counter" + _instrument_afe_latency: "Histogram" + _instrument_afe_missing_header_count: "Counter" current_op: MetricOpTracer enabled: bool gfe_enabled: bool @@ -201,6 +203,8 @@ def __init__( gfe_enabled: bool = False, instrument_gfe_latency: Optional["Histogram"] = None, instrument_gfe_missing_header_count: Optional["Counter"] = None, + instrument_afe_latency: Optional["Histogram"] = None, + instrument_afe_missing_header_count: Optional["Counter"] = None, ): """ Initialize a MetricsTracer instance with the given parameters. @@ -219,6 +223,8 @@ def __init__( gfe_enabled (bool, optional): Indicates if GFE metrics are enabled. Defaults to False. instrument_gfe_latency (Histogram, optional): Instrument for measuring GFE latency. instrument_gfe_missing_header_count (Counter, optional): Instrument for counting missing GFE headers. + instrument_afe_latency (Histogram, optional): Instrument for measuring AFE latency. + instrument_afe_missing_header_count (Counter, optional): Instrument for counting missing AFE headers. """ self.current_op = MetricOpTracer() self._client_attributes = client_attributes @@ -228,6 +234,8 @@ def __init__( self._instrument_operation_counter = instrument_operation_counter self._instrument_gfe_latency = instrument_gfe_latency self._instrument_gfe_missing_header_count = instrument_gfe_missing_header_count + self._instrument_afe_latency = instrument_afe_latency + self._instrument_afe_missing_header_count = instrument_afe_missing_header_count self.enabled = enabled self.gfe_enabled = True @@ -430,6 +438,37 @@ def record_gfe_missing_header_count(self) -> None: amount=1, attributes=self.client_attributes ) + def record_afe_latency(self, latency: int) -> None: + """ + Records the AFE latency using the Histogram instrument. + + Args: + latency (int): The latency duration to be recorded. + """ + if ( + not self.enabled + or not HAS_OPENTELEMETRY_INSTALLED + or not getattr(self, "_instrument_afe_latency", None) + ): + return + self._instrument_afe_latency.record( + amount=latency, attributes=self.client_attributes + ) + + def record_afe_missing_header_count(self) -> None: + """ + Increments the counter for missing AFE headers. + """ + if ( + not self.enabled + or not HAS_OPENTELEMETRY_INSTALLED + or not getattr(self, "_instrument_afe_missing_header_count", None) + ): + return + self._instrument_afe_missing_header_count.add( + amount=1, attributes=self.client_attributes + ) + @staticmethod def extract_gfe_latency(metadata: Any) -> Optional[int]: """ @@ -486,6 +525,62 @@ def record_gfe_metrics(self, metadata: Any) -> None: else: self.record_gfe_missing_header_count() + @staticmethod + def extract_afe_latency(metadata: Any) -> Optional[int]: + """ + Extracts the AFE latency value (in milliseconds) from response metadata. + """ + if not metadata: + return None + + header_vals = [] + if isinstance(metadata, dict): + for key, val in metadata.items(): + if key and str(key).lower() in ("server-timing", "server_timing"): + if isinstance(val, (list, tuple)): + header_vals.extend(val) + else: + header_vals.append(val) + elif isinstance(metadata, (list, tuple)): + for item in metadata: + if isinstance(item, (list, tuple)) and len(item) == 2: + key, val = item + if key and str(key).lower() in ("server-timing", "server_timing"): + if isinstance(val, (list, tuple)): + header_vals.extend(val) + else: + header_vals.append(val) + + for header_val in header_vals: + if not header_val: + continue + if isinstance(header_val, bytes): + try: + header_val = header_val.decode("utf-8") + except Exception: + header_val = str(header_val) + elif not isinstance(header_val, str): + header_val = str(header_val) + match = re.search(r"afe(?:t4t7)?;\s*dur=([0-9.]+)", header_val) + if match: + try: + return int(float(match.group(1))) + except ValueError: + pass + return None + + def record_afe_metrics(self, metadata: Any) -> None: + """ + Extracts and records AFE metrics from the RPC response metadata. + """ + if not self.enabled or not HAS_OPENTELEMETRY_INSTALLED: + return + latency = self.extract_afe_latency(metadata) + if latency is not None: + self.record_afe_latency(latency) + else: + self.record_afe_missing_header_count() + def _create_operation_otel_attributes(self) -> dict: """ Create additional attributes for operation metrics tracing. diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py index 029dddfaa15a..ee6360f47aef 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py @@ -22,6 +22,8 @@ METRIC_LABEL_KEY_CLIENT_UID, METRIC_LABEL_KEY_DATABASE, METRIC_LABEL_KEY_DIRECT_PATH_ENABLED, + METRIC_NAME_AFE_LATENCY, + METRIC_NAME_AFE_MISSING_HEADER_COUNT, METRIC_NAME_ATTEMPT_COUNT, METRIC_NAME_ATTEMPT_LATENCIES, METRIC_NAME_GFE_LATENCY, @@ -57,6 +59,8 @@ class MetricsTracerFactory: _instrument_operation_counter: "Counter" _instrument_gfe_latency: "Histogram" _instrument_gfe_missing_header_count: "Counter" + _instrument_afe_latency: "Histogram" + _instrument_afe_missing_header_count: "Counter" _client_attributes: Dict[str, str] @property @@ -274,6 +278,10 @@ def create_metrics_tracer(self) -> MetricsTracer: instrument_gfe_missing_header_count=getattr( self, "_instrument_gfe_missing_header_count", None ), + instrument_afe_latency=getattr(self, "_instrument_afe_latency", None), + instrument_afe_missing_header_count=getattr( + self, "_instrument_afe_missing_header_count", None + ), ) return metrics_tracer @@ -331,3 +339,15 @@ def _create_metric_instruments(self, service_name: str) -> None: unit="1", description="GFE missing header count.", ) + + self._instrument_afe_latency = meter.create_histogram( + name=METRIC_NAME_AFE_LATENCY, + unit="ms", + description="AFE Latency.", + ) + + self._instrument_afe_missing_header_count = meter.create_counter( + name=METRIC_NAME_AFE_MISSING_HEADER_COUNT, + unit="1", + description="AFE missing header count.", + ) diff --git a/packages/google-cloud-spanner/tests/mockserver_tests/test_gfe_metrics.py b/packages/google-cloud-spanner/tests/mockserver_tests/test_frontend_metrics.py similarity index 89% rename from packages/google-cloud-spanner/tests/mockserver_tests/test_gfe_metrics.py rename to packages/google-cloud-spanner/tests/mockserver_tests/test_frontend_metrics.py index 6d4cf51bff4e..24be3159924c 100644 --- a/packages/google-cloud-spanner/tests/mockserver_tests/test_gfe_metrics.py +++ b/packages/google-cloud-spanner/tests/mockserver_tests/test_frontend_metrics.py @@ -34,7 +34,7 @@ ) -class TestGFEMetricsIntegration(MockServerTestBase): +class TestFrontendMetricsIntegration(MockServerTestBase): def setUp(self): super().setUp() os.environ["SPANNER_DISABLE_BUILTIN_METRICS"] = "false" @@ -61,13 +61,13 @@ def test_gfe_metrics_exported(self): def custom_initial_metadata(self): mocked = getattr(self, "_is_execute_streaming_sql_mock", False) if mocked: - return (("server-timing", "gfet4t7; dur=55"),) + return (("server-timing", "gfet4t7; dur=55, afe; dur=23"),) return orig_initial_metadata(self) def custom_trailing_metadata(self): mocked = getattr(self, "_is_execute_streaming_sql_mock", False) if mocked: - return (("server-timing", "gfet4t7; dur=55"),) + return (("server-timing", "gfet4t7; dur=55, afe; dur=23"),) return orig_trailing_metadata(self) def custom_call(self_callable, request, *args, **kwargs): @@ -137,11 +137,16 @@ def custom_call(self_callable, request, *args, **kwargs): for metric in sm.metrics } - self.assertIn("gfe_latency", metrics, f"Metrics: {list(metrics.keys())}") - gfe_metric = metrics["gfe_latency"] + self.assertIn("gfe_latencies", metrics, f"Metrics: {list(metrics.keys())}") + gfe_metric = metrics["gfe_latencies"] point = next(iter(gfe_metric.data.data_points)) self.assertEqual(point.sum, 55) + self.assertIn("afe_latencies", metrics, f"Metrics: {list(metrics.keys())}") + afe_metric = metrics["afe_latencies"] + point = next(iter(afe_metric.data.data_points)) + self.assertEqual(point.sum, 23) + finally: pass @@ -202,5 +207,12 @@ def test_gfe_missing_header_count_exported(self): missing_metric = metrics["gfe_missing_header_count"] point = next(iter(missing_metric.data.data_points)) self.assertGreaterEqual(point.value, 1) + + self.assertIn( + "afe_missing_header_count", metrics, f"Metrics: {list(metrics.keys())}" + ) + afe_missing_metric = metrics["afe_missing_header_count"] + afe_point = next(iter(afe_missing_metric.data.data_points)) + self.assertGreaterEqual(afe_point.value, 1) finally: pass diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py index efa080191c9e..5cfc46143ac4 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py @@ -46,6 +46,7 @@ def __init__(self): self.record_attempt_completion = MagicMock() self.set_method = MagicMock() self.record_gfe_metrics = MagicMock() + self.record_afe_metrics = MagicMock() self.set_project = MagicMock() self.set_instance = MagicMock() self.set_database = MagicMock() @@ -118,4 +119,5 @@ def test_intercept_with_tracer(interceptor, mock_tracer_ctx): mock_tracer_ctx.record_attempt_start.assert_called() mock_tracer_ctx.record_attempt_completion.assert_called_once() mock_tracer_ctx.record_gfe_metrics.assert_called_once() + mock_tracer_ctx.record_afe_metrics.assert_called_once() mock_invoked_method.assert_called_once_with("request", call_details) diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py index 4769974f0c8a..7e79fd124f5c 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py @@ -297,3 +297,74 @@ def test_record_gfe_metrics(metrics_tracer): metrics_tracer.record_gfe_metrics([("other", "1")]) assert mock_gfe_latency.record.call_count == 1 assert mock_gfe_missing.add.call_count == 1 + + +def test_record_afe_latency(metrics_tracer): + mock_afe_latency = mock.create_autospec(Histogram, instance=True) + metrics_tracer._instrument_afe_latency = mock_afe_latency + metrics_tracer.gfe_enabled = True + + metrics_tracer.record_afe_latency(100) + assert mock_afe_latency.record.call_count == 1 + assert mock_afe_latency.record.call_args[1]["amount"] == 100 + assert ( + mock_afe_latency.record.call_args[1]["attributes"] + == metrics_tracer.client_attributes + ) + + metrics_tracer.enabled = False + metrics_tracer.record_afe_latency(200) + assert mock_afe_latency.record.call_count == 1 + metrics_tracer.enabled = True + + +def test_record_afe_missing_header_count(metrics_tracer): + mock_afe_missing = mock.create_autospec(Counter, instance=True) + metrics_tracer._instrument_afe_missing_header_count = mock_afe_missing + metrics_tracer.gfe_enabled = True + + metrics_tracer.record_afe_missing_header_count() + assert mock_afe_missing.add.call_count == 1 + assert mock_afe_missing.add.call_args[1]["amount"] == 1 + assert ( + mock_afe_missing.add.call_args[1]["attributes"] + == metrics_tracer.client_attributes + ) + + metrics_tracer.enabled = False + metrics_tracer.record_afe_missing_header_count() + assert mock_afe_missing.add.call_count == 1 + metrics_tracer.enabled = True + + +def test_extract_afe_latency(): + # Valid trailing metadata list of tuples + metadata_list = [("server-timing", "afe; dur=123")] + assert MetricsTracer.extract_afe_latency(metadata_list) == 123 + + # Valid metadata dict + metadata_dict = {"server-timing": "afet4t7; dur=456"} + assert MetricsTracer.extract_afe_latency(metadata_dict) == 456 + + # Missing header + assert MetricsTracer.extract_afe_latency([("other-header", "val")]) is None + assert MetricsTracer.extract_afe_latency(None) is None + + +def test_record_afe_metrics(metrics_tracer): + mock_afe_latency = mock.create_autospec(Histogram, instance=True) + mock_afe_missing = mock.create_autospec(Counter, instance=True) + metrics_tracer._instrument_afe_latency = mock_afe_latency + metrics_tracer._instrument_afe_missing_header_count = mock_afe_missing + metrics_tracer.gfe_enabled = True + + # With header + metrics_tracer.record_afe_metrics([("server-timing", "afe; dur=88")]) + assert mock_afe_latency.record.call_count == 1 + assert mock_afe_latency.record.call_args[1]["amount"] == 88 + assert mock_afe_missing.add.call_count == 0 + + # Without header + metrics_tracer.record_afe_metrics([("other", "1")]) + assert mock_afe_latency.record.call_count == 1 + assert mock_afe_missing.add.call_count == 1 From 30e3122e7ac36e5357d7f7006b81fec31bc0ef92 Mon Sep 17 00:00:00 2001 From: Subham Sinha <35077434+sinhasubham@users.noreply.github.com> Date: Mon, 20 Jul 2026 21:09:12 +0530 Subject: [PATCH 3/6] fix(spanner): align GFE/AFE metric names and instruments with backend specs --- .../google/cloud/spanner_v1/metrics/README.md | 8 +- .../cloud/spanner_v1/metrics/constants.py | 8 +- .../spanner_v1/metrics/metrics_interceptor.py | 33 ++-- .../spanner_v1/metrics/metrics_tracer.py | 180 ++++++++---------- .../metrics/metrics_tracer_factory.py | 32 ++-- .../tests/unit/test_metrics_interceptor.py | 6 +- .../tests/unit/test_metrics_tracer.py | 106 +++++------ 7 files changed, 168 insertions(+), 205 deletions(-) diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/README.md b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/README.md index 9619715c8531..179e9b2d9464 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/README.md +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/README.md @@ -1,4 +1,4 @@ -# Custom Metric Exporter +# Custom Metric Exporter The custom metric exporter, as defined in [metrics_exporter.py](./metrics_exporter.py), is designed to work in conjunction with OpenTelemetry and the Spanner client. It converts data into its protobuf equivalent and sends it to Google Cloud Monitoring. ## Filtering Criteria @@ -10,8 +10,10 @@ The exporter filters metrics based on the following conditions, utilizing values * `attempt_count` * `operation_latencies` * `operation_count` - * `gfe_latency` - * `gfe_missing_header_count` + * `gfe_latencies` + * `gfe_connectivity_error_count` + * `afe_latencies` + * `afe_connectivity_error_count` ## Service Endpoint The exporter sends metrics to the Google Cloud Monitoring [service endpoint](https://cloud.google.com/python/docs/reference/monitoring/latest/google.cloud.monitoring_v3.services.metric_service.MetricServiceClient#google_cloud_monitoring_v3_services_metric_service_MetricServiceClient_create_service_time_series), distinct from the regular client endpoint. This service endpoint operates under a different quota limit than the user endpoint and features an additional server-side filter that only permits a predefined set of metrics to pass through. diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/constants.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/constants.py index 2e213dabccc2..fa5f5ca4d98d 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/constants.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/constants.py @@ -59,18 +59,18 @@ METRIC_NAME_OPERATION_COUNT = "operation_count" METRIC_NAME_ATTEMPT_COUNT = "attempt_count" METRIC_NAME_GFE_LATENCY = "gfe_latencies" -METRIC_NAME_GFE_MISSING_HEADER_COUNT = "gfe_missing_header_count" +METRIC_NAME_GFE_CONNECTIVITY_ERROR_COUNT = "gfe_connectivity_error_count" METRIC_NAME_AFE_LATENCY = "afe_latencies" -METRIC_NAME_AFE_MISSING_HEADER_COUNT = "afe_missing_header_count" +METRIC_NAME_AFE_CONNECTIVITY_ERROR_COUNT = "afe_connectivity_error_count" METRIC_NAMES = [ METRIC_NAME_OPERATION_LATENCIES, METRIC_NAME_ATTEMPT_LATENCIES, METRIC_NAME_OPERATION_COUNT, METRIC_NAME_ATTEMPT_COUNT, METRIC_NAME_GFE_LATENCY, - METRIC_NAME_GFE_MISSING_HEADER_COUNT, + METRIC_NAME_GFE_CONNECTIVITY_ERROR_COUNT, METRIC_NAME_AFE_LATENCY, - METRIC_NAME_AFE_MISSING_HEADER_COUNT, + METRIC_NAME_AFE_CONNECTIVITY_ERROR_COUNT, ] METRIC_EXPORT_INTERVAL_MS = 60000 # 1 Minute diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py index 6e865dd1272b..24ff9aa11ec4 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py @@ -52,22 +52,27 @@ def _parse_resource_path(path: str) -> dict: return {} @staticmethod - def _extract_resource_from_path(metadata: Dict[str, str]) -> Dict[str, str]: + def _extract_resource_from_path(metadata: Any) -> Dict[str, str]: """ Extracts resource information from the metadata based on the path. - This method iterates through the metadata dictionary to find the first tuple containing the key 'google-cloud-resource-prefix'. It then extracts the path from this tuple and parses it to extract project, instance, and database information using the _parse_resource_path method. - Args: - metadata (Dict[str, str]): A dictionary containing metadata information. + metadata (Any): A sequence or dictionary containing metadata information. Returns: Dict[str, str]: A dictionary containing extracted project, instance, and database information. """ - # Extract resource info from the first metadata tuple containing :path - path = next( - (value for key, value in metadata if key == GOOGLE_CLOUD_RESOURCE_KEY), "" - ) + if not metadata: + return {} + + items = metadata.items() if isinstance(metadata, dict) else metadata + path = "" + + for key, value in items: + key_str = key.decode("utf-8") if isinstance(key, bytes) else str(key) + if key_str == GOOGLE_CLOUD_RESOURCE_KEY: + path = value.decode("utf-8") if isinstance(value, bytes) else str(value) + break resources = MetricsInterceptor._parse_resource_path(path) return resources @@ -147,8 +152,7 @@ def _wrap_response(response: Any, tracer: Any) -> Any: metadata.extend(response.initial_metadata() or []) except Exception as e: logger.warning(f"Failed to retrieve initial metadata: {e}") - tracer.record_gfe_metrics(metadata) - tracer.record_afe_metrics(metadata) + tracer.record_front_end_metrics(metadata) except Exception as e: logger.warning(f"Failed to record metrics: {e}") return response @@ -256,8 +260,7 @@ def _record_metrics(self): metadata.extend(self._response.initial_metadata() or []) except Exception as e: logger.warning(f"Failed to retrieve initial metadata: {e}") - self._tracer.record_gfe_metrics(metadata) - self._tracer.record_afe_metrics(metadata) + self._tracer.record_front_end_metrics(metadata) except Exception as e: logger.warning(f"Failed to record metrics: {e}") @@ -333,8 +336,7 @@ async def _record_metrics(self): metadata.extend(res or []) except Exception as e: logger.warning(f"Failed to retrieve initial metadata: {e}") - self._tracer.record_gfe_metrics(metadata) - self._tracer.record_afe_metrics(metadata) + self._tracer.record_front_end_metrics(metadata) except Exception as e: logger.warning(f"Failed to record metrics: {e}") @@ -439,8 +441,7 @@ async def _record_metrics(self): metadata.extend(res or []) except Exception as e: logger.warning(f"Failed to retrieve initial metadata: {e}") - self._tracer.record_gfe_metrics(metadata) - self._tracer.record_afe_metrics(metadata) + self._tracer.record_front_end_metrics(metadata) except Exception as e: logger.warning(f"Failed to record metrics: {e}") diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py index 0146b35ad14d..a357682357ec 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py @@ -184,9 +184,9 @@ class should not have any knowledge about the observability framework used for m _instrument_operation_counter: "Counter" _instrument_operation_latency: "Histogram" _instrument_gfe_latency: "Histogram" - _instrument_gfe_missing_header_count: "Counter" + _instrument_gfe_connectivity_error_count: "Counter" _instrument_afe_latency: "Histogram" - _instrument_afe_missing_header_count: "Counter" + _instrument_afe_connectivity_error_count: "Counter" current_op: MetricOpTracer enabled: bool gfe_enabled: bool @@ -200,11 +200,11 @@ def __init__( instrument_operation_latency: "Histogram", instrument_operation_counter: "Counter", client_attributes: Dict[str, str], + instrument_gfe_latency: "Histogram", + instrument_gfe_connectivity_error_count: "Counter", + instrument_afe_latency: "Histogram", + instrument_afe_connectivity_error_count: "Counter", gfe_enabled: bool = False, - instrument_gfe_latency: Optional["Histogram"] = None, - instrument_gfe_missing_header_count: Optional["Counter"] = None, - instrument_afe_latency: Optional["Histogram"] = None, - instrument_afe_missing_header_count: Optional["Counter"] = None, ): """ Initialize a MetricsTracer instance with the given parameters. @@ -221,10 +221,10 @@ def __init__( instrument_operation_counter (Counter): Instrument for counting operations. client_attributes (Dict[str, str]): Dictionary of client attributes used for metrics tracing. gfe_enabled (bool, optional): Indicates if GFE metrics are enabled. Defaults to False. - instrument_gfe_latency (Histogram, optional): Instrument for measuring GFE latency. - instrument_gfe_missing_header_count (Counter, optional): Instrument for counting missing GFE headers. - instrument_afe_latency (Histogram, optional): Instrument for measuring AFE latency. - instrument_afe_missing_header_count (Counter, optional): Instrument for counting missing AFE headers. + instrument_gfe_latency (Histogram): Instrument for measuring GFE latency. + instrument_gfe_connectivity_error_count (Counter): Instrument for counting GFE connectivity errors. + instrument_afe_latency (Histogram): Instrument for measuring AFE latency. + instrument_afe_connectivity_error_count (Counter): Instrument for counting AFE connectivity errors. """ self.current_op = MetricOpTracer() self._client_attributes = client_attributes @@ -233,9 +233,13 @@ def __init__( self._instrument_operation_latency = instrument_operation_latency self._instrument_operation_counter = instrument_operation_counter self._instrument_gfe_latency = instrument_gfe_latency - self._instrument_gfe_missing_header_count = instrument_gfe_missing_header_count + self._instrument_gfe_connectivity_error_count = ( + instrument_gfe_connectivity_error_count + ) self._instrument_afe_latency = instrument_afe_latency - self._instrument_afe_missing_header_count = instrument_afe_missing_header_count + self._instrument_afe_connectivity_error_count = ( + instrument_afe_connectivity_error_count + ) self.enabled = enabled self.gfe_enabled = True @@ -424,17 +428,17 @@ def record_gfe_latency(self, latency: int) -> None: amount=latency, attributes=self.client_attributes ) - def record_gfe_missing_header_count(self) -> None: + def record_gfe_connectivity_error_count(self) -> None: """ - Increments the counter for missing GFE headers. + Increments the counter for GFE connectivity errors. """ if ( not self.enabled or not HAS_OPENTELEMETRY_INSTALLED - or not getattr(self, "_instrument_gfe_missing_header_count", None) + or not getattr(self, "_instrument_gfe_connectivity_error_count", None) ): return - self._instrument_gfe_missing_header_count.add( + self._instrument_gfe_connectivity_error_count.add( amount=1, attributes=self.client_attributes ) @@ -455,101 +459,52 @@ def record_afe_latency(self, latency: int) -> None: amount=latency, attributes=self.client_attributes ) - def record_afe_missing_header_count(self) -> None: + def record_afe_connectivity_error_count(self) -> None: """ - Increments the counter for missing AFE headers. + Increments the counter for AFE connectivity errors. """ if ( not self.enabled or not HAS_OPENTELEMETRY_INSTALLED - or not getattr(self, "_instrument_afe_missing_header_count", None) + or not getattr(self, "_instrument_afe_connectivity_error_count", None) ): return - self._instrument_afe_missing_header_count.add( + self._instrument_afe_connectivity_error_count.add( amount=1, attributes=self.client_attributes ) @staticmethod - def extract_gfe_latency(metadata: Any) -> Optional[int]: + def extract_front_end_latencies( + metadata: Any, + ) -> tuple[Optional[int], Optional[int]]: """ - Extracts the GFE latency value (in milliseconds) from response metadata. + Extracts both GFE and AFE latency values (in milliseconds) from response metadata. """ if not metadata: - return None + return None, None - header_vals = [] if isinstance(metadata, dict): - for key, val in metadata.items(): - if key and str(key).lower() in ("server-timing", "server_timing"): - if isinstance(val, (list, tuple)): - header_vals.extend(val) - else: - header_vals.append(val) + items = metadata.items() elif isinstance(metadata, (list, tuple)): - for item in metadata: - if isinstance(item, (list, tuple)) and len(item) == 2: - key, val = item - if key and str(key).lower() in ("server-timing", "server_timing"): - if isinstance(val, (list, tuple)): - header_vals.extend(val) - else: - header_vals.append(val) - - for header_val in header_vals: - if not header_val: - continue - if isinstance(header_val, bytes): - try: - header_val = header_val.decode("utf-8") - except Exception: - header_val = str(header_val) - elif not isinstance(header_val, str): - header_val = str(header_val) - match = re.search(r"gfet4t7;\s*dur=([0-9.]+)", header_val) - if match: - try: - return int(float(match.group(1))) - except ValueError: - pass - return None - - def record_gfe_metrics(self, metadata: Any) -> None: - """ - Extracts and records GFE metrics from the RPC response metadata. - """ - if not self.enabled or not HAS_OPENTELEMETRY_INSTALLED: - return - latency = self.extract_gfe_latency(metadata) - if latency is not None: - self.record_gfe_latency(latency) + items = [ + item + for item in metadata + if isinstance(item, (list, tuple)) and len(item) == 2 + ] else: - self.record_gfe_missing_header_count() - - @staticmethod - def extract_afe_latency(metadata: Any) -> Optional[int]: - """ - Extracts the AFE latency value (in milliseconds) from response metadata. - """ - if not metadata: - return None + items = [] header_vals = [] - if isinstance(metadata, dict): - for key, val in metadata.items(): - if key and str(key).lower() in ("server-timing", "server_timing"): - if isinstance(val, (list, tuple)): - header_vals.extend(val) - else: - header_vals.append(val) - elif isinstance(metadata, (list, tuple)): - for item in metadata: - if isinstance(item, (list, tuple)) and len(item) == 2: - key, val = item - if key and str(key).lower() in ("server-timing", "server_timing"): - if isinstance(val, (list, tuple)): - header_vals.extend(val) - else: - header_vals.append(val) + for key, val in items: + key_str = key.decode("utf-8") if isinstance(key, bytes) else str(key) + if key_str and key_str.lower() in ("server-timing", "server_timing"): + if isinstance(val, (list, tuple)): + header_vals.extend(val) + else: + header_vals.append(val) + + gfe_latency = None + afe_latency = None for header_val in header_vals: if not header_val: @@ -561,25 +516,42 @@ def extract_afe_latency(metadata: Any) -> Optional[int]: header_val = str(header_val) elif not isinstance(header_val, str): header_val = str(header_val) - match = re.search(r"afe(?:t4t7)?;\s*dur=([0-9.]+)", header_val) - if match: - try: - return int(float(match.group(1))) - except ValueError: - pass - return None - def record_afe_metrics(self, metadata: Any) -> None: + if gfe_latency is None: + match = re.search(r"gfet4t7;\s*dur=([0-9.]+)", header_val) + if match: + try: + gfe_latency = int(float(match.group(1))) + except ValueError: + pass + + if afe_latency is None: + match = re.search(r"afe(?:t4t7)?;\s*dur=([0-9.]+)", header_val) + if match: + try: + afe_latency = int(float(match.group(1))) + except ValueError: + pass + + return gfe_latency, afe_latency + + def record_front_end_metrics(self, metadata: Any) -> None: """ - Extracts and records AFE metrics from the RPC response metadata. + Extracts and records both GFE and AFE metrics from the RPC response metadata. """ if not self.enabled or not HAS_OPENTELEMETRY_INSTALLED: return - latency = self.extract_afe_latency(metadata) - if latency is not None: - self.record_afe_latency(latency) + gfe_latency, afe_latency = self.extract_front_end_latencies(metadata) + + if gfe_latency is not None: + self.record_gfe_latency(gfe_latency) + else: + self.record_gfe_connectivity_error_count() + + if afe_latency is not None: + self.record_afe_latency(afe_latency) else: - self.record_afe_missing_header_count() + self.record_afe_connectivity_error_count() def _create_operation_otel_attributes(self) -> dict: """ diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py index ee6360f47aef..804fb46f1faf 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py @@ -22,12 +22,12 @@ METRIC_LABEL_KEY_CLIENT_UID, METRIC_LABEL_KEY_DATABASE, METRIC_LABEL_KEY_DIRECT_PATH_ENABLED, + METRIC_NAME_AFE_CONNECTIVITY_ERROR_COUNT, METRIC_NAME_AFE_LATENCY, - METRIC_NAME_AFE_MISSING_HEADER_COUNT, METRIC_NAME_ATTEMPT_COUNT, METRIC_NAME_ATTEMPT_LATENCIES, + METRIC_NAME_GFE_CONNECTIVITY_ERROR_COUNT, METRIC_NAME_GFE_LATENCY, - METRIC_NAME_GFE_MISSING_HEADER_COUNT, METRIC_NAME_OPERATION_COUNT, METRIC_NAME_OPERATION_LATENCIES, MONITORED_RES_LABEL_KEY_CLIENT_HASH, @@ -58,9 +58,9 @@ class MetricsTracerFactory: _instrument_operation_latency: "Histogram" _instrument_operation_counter: "Counter" _instrument_gfe_latency: "Histogram" - _instrument_gfe_missing_header_count: "Counter" + _instrument_gfe_connectivity_error_count: "Counter" _instrument_afe_latency: "Histogram" - _instrument_afe_missing_header_count: "Counter" + _instrument_afe_connectivity_error_count: "Counter" _client_attributes: Dict[str, str] @property @@ -274,14 +274,10 @@ def create_metrics_tracer(self) -> MetricsTracer: instrument_operation_counter=self._instrument_operation_counter, client_attributes=self._client_attributes.copy(), gfe_enabled=True, - instrument_gfe_latency=getattr(self, "_instrument_gfe_latency", None), - instrument_gfe_missing_header_count=getattr( - self, "_instrument_gfe_missing_header_count", None - ), - instrument_afe_latency=getattr(self, "_instrument_afe_latency", None), - instrument_afe_missing_header_count=getattr( - self, "_instrument_afe_missing_header_count", None - ), + instrument_gfe_latency=self._instrument_gfe_latency, + instrument_gfe_connectivity_error_count=self._instrument_gfe_connectivity_error_count, + instrument_afe_latency=self._instrument_afe_latency, + instrument_afe_connectivity_error_count=self._instrument_afe_connectivity_error_count, ) return metrics_tracer @@ -334,10 +330,10 @@ def _create_metric_instruments(self, service_name: str) -> None: description="GFE Latency.", ) - self._instrument_gfe_missing_header_count = meter.create_counter( - name=METRIC_NAME_GFE_MISSING_HEADER_COUNT, + self._instrument_gfe_connectivity_error_count = meter.create_counter( + name=METRIC_NAME_GFE_CONNECTIVITY_ERROR_COUNT, unit="1", - description="GFE missing header count.", + description="GFE connectivity error count.", ) self._instrument_afe_latency = meter.create_histogram( @@ -346,8 +342,8 @@ def _create_metric_instruments(self, service_name: str) -> None: description="AFE Latency.", ) - self._instrument_afe_missing_header_count = meter.create_counter( - name=METRIC_NAME_AFE_MISSING_HEADER_COUNT, + self._instrument_afe_connectivity_error_count = meter.create_counter( + name=METRIC_NAME_AFE_CONNECTIVITY_ERROR_COUNT, unit="1", - description="AFE missing header count.", + description="AFE connectivity error count.", ) diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py index 5cfc46143ac4..d2960d5bc89a 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py @@ -45,8 +45,7 @@ def __init__(self): self.record_attempt_start = MagicMock() self.record_attempt_completion = MagicMock() self.set_method = MagicMock() - self.record_gfe_metrics = MagicMock() - self.record_afe_metrics = MagicMock() + self.record_front_end_metrics = MagicMock() self.set_project = MagicMock() self.set_instance = MagicMock() self.set_database = MagicMock() @@ -118,6 +117,5 @@ def test_intercept_with_tracer(interceptor, mock_tracer_ctx): assert response == invoked_response mock_tracer_ctx.record_attempt_start.assert_called() mock_tracer_ctx.record_attempt_completion.assert_called_once() - mock_tracer_ctx.record_gfe_metrics.assert_called_once() - mock_tracer_ctx.record_afe_metrics.assert_called_once() + mock_tracer_ctx.record_front_end_metrics.assert_called_once() mock_invoked_method.assert_called_once_with("request", call_details) diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py index 7e79fd124f5c..89c13a3b77bd 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py @@ -29,6 +29,10 @@ def metrics_tracer(): mock_attempt_latency = mock.create_autospec(Histogram, instance=True) mock_operation_counter = mock.create_autospec(Counter, instance=True) mock_operation_latency = mock.create_autospec(Histogram, instance=True) + mock_gfe_latency = mock.create_autospec(Histogram, instance=True) + mock_gfe_missing = mock.create_autospec(Counter, instance=True) + mock_afe_latency = mock.create_autospec(Histogram, instance=True) + mock_afe_missing = mock.create_autospec(Counter, instance=True) return MetricsTracer( enabled=True, instrument_attempt_latency=mock_attempt_latency, @@ -36,6 +40,10 @@ def metrics_tracer(): instrument_operation_latency=mock_operation_latency, instrument_operation_counter=mock_operation_counter, client_attributes={"project_id": "test_project"}, + instrument_gfe_latency=mock_gfe_latency, + instrument_gfe_connectivity_error_count=mock_gfe_missing, + instrument_afe_latency=mock_afe_latency, + instrument_afe_connectivity_error_count=mock_afe_missing, ) @@ -245,58 +253,77 @@ def test_record_gfe_latency(metrics_tracer): metrics_tracer.enabled = True # Reset for next test -def test_record_gfe_missing_header_count(metrics_tracer): - mock_gfe_missing_header_count = mock.create_autospec(Counter, instance=True) - metrics_tracer._instrument_gfe_missing_header_count = mock_gfe_missing_header_count +def test_record_gfe_connectivity_error_count(metrics_tracer): + mock_gfe_connectivity_error_count = mock.create_autospec(Counter, instance=True) + metrics_tracer._instrument_gfe_connectivity_error_count = ( + mock_gfe_connectivity_error_count + ) metrics_tracer.gfe_enabled = True # Ensure GFE is enabled # Test when tracing is enabled - metrics_tracer.record_gfe_missing_header_count() - assert mock_gfe_missing_header_count.add.call_count == 1 - assert mock_gfe_missing_header_count.add.call_args[1]["amount"] == 1 + metrics_tracer.record_gfe_connectivity_error_count() + assert mock_gfe_connectivity_error_count.add.call_count == 1 + assert mock_gfe_connectivity_error_count.add.call_args[1]["amount"] == 1 assert ( - mock_gfe_missing_header_count.add.call_args[1]["attributes"] + mock_gfe_connectivity_error_count.add.call_args[1]["attributes"] == metrics_tracer.client_attributes ) # Test when tracing is disabled metrics_tracer.enabled = False - metrics_tracer.record_gfe_missing_header_count() - assert mock_gfe_missing_header_count.add.call_count == 1 # Should not increment + metrics_tracer.record_gfe_connectivity_error_count() + assert mock_gfe_connectivity_error_count.add.call_count == 1 # Should not increment metrics_tracer.enabled = True # Reset for next test -def test_extract_gfe_latency(): +def test_extract_front_end_latencies(): # Valid trailing metadata list of tuples - metadata_list = [("server-timing", "gfet4t7; dur=123")] - assert MetricsTracer.extract_gfe_latency(metadata_list) == 123 + metadata_list = [ + ("server-timing", "gfet4t7; dur=123"), + ("server-timing", "afe; dur=100"), + ] + assert MetricsTracer.extract_front_end_latencies(metadata_list) == (123, 100) # Valid metadata dict metadata_dict = {"server-timing": "gfet4t7; dur=456"} - assert MetricsTracer.extract_gfe_latency(metadata_dict) == 456 + assert MetricsTracer.extract_front_end_latencies(metadata_dict) == (456, None) # Missing header - assert MetricsTracer.extract_gfe_latency([("other-header", "val")]) is None - assert MetricsTracer.extract_gfe_latency(None) is None + assert MetricsTracer.extract_front_end_latencies([("other-header", "val")]) == ( + None, + None, + ) + assert MetricsTracer.extract_front_end_latencies(None) == (None, None) -def test_record_gfe_metrics(metrics_tracer): +def test_record_front_end_metrics(metrics_tracer): mock_gfe_latency = mock.create_autospec(Histogram, instance=True) mock_gfe_missing = mock.create_autospec(Counter, instance=True) + mock_afe_latency = mock.create_autospec(Histogram, instance=True) + mock_afe_missing = mock.create_autospec(Counter, instance=True) metrics_tracer._instrument_gfe_latency = mock_gfe_latency - metrics_tracer._instrument_gfe_missing_header_count = mock_gfe_missing + metrics_tracer._instrument_gfe_connectivity_error_count = mock_gfe_missing + metrics_tracer._instrument_afe_latency = mock_afe_latency + metrics_tracer._instrument_afe_connectivity_error_count = mock_afe_missing metrics_tracer.gfe_enabled = True # With header - metrics_tracer.record_gfe_metrics([("server-timing", "gfet4t7; dur=88")]) + metrics_tracer.record_front_end_metrics( + [("server-timing", "gfet4t7; dur=88"), ("server-timing", "afe; dur=90")] + ) assert mock_gfe_latency.record.call_count == 1 assert mock_gfe_latency.record.call_args[1]["amount"] == 88 assert mock_gfe_missing.add.call_count == 0 + assert mock_afe_latency.record.call_count == 1 + assert mock_afe_latency.record.call_args[1]["amount"] == 90 + assert mock_afe_missing.add.call_count == 0 # Without header - metrics_tracer.record_gfe_metrics([("other", "1")]) + metrics_tracer.record_front_end_metrics([("other", "1")]) assert mock_gfe_latency.record.call_count == 1 assert mock_gfe_missing.add.call_count == 1 + assert mock_afe_latency.record.call_count == 1 + assert mock_afe_missing.add.call_count == 1 def test_record_afe_latency(metrics_tracer): @@ -318,12 +345,12 @@ def test_record_afe_latency(metrics_tracer): metrics_tracer.enabled = True -def test_record_afe_missing_header_count(metrics_tracer): +def test_record_afe_connectivity_error_count(metrics_tracer): mock_afe_missing = mock.create_autospec(Counter, instance=True) - metrics_tracer._instrument_afe_missing_header_count = mock_afe_missing + metrics_tracer._instrument_afe_connectivity_error_count = mock_afe_missing metrics_tracer.gfe_enabled = True - metrics_tracer.record_afe_missing_header_count() + metrics_tracer.record_afe_connectivity_error_count() assert mock_afe_missing.add.call_count == 1 assert mock_afe_missing.add.call_args[1]["amount"] == 1 assert ( @@ -332,39 +359,6 @@ def test_record_afe_missing_header_count(metrics_tracer): ) metrics_tracer.enabled = False - metrics_tracer.record_afe_missing_header_count() + metrics_tracer.record_afe_connectivity_error_count() assert mock_afe_missing.add.call_count == 1 metrics_tracer.enabled = True - - -def test_extract_afe_latency(): - # Valid trailing metadata list of tuples - metadata_list = [("server-timing", "afe; dur=123")] - assert MetricsTracer.extract_afe_latency(metadata_list) == 123 - - # Valid metadata dict - metadata_dict = {"server-timing": "afet4t7; dur=456"} - assert MetricsTracer.extract_afe_latency(metadata_dict) == 456 - - # Missing header - assert MetricsTracer.extract_afe_latency([("other-header", "val")]) is None - assert MetricsTracer.extract_afe_latency(None) is None - - -def test_record_afe_metrics(metrics_tracer): - mock_afe_latency = mock.create_autospec(Histogram, instance=True) - mock_afe_missing = mock.create_autospec(Counter, instance=True) - metrics_tracer._instrument_afe_latency = mock_afe_latency - metrics_tracer._instrument_afe_missing_header_count = mock_afe_missing - metrics_tracer.gfe_enabled = True - - # With header - metrics_tracer.record_afe_metrics([("server-timing", "afe; dur=88")]) - assert mock_afe_latency.record.call_count == 1 - assert mock_afe_latency.record.call_args[1]["amount"] == 88 - assert mock_afe_missing.add.call_count == 0 - - # Without header - metrics_tracer.record_afe_metrics([("other", "1")]) - assert mock_afe_latency.record.call_count == 1 - assert mock_afe_missing.add.call_count == 1 From d938f30b992c35f9c34335c51a7eab0a7f74ddb3 Mon Sep 17 00:00:00 2001 From: Subham Sinha <35077434+sinhasubham@users.noreply.github.com> Date: Thu, 23 Jul 2026 12:34:03 +0530 Subject: [PATCH 4/6] feat(metrics): add support for AFE server timings --- .../spanner_v1/metrics/metrics_interceptor.py | 14 +++++++++++++- .../cloud/spanner_v1/metrics/metrics_tracer.py | 11 +++++++---- .../mockserver_tests/test_frontend_metrics.py | 12 ++++++++---- .../tests/unit/test_metrics_interceptor.py | 3 ++- .../tests/unit/test_metrics_tracer.py | 16 ++++++++++++---- 5 files changed, 42 insertions(+), 14 deletions(-) diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py index 24ff9aa11ec4..1a4e1863ea4b 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py @@ -16,6 +16,7 @@ import inspect import logging +import os import re from typing import Any, Dict @@ -133,6 +134,12 @@ def intercept(self, invoked_method, request_or_iterator, call_details): tracer.set_method(method_name) tracer.record_attempt_start() + + if os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true": + metadata = list(call_details.metadata or []) + metadata.append(("x-goog-spanner-enable-afe-server-timing", "true")) + call_details = call_details._replace(metadata=metadata) + response = invoked_method(request_or_iterator, call_details) return _wrap_response(response, tracer) @@ -215,8 +222,13 @@ async def _async_intercept( tracer.set_method(method_name) tracer.record_attempt_start() - response = await continuation(call_details, request_or_iterator) + if os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true": + metadata = list(call_details.metadata or []) + metadata.append(("x-goog-spanner-enable-afe-server-timing", "true")) + call_details = call_details._replace(metadata=metadata) + + response = await continuation(call_details, request_or_iterator) if hasattr(response, "__anext__"): return _AsyncStreamingResponseWrapper(response, tracer) else: diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py index a357682357ec..f24a09c594c0 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py @@ -19,6 +19,7 @@ while the helper classes provide additional functionality and context for the metrics being traced. """ +import os import re from datetime import datetime from typing import Any, Dict, Optional @@ -425,7 +426,7 @@ def record_gfe_latency(self, latency: int) -> None: ): return self._instrument_gfe_latency.record( - amount=latency, attributes=self.client_attributes + amount=latency, attributes=self._create_attempt_otel_attributes() ) def record_gfe_connectivity_error_count(self) -> None: @@ -439,7 +440,7 @@ def record_gfe_connectivity_error_count(self) -> None: ): return self._instrument_gfe_connectivity_error_count.add( - amount=1, attributes=self.client_attributes + amount=1, attributes=self._create_attempt_otel_attributes() ) def record_afe_latency(self, latency: int) -> None: @@ -453,10 +454,11 @@ def record_afe_latency(self, latency: int) -> None: not self.enabled or not HAS_OPENTELEMETRY_INSTALLED or not getattr(self, "_instrument_afe_latency", None) + or os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() == "true" ): return self._instrument_afe_latency.record( - amount=latency, attributes=self.client_attributes + amount=latency, attributes=self._create_attempt_otel_attributes() ) def record_afe_connectivity_error_count(self) -> None: @@ -467,10 +469,11 @@ def record_afe_connectivity_error_count(self) -> None: not self.enabled or not HAS_OPENTELEMETRY_INSTALLED or not getattr(self, "_instrument_afe_connectivity_error_count", None) + or os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() == "true" ): return self._instrument_afe_connectivity_error_count.add( - amount=1, attributes=self.client_attributes + amount=1, attributes=self._create_attempt_otel_attributes() ) @staticmethod diff --git a/packages/google-cloud-spanner/tests/mockserver_tests/test_frontend_metrics.py b/packages/google-cloud-spanner/tests/mockserver_tests/test_frontend_metrics.py index 24be3159924c..367104260c4d 100644 --- a/packages/google-cloud-spanner/tests/mockserver_tests/test_frontend_metrics.py +++ b/packages/google-cloud-spanner/tests/mockserver_tests/test_frontend_metrics.py @@ -202,16 +202,20 @@ def test_gfe_missing_header_count_exported(self): } self.assertIn( - "gfe_missing_header_count", metrics, f"Metrics: {list(metrics.keys())}" + "gfe_connectivity_error_count", + metrics, + f"Metrics: {list(metrics.keys())}", ) - missing_metric = metrics["gfe_missing_header_count"] + missing_metric = metrics["gfe_connectivity_error_count"] point = next(iter(missing_metric.data.data_points)) self.assertGreaterEqual(point.value, 1) self.assertIn( - "afe_missing_header_count", metrics, f"Metrics: {list(metrics.keys())}" + "afe_connectivity_error_count", + metrics, + f"Metrics: {list(metrics.keys())}", ) - afe_missing_metric = metrics["afe_missing_header_count"] + afe_missing_metric = metrics["afe_connectivity_error_count"] afe_point = next(iter(afe_missing_metric.data.data_points)) self.assertGreaterEqual(afe_point.value, 1) finally: diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py index d2960d5bc89a..2d4c1bbcbe20 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py @@ -113,9 +113,10 @@ def test_intercept_with_tracer(interceptor, mock_tracer_ctx): ], ) + replaced_call_details = call_details._replace.return_value response = interceptor.intercept(mock_invoked_method, "request", call_details) assert response == invoked_response mock_tracer_ctx.record_attempt_start.assert_called() mock_tracer_ctx.record_attempt_completion.assert_called_once() mock_tracer_ctx.record_front_end_metrics.assert_called_once() - mock_invoked_method.assert_called_once_with("request", call_details) + mock_invoked_method.assert_called_once_with("request", replaced_call_details) diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py index 89c13a3b77bd..a35284703d87 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py @@ -243,7 +243,7 @@ def test_record_gfe_latency(metrics_tracer): assert mock_gfe_latency.record.call_args[1]["amount"] == 100 assert ( mock_gfe_latency.record.call_args[1]["attributes"] - == metrics_tracer.client_attributes + == metrics_tracer._create_attempt_otel_attributes() ) # Test when tracing is disabled @@ -266,7 +266,7 @@ def test_record_gfe_connectivity_error_count(metrics_tracer): assert mock_gfe_connectivity_error_count.add.call_args[1]["amount"] == 1 assert ( mock_gfe_connectivity_error_count.add.call_args[1]["attributes"] - == metrics_tracer.client_attributes + == metrics_tracer._create_attempt_otel_attributes() ) # Test when tracing is disabled @@ -336,9 +336,13 @@ def test_record_afe_latency(metrics_tracer): assert mock_afe_latency.record.call_args[1]["amount"] == 100 assert ( mock_afe_latency.record.call_args[1]["attributes"] - == metrics_tracer.client_attributes + == metrics_tracer._create_attempt_otel_attributes() ) + with mock.patch.dict("os.environ", {"SPANNER_DISABLE_AFE_SERVER_TIMING": "true"}): + metrics_tracer.record_afe_latency(300) + assert mock_afe_latency.record.call_count == 1 + metrics_tracer.enabled = False metrics_tracer.record_afe_latency(200) assert mock_afe_latency.record.call_count == 1 @@ -355,9 +359,13 @@ def test_record_afe_connectivity_error_count(metrics_tracer): assert mock_afe_missing.add.call_args[1]["amount"] == 1 assert ( mock_afe_missing.add.call_args[1]["attributes"] - == metrics_tracer.client_attributes + == metrics_tracer._create_attempt_otel_attributes() ) + with mock.patch.dict("os.environ", {"SPANNER_DISABLE_AFE_SERVER_TIMING": "true"}): + metrics_tracer.record_afe_connectivity_error_count() + assert mock_afe_missing.add.call_count == 1 + metrics_tracer.enabled = False metrics_tracer.record_afe_connectivity_error_count() assert mock_afe_missing.add.call_count == 1 From df804281d2b490e83e5d1245d5297db05d5f421d Mon Sep 17 00:00:00 2001 From: Subham Sinha <35077434+sinhasubham@users.noreply.github.com> Date: Thu, 23 Jul 2026 14:54:34 +0530 Subject: [PATCH 5/6] fix(metrics): fix OpenTelemetry resource detector import path --- packages/google-cloud-spanner/tests/_helpers.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/packages/google-cloud-spanner/tests/_helpers.py b/packages/google-cloud-spanner/tests/_helpers.py index 83aecfd8b2f6..519a79f5a447 100644 --- a/packages/google-cloud-spanner/tests/_helpers.py +++ b/packages/google-cloud-spanner/tests/_helpers.py @@ -13,6 +13,9 @@ try: from opentelemetry import trace + from opentelemetry.resourcedetector.gcp_resource_detector import ( # noqa: F401 + GoogleCloudResourceDetector, + ) from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import SimpleSpanProcessor from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( From a1532cea1234fee98a5d4f22506c98dfe691f476 Mon Sep 17 00:00:00 2001 From: Subham Sinha <35077434+sinhasubham@users.noreply.github.com> Date: Tue, 28 Jul 2026 15:00:58 +0530 Subject: [PATCH 6/6] fix(metrics): remove obsolete gfe_enabled flag and refine AFE timing logic --- .../google/cloud/spanner_v1/_helpers.py | 16 +++++++++++++++- .../spanner_v1/metrics/metrics_interceptor.py | 15 --------------- .../cloud/spanner_v1/metrics/metrics_tracer.py | 15 +++++++-------- .../spanner_v1/metrics/metrics_tracer_factory.py | 3 --- .../metrics/spanner_metrics_tracer_factory.py | 1 - .../tests/system/test_metrics.py | 1 + .../tests/unit/test_metrics_interceptor.py | 4 +--- .../tests/unit/test_metrics_tracer.py | 13 ++++--------- 8 files changed, 28 insertions(+), 40 deletions(-) diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/_helpers.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/_helpers.py index 6026e2db81f6..06b137db1e28 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/_helpers.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/_helpers.py @@ -20,6 +20,7 @@ import logging import math import operator +import os import threading import time import uuid @@ -69,6 +70,11 @@ import random from typing import List, Tuple +ENABLE_AFE_SERVER_TIMING = ( + os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true" + and os.environ.get("SPANNER_DISABLE_BUILTIN_METRICS", "").lower() != "true" +) + # Validation error messages NUMERIC_MAX_SCALE_ERR_MSG = ( "Max scale for a numeric is 9. The requested numeric has scale {}" @@ -707,6 +713,13 @@ def __init__(self, session): self._session = session +def _append_routing_headers(metadata): + """Appends routing and backend-specific headers to the metadata.""" + if ENABLE_AFE_SERVER_TIMING: + metadata.append(("x-goog-spanner-enable-afe-server-timing", "true")) + return metadata + + def _metadata_with_prefix(prefix, **kw): """Create RPC metadata containing a prefix. @@ -716,7 +729,8 @@ def _metadata_with_prefix(prefix, **kw): Returns: List[Tuple[str, str]]: RPC metadata with supplied prefix """ - return [("google-cloud-resource-prefix", prefix)] + metadata = [("google-cloud-resource-prefix", prefix)] + return _append_routing_headers(metadata) def _retry_on_aborted_exception( diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py index 1a4e1863ea4b..fa3808d034cb 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py @@ -16,7 +16,6 @@ import inspect import logging -import os import re from typing import Any, Dict @@ -128,18 +127,11 @@ def intercept(self, invoked_method, request_or_iterator, call_details): ## Format method to be be spanner. method_str = call_details.method - if isinstance(method_str, bytes): - method_str = method_str.decode("utf-8") method_name = method_str.removeprefix(SPANNER_METHOD_PREFIX).replace("/", ".") tracer.set_method(method_name) tracer.record_attempt_start() - if os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true": - metadata = list(call_details.metadata or []) - metadata.append(("x-goog-spanner-enable-afe-server-timing", "true")) - call_details = call_details._replace(metadata=metadata) - response = invoked_method(request_or_iterator, call_details) return _wrap_response(response, tracer) @@ -216,18 +208,11 @@ async def _async_intercept( MetricsInterceptor._set_metrics_tracer_attributes(resources) method_str = call_details.method - if isinstance(method_str, bytes): - method_str = method_str.decode("utf-8") method_name = method_str.removeprefix(SPANNER_METHOD_PREFIX).replace("/", ".") tracer.set_method(method_name) tracer.record_attempt_start() - if os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true": - metadata = list(call_details.metadata or []) - metadata.append(("x-goog-spanner-enable-afe-server-timing", "true")) - call_details = call_details._replace(metadata=metadata) - response = await continuation(call_details, request_or_iterator) if hasattr(response, "__anext__"): return _AsyncStreamingResponseWrapper(response, tracer) diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py index f24a09c594c0..6106fa6e18b0 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py @@ -190,7 +190,6 @@ class should not have any knowledge about the observability framework used for m _instrument_afe_connectivity_error_count: "Counter" current_op: MetricOpTracer enabled: bool - gfe_enabled: bool method: str def __init__( @@ -205,7 +204,6 @@ def __init__( instrument_gfe_connectivity_error_count: "Counter", instrument_afe_latency: "Histogram", instrument_afe_connectivity_error_count: "Counter", - gfe_enabled: bool = False, ): """ Initialize a MetricsTracer instance with the given parameters. @@ -221,7 +219,6 @@ def __init__( instrument_operation_latency (Histogram): Instrument for measuring operation latency. instrument_operation_counter (Counter): Instrument for counting operations. client_attributes (Dict[str, str]): Dictionary of client attributes used for metrics tracing. - gfe_enabled (bool, optional): Indicates if GFE metrics are enabled. Defaults to False. instrument_gfe_latency (Histogram): Instrument for measuring GFE latency. instrument_gfe_connectivity_error_count (Counter): Instrument for counting GFE connectivity errors. instrument_afe_latency (Histogram): Instrument for measuring AFE latency. @@ -242,7 +239,9 @@ def __init__( instrument_afe_connectivity_error_count ) self.enabled = enabled - self.gfe_enabled = True + self.afe_server_timing_enabled = ( + os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true" + ) @staticmethod def _get_ms_time_diff(start: datetime, end: datetime) -> float: @@ -454,7 +453,7 @@ def record_afe_latency(self, latency: int) -> None: not self.enabled or not HAS_OPENTELEMETRY_INSTALLED or not getattr(self, "_instrument_afe_latency", None) - or os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() == "true" + or not getattr(self, "afe_server_timing_enabled", True) ): return self._instrument_afe_latency.record( @@ -469,7 +468,7 @@ def record_afe_connectivity_error_count(self) -> None: not self.enabled or not HAS_OPENTELEMETRY_INSTALLED or not getattr(self, "_instrument_afe_connectivity_error_count", None) - or os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() == "true" + or not getattr(self, "afe_server_timing_enabled", True) ): return self._instrument_afe_connectivity_error_count.add( @@ -500,7 +499,7 @@ def extract_front_end_latencies( header_vals = [] for key, val in items: key_str = key.decode("utf-8") if isinstance(key, bytes) else str(key) - if key_str and key_str.lower() in ("server-timing", "server_timing"): + if key_str and key_str.lower() == "server-timing": if isinstance(val, (list, tuple)): header_vals.extend(val) else: @@ -529,7 +528,7 @@ def extract_front_end_latencies( pass if afe_latency is None: - match = re.search(r"afe(?:t4t7)?;\s*dur=([0-9.]+)", header_val) + match = re.search(r"afe;\s*dur=([0-9.]+)", header_val) if match: try: afe_latency = int(float(match.group(1))) diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py index 804fb46f1faf..0363c08898c0 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py @@ -52,7 +52,6 @@ class MetricsTracerFactory: """Factory class for creating MetricTracer instances. This class facilitates the creation of MetricTracer objects, which are responsible for collecting and tracing metrics.""" enabled: bool - gfe_enabled: bool _instrument_attempt_latency: "Histogram" _instrument_attempt_counter: "Counter" _instrument_operation_latency: "Histogram" @@ -89,7 +88,6 @@ def __init__(self, enabled: bool, service_name: str): project (str): The project ID for the monitored resource. """ self.enabled = enabled - self.gfe_enabled = True self._create_metric_instruments(service_name) self._client_attributes = {} @@ -273,7 +271,6 @@ def create_metrics_tracer(self) -> MetricsTracer: instrument_operation_latency=self._instrument_operation_latency, instrument_operation_counter=self._instrument_operation_counter, client_attributes=self._client_attributes.copy(), - gfe_enabled=True, instrument_gfe_latency=self._instrument_gfe_latency, instrument_gfe_connectivity_error_count=self._instrument_gfe_connectivity_error_count, instrument_afe_latency=self._instrument_afe_latency, diff --git a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/spanner_metrics_tracer_factory.py b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/spanner_metrics_tracer_factory.py index 7886e555f120..5342d6474689 100644 --- a/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/spanner_metrics_tracer_factory.py +++ b/packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/spanner_metrics_tracer_factory.py @@ -80,7 +80,6 @@ def __new__(cls, enabled: bool = True) -> "SpannerMetricsTracerFactory": cls._generate_client_hash(client_uid) ) cls._metrics_tracer_factory.set_location(_get_cloud_region()) - cls._metrics_tracer_factory.gfe_enabled = True if cls._metrics_tracer_factory.enabled != enabled: cls._metrics_tracer_factory.enabled = enabled diff --git a/packages/google-cloud-spanner/tests/system/test_metrics.py b/packages/google-cloud-spanner/tests/system/test_metrics.py index 09ee81ee0113..7877f9a1d70e 100644 --- a/packages/google-cloud-spanner/tests/system/test_metrics.py +++ b/packages/google-cloud-spanner/tests/system/test_metrics.py @@ -79,6 +79,7 @@ def test_builtin_metrics_with_default_otel(metrics_database): "spanner/operation_count", "spanner/attempt_count", "spanner/gfe_latencies", + "spanner/afe_latencies", } assert expected_metrics.issubset(collected_metrics) diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py index 2d4c1bbcbe20..aa31dc5f1210 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py @@ -41,7 +41,6 @@ def __init__(self): self.project = None self.instance = None self.database = None - self.gfe_enabled = True self.record_attempt_start = MagicMock() self.record_attempt_completion = MagicMock() self.set_method = MagicMock() @@ -113,10 +112,9 @@ def test_intercept_with_tracer(interceptor, mock_tracer_ctx): ], ) - replaced_call_details = call_details._replace.return_value response = interceptor.intercept(mock_invoked_method, "request", call_details) assert response == invoked_response mock_tracer_ctx.record_attempt_start.assert_called() mock_tracer_ctx.record_attempt_completion.assert_called_once() mock_tracer_ctx.record_front_end_metrics.assert_called_once() - mock_invoked_method.assert_called_once_with("request", replaced_call_details) + mock_invoked_method.assert_called_once_with("request", call_details) diff --git a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py index a35284703d87..ce5d97ac2ec7 100644 --- a/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py +++ b/packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py @@ -235,7 +235,6 @@ def test_set_method(metrics_tracer): def test_record_gfe_latency(metrics_tracer): mock_gfe_latency = mock.create_autospec(Histogram, instance=True) metrics_tracer._instrument_gfe_latency = mock_gfe_latency - metrics_tracer.gfe_enabled = True # Ensure GFE is enabled # Test when tracing is enabled metrics_tracer.record_gfe_latency(100) @@ -258,7 +257,6 @@ def test_record_gfe_connectivity_error_count(metrics_tracer): metrics_tracer._instrument_gfe_connectivity_error_count = ( mock_gfe_connectivity_error_count ) - metrics_tracer.gfe_enabled = True # Ensure GFE is enabled # Test when tracing is enabled metrics_tracer.record_gfe_connectivity_error_count() @@ -305,7 +303,6 @@ def test_record_front_end_metrics(metrics_tracer): metrics_tracer._instrument_gfe_connectivity_error_count = mock_gfe_missing metrics_tracer._instrument_afe_latency = mock_afe_latency metrics_tracer._instrument_afe_connectivity_error_count = mock_afe_missing - metrics_tracer.gfe_enabled = True # With header metrics_tracer.record_front_end_metrics( @@ -329,7 +326,6 @@ def test_record_front_end_metrics(metrics_tracer): def test_record_afe_latency(metrics_tracer): mock_afe_latency = mock.create_autospec(Histogram, instance=True) metrics_tracer._instrument_afe_latency = mock_afe_latency - metrics_tracer.gfe_enabled = True metrics_tracer.record_afe_latency(100) assert mock_afe_latency.record.call_count == 1 @@ -339,8 +335,8 @@ def test_record_afe_latency(metrics_tracer): == metrics_tracer._create_attempt_otel_attributes() ) - with mock.patch.dict("os.environ", {"SPANNER_DISABLE_AFE_SERVER_TIMING": "true"}): - metrics_tracer.record_afe_latency(300) + metrics_tracer.afe_server_timing_enabled = False + metrics_tracer.record_afe_latency(300) assert mock_afe_latency.record.call_count == 1 metrics_tracer.enabled = False @@ -352,7 +348,6 @@ def test_record_afe_latency(metrics_tracer): def test_record_afe_connectivity_error_count(metrics_tracer): mock_afe_missing = mock.create_autospec(Counter, instance=True) metrics_tracer._instrument_afe_connectivity_error_count = mock_afe_missing - metrics_tracer.gfe_enabled = True metrics_tracer.record_afe_connectivity_error_count() assert mock_afe_missing.add.call_count == 1 @@ -362,8 +357,8 @@ def test_record_afe_connectivity_error_count(metrics_tracer): == metrics_tracer._create_attempt_otel_attributes() ) - with mock.patch.dict("os.environ", {"SPANNER_DISABLE_AFE_SERVER_TIMING": "true"}): - metrics_tracer.record_afe_connectivity_error_count() + metrics_tracer.afe_server_timing_enabled = False + metrics_tracer.record_afe_connectivity_error_count() assert mock_afe_missing.add.call_count == 1 metrics_tracer.enabled = False