Skip to content

Commit 3f91b5f

Browse files
committed
fix(metrics): remove obsolete gfe_enabled flag and refine AFE timing logic
1 parent df80428 commit 3f91b5f

8 files changed

Lines changed: 28 additions & 38 deletions

File tree

packages/google-cloud-spanner/google/cloud/spanner_v1/_helpers.py

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import logging
2121
import math
2222
import operator
23+
import os
2324
import threading
2425
import time
2526
import uuid
@@ -69,6 +70,11 @@
6970
import random
7071
from typing import List, Tuple
7172

73+
ENABLE_AFE_SERVER_TIMING = (
74+
os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true"
75+
and os.environ.get("SPANNER_DISABLE_BUILTIN_METRICS", "").lower() != "true"
76+
)
77+
7278
# Validation error messages
7379
NUMERIC_MAX_SCALE_ERR_MSG = (
7480
"Max scale for a numeric is 9. The requested numeric has scale {}"
@@ -707,6 +713,13 @@ def __init__(self, session):
707713
self._session = session
708714

709715

716+
def _append_routing_headers(metadata):
717+
"""Appends routing and backend-specific headers to the metadata."""
718+
if ENABLE_AFE_SERVER_TIMING:
719+
metadata.append(("x-goog-spanner-enable-afe-server-timing", "true"))
720+
return metadata
721+
722+
710723
def _metadata_with_prefix(prefix, **kw):
711724
"""Create RPC metadata containing a prefix.
712725
@@ -716,7 +729,8 @@ def _metadata_with_prefix(prefix, **kw):
716729
Returns:
717730
List[Tuple[str, str]]: RPC metadata with supplied prefix
718731
"""
719-
return [("google-cloud-resource-prefix", prefix)]
732+
metadata = [("google-cloud-resource-prefix", prefix)]
733+
return _append_routing_headers(metadata)
720734

721735

722736
def _retry_on_aborted_exception(

packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_interceptor.py

Lines changed: 0 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,6 @@
1616

1717
import inspect
1818
import logging
19-
import os
2019
import re
2120
from typing import Any, Dict
2221

@@ -128,18 +127,11 @@ def intercept(self, invoked_method, request_or_iterator, call_details):
128127

129128
## Format method to be be spanner.<method name>
130129
method_str = call_details.method
131-
if isinstance(method_str, bytes):
132-
method_str = method_str.decode("utf-8")
133130
method_name = method_str.removeprefix(SPANNER_METHOD_PREFIX).replace("/", ".")
134131

135132
tracer.set_method(method_name)
136133
tracer.record_attempt_start()
137134

138-
if os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true":
139-
metadata = list(call_details.metadata or [])
140-
metadata.append(("x-goog-spanner-enable-afe-server-timing", "true"))
141-
call_details = call_details._replace(metadata=metadata)
142-
143135
response = invoked_method(request_or_iterator, call_details)
144136

145137
return _wrap_response(response, tracer)
@@ -223,11 +215,6 @@ async def _async_intercept(
223215
tracer.set_method(method_name)
224216
tracer.record_attempt_start()
225217

226-
if os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true":
227-
metadata = list(call_details.metadata or [])
228-
metadata.append(("x-goog-spanner-enable-afe-server-timing", "true"))
229-
call_details = call_details._replace(metadata=metadata)
230-
231218
response = await continuation(call_details, request_or_iterator)
232219
if hasattr(response, "__anext__"):
233220
return _AsyncStreamingResponseWrapper(response, tracer)

packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer.py

Lines changed: 7 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -190,7 +190,6 @@ class should not have any knowledge about the observability framework used for m
190190
_instrument_afe_connectivity_error_count: "Counter"
191191
current_op: MetricOpTracer
192192
enabled: bool
193-
gfe_enabled: bool
194193
method: str
195194

196195
def __init__(
@@ -205,7 +204,6 @@ def __init__(
205204
instrument_gfe_connectivity_error_count: "Counter",
206205
instrument_afe_latency: "Histogram",
207206
instrument_afe_connectivity_error_count: "Counter",
208-
gfe_enabled: bool = False,
209207
):
210208
"""
211209
Initialize a MetricsTracer instance with the given parameters.
@@ -221,7 +219,6 @@ def __init__(
221219
instrument_operation_latency (Histogram): Instrument for measuring operation latency.
222220
instrument_operation_counter (Counter): Instrument for counting operations.
223221
client_attributes (Dict[str, str]): Dictionary of client attributes used for metrics tracing.
224-
gfe_enabled (bool, optional): Indicates if GFE metrics are enabled. Defaults to False.
225222
instrument_gfe_latency (Histogram): Instrument for measuring GFE latency.
226223
instrument_gfe_connectivity_error_count (Counter): Instrument for counting GFE connectivity errors.
227224
instrument_afe_latency (Histogram): Instrument for measuring AFE latency.
@@ -242,7 +239,9 @@ def __init__(
242239
instrument_afe_connectivity_error_count
243240
)
244241
self.enabled = enabled
245-
self.gfe_enabled = True
242+
self.afe_server_timing_enabled = (
243+
os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() != "true"
244+
)
246245

247246
@staticmethod
248247
def _get_ms_time_diff(start: datetime, end: datetime) -> float:
@@ -454,7 +453,7 @@ def record_afe_latency(self, latency: int) -> None:
454453
not self.enabled
455454
or not HAS_OPENTELEMETRY_INSTALLED
456455
or not getattr(self, "_instrument_afe_latency", None)
457-
or os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() == "true"
456+
or not getattr(self, "afe_server_timing_enabled", True)
458457
):
459458
return
460459
self._instrument_afe_latency.record(
@@ -469,7 +468,7 @@ def record_afe_connectivity_error_count(self) -> None:
469468
not self.enabled
470469
or not HAS_OPENTELEMETRY_INSTALLED
471470
or not getattr(self, "_instrument_afe_connectivity_error_count", None)
472-
or os.environ.get("SPANNER_DISABLE_AFE_SERVER_TIMING", "").lower() == "true"
471+
or not getattr(self, "afe_server_timing_enabled", True)
473472
):
474473
return
475474
self._instrument_afe_connectivity_error_count.add(
@@ -500,7 +499,7 @@ def extract_front_end_latencies(
500499
header_vals = []
501500
for key, val in items:
502501
key_str = key.decode("utf-8") if isinstance(key, bytes) else str(key)
503-
if key_str and key_str.lower() in ("server-timing", "server_timing"):
502+
if key_str and key_str.lower() == "server-timing":
504503
if isinstance(val, (list, tuple)):
505504
header_vals.extend(val)
506505
else:
@@ -529,7 +528,7 @@ def extract_front_end_latencies(
529528
pass
530529

531530
if afe_latency is None:
532-
match = re.search(r"afe(?:t4t7)?;\s*dur=([0-9.]+)", header_val)
531+
match = re.search(r"afe;\s*dur=([0-9.]+)", header_val)
533532
if match:
534533
try:
535534
afe_latency = int(float(match.group(1)))

packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/metrics_tracer_factory.py

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,6 @@ class MetricsTracerFactory:
5252
"""Factory class for creating MetricTracer instances. This class facilitates the creation of MetricTracer objects, which are responsible for collecting and tracing metrics."""
5353

5454
enabled: bool
55-
gfe_enabled: bool
5655
_instrument_attempt_latency: "Histogram"
5756
_instrument_attempt_counter: "Counter"
5857
_instrument_operation_latency: "Histogram"
@@ -89,7 +88,6 @@ def __init__(self, enabled: bool, service_name: str):
8988
project (str): The project ID for the monitored resource.
9089
"""
9190
self.enabled = enabled
92-
self.gfe_enabled = True
9391
self._create_metric_instruments(service_name)
9492
self._client_attributes = {}
9593

@@ -273,7 +271,6 @@ def create_metrics_tracer(self) -> MetricsTracer:
273271
instrument_operation_latency=self._instrument_operation_latency,
274272
instrument_operation_counter=self._instrument_operation_counter,
275273
client_attributes=self._client_attributes.copy(),
276-
gfe_enabled=True,
277274
instrument_gfe_latency=self._instrument_gfe_latency,
278275
instrument_gfe_connectivity_error_count=self._instrument_gfe_connectivity_error_count,
279276
instrument_afe_latency=self._instrument_afe_latency,

packages/google-cloud-spanner/google/cloud/spanner_v1/metrics/spanner_metrics_tracer_factory.py

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -80,7 +80,6 @@ def __new__(cls, enabled: bool = True) -> "SpannerMetricsTracerFactory":
8080
cls._generate_client_hash(client_uid)
8181
)
8282
cls._metrics_tracer_factory.set_location(_get_cloud_region())
83-
cls._metrics_tracer_factory.gfe_enabled = True
8483

8584
if cls._metrics_tracer_factory.enabled != enabled:
8685
cls._metrics_tracer_factory.enabled = enabled

packages/google-cloud-spanner/tests/system/test_metrics.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -79,6 +79,7 @@ def test_builtin_metrics_with_default_otel(metrics_database):
7979
"spanner/operation_count",
8080
"spanner/attempt_count",
8181
"spanner/gfe_latencies",
82+
"spanner/afe_latencies",
8283
}
8384
assert expected_metrics.issubset(collected_metrics)
8485

packages/google-cloud-spanner/tests/unit/test_metrics_interceptor.py

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -41,7 +41,6 @@ def __init__(self):
4141
self.project = None
4242
self.instance = None
4343
self.database = None
44-
self.gfe_enabled = True
4544
self.record_attempt_start = MagicMock()
4645
self.record_attempt_completion = MagicMock()
4746
self.set_method = MagicMock()
@@ -113,10 +112,9 @@ def test_intercept_with_tracer(interceptor, mock_tracer_ctx):
113112
],
114113
)
115114

116-
replaced_call_details = call_details._replace.return_value
117115
response = interceptor.intercept(mock_invoked_method, "request", call_details)
118116
assert response == invoked_response
119117
mock_tracer_ctx.record_attempt_start.assert_called()
120118
mock_tracer_ctx.record_attempt_completion.assert_called_once()
121119
mock_tracer_ctx.record_front_end_metrics.assert_called_once()
122-
mock_invoked_method.assert_called_once_with("request", replaced_call_details)
120+
mock_invoked_method.assert_called_once_with("request", call_details)

packages/google-cloud-spanner/tests/unit/test_metrics_tracer.py

Lines changed: 4 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -235,7 +235,6 @@ def test_set_method(metrics_tracer):
235235
def test_record_gfe_latency(metrics_tracer):
236236
mock_gfe_latency = mock.create_autospec(Histogram, instance=True)
237237
metrics_tracer._instrument_gfe_latency = mock_gfe_latency
238-
metrics_tracer.gfe_enabled = True # Ensure GFE is enabled
239238

240239
# Test when tracing is enabled
241240
metrics_tracer.record_gfe_latency(100)
@@ -258,7 +257,6 @@ def test_record_gfe_connectivity_error_count(metrics_tracer):
258257
metrics_tracer._instrument_gfe_connectivity_error_count = (
259258
mock_gfe_connectivity_error_count
260259
)
261-
metrics_tracer.gfe_enabled = True # Ensure GFE is enabled
262260

263261
# Test when tracing is enabled
264262
metrics_tracer.record_gfe_connectivity_error_count()
@@ -305,7 +303,6 @@ def test_record_front_end_metrics(metrics_tracer):
305303
metrics_tracer._instrument_gfe_connectivity_error_count = mock_gfe_missing
306304
metrics_tracer._instrument_afe_latency = mock_afe_latency
307305
metrics_tracer._instrument_afe_connectivity_error_count = mock_afe_missing
308-
metrics_tracer.gfe_enabled = True
309306

310307
# With header
311308
metrics_tracer.record_front_end_metrics(
@@ -329,7 +326,6 @@ def test_record_front_end_metrics(metrics_tracer):
329326
def test_record_afe_latency(metrics_tracer):
330327
mock_afe_latency = mock.create_autospec(Histogram, instance=True)
331328
metrics_tracer._instrument_afe_latency = mock_afe_latency
332-
metrics_tracer.gfe_enabled = True
333329

334330
metrics_tracer.record_afe_latency(100)
335331
assert mock_afe_latency.record.call_count == 1
@@ -339,8 +335,8 @@ def test_record_afe_latency(metrics_tracer):
339335
== metrics_tracer._create_attempt_otel_attributes()
340336
)
341337

342-
with mock.patch.dict("os.environ", {"SPANNER_DISABLE_AFE_SERVER_TIMING": "true"}):
343-
metrics_tracer.record_afe_latency(300)
338+
metrics_tracer.afe_server_timing_enabled = False
339+
metrics_tracer.record_afe_latency(300)
344340
assert mock_afe_latency.record.call_count == 1
345341

346342
metrics_tracer.enabled = False
@@ -352,7 +348,6 @@ def test_record_afe_latency(metrics_tracer):
352348
def test_record_afe_connectivity_error_count(metrics_tracer):
353349
mock_afe_missing = mock.create_autospec(Counter, instance=True)
354350
metrics_tracer._instrument_afe_connectivity_error_count = mock_afe_missing
355-
metrics_tracer.gfe_enabled = True
356351

357352
metrics_tracer.record_afe_connectivity_error_count()
358353
assert mock_afe_missing.add.call_count == 1
@@ -362,8 +357,8 @@ def test_record_afe_connectivity_error_count(metrics_tracer):
362357
== metrics_tracer._create_attempt_otel_attributes()
363358
)
364359

365-
with mock.patch.dict("os.environ", {"SPANNER_DISABLE_AFE_SERVER_TIMING": "true"}):
366-
metrics_tracer.record_afe_connectivity_error_count()
360+
metrics_tracer.afe_server_timing_enabled = False
361+
metrics_tracer.record_afe_connectivity_error_count()
367362
assert mock_afe_missing.add.call_count == 1
368363

369364
metrics_tracer.enabled = False

0 commit comments

Comments
 (0)