Repository navigation
feat(observability): suppress redundant gRPC auto-instrumentation spans - #18592
chalmerlowe wants to merge 15 commits into
Conversation
…error recording, and templates - Hoist sync channel interceptors (apply_channel_interceptors) and async channel interceptors (apply_async_channel_interceptors) to _compat with active fallback wrapping, eliminating the noop lambda fallback in grpc.py.j2 - Hoist wrap_method tracing and kind support constants to _compat and eliminate inspect.signature module-load checks in base transports - Adopt hybrid error recording in _TraceContext.record_error to enrich span attributes (error.type, status.message) while allowing upstream OpenTelemetry defaults to handle exception events - Pass parameterized proto route template from http_options into trace_http_request url_template - Add rest_transport.kind assertion for AsyncResumableUploadServiceRestTransport to maintain 100% test coverage - Add unit test isolation for custom tracer provider and test cases for _compat interceptors
When OpenTelemetry gRPC auto-instrumentation (GrpcInstrumentorClient().instrument()) is enabled globally alongside Google Cloud SDK tracing, downstream channel calls emit redundant, un-enriched gRPC wire spans. Introduce targeted downstream suppression interceptors for both sync (grpc) and async (grpc.aio) channels using opentelemetry.instrumentation.utils.suppress_instrumentation. This ensures our feature-rich Google Cloud SDK T4 span is emitted while redundant generic auto-instrumentation spans are suppressed.
There was a problem hiding this comment.
Code Review
This pull request introduces client interceptors (_SuppressingClientInterceptor and _AsyncSuppressingClientInterceptor) to suppress redundant downstream generic auto-instrumentation spans when OpenTelemetry is enabled, ensuring that only the feature-rich Google Cloud SDK spans are emitted. It also updates the interceptor pipeline logic to prepend OpenTelemetry interceptors in order. The review feedback suggests critical performance optimizations, specifically moving imports out of the hot path in _suppress_instrumentation to avoid overhead on every RPC call, and removing a redundant inline import of grpc in otel_interceptor.
…typo Address review feedback from Gemini Code Assist and Code Review Council: - Move contextlib and _suppress_instrumentation to module level to eliminate per-RPC dynamic import overhead and ImportError catching in the hot path. - Remove redundant inline import of grpc inside otel_interceptor, safely catching (NameError, AttributeError). - Fix syntax error typo (try:x -> try:) in test_tracing.py.
Replace tuple unpacking (*_GRPC_INTERCEPTOR_BASE) in client interceptor class definitions with explicit base classes and distinct fallback dummy classes when grpc is not installed. Resolves mypy 'Invalid base class [misc]' errors.
…ate test fixtures - Flatten interceptor dispatch in apply_channel_interceptors using upfront target list - Bind grpc.intercept_channel directly to channel and set grpc=None fallback - Add mock_otel_grpc fixture in conftest.py and reuse across test_observability.py - Parametrize sync and async suppression invocation tests across all RPC forms - Parametrize negative interceptor getter tests for disabled and missing OTel
…on tests - Add # type: ignore[misc] with explanatory comment on fallback trace_http_request in _compat.py.j2 - Regenerate Bazel integration goldens with updated _compat.py typing pragma - Gate test_auto_instrumentation_suppression_sync and async with hasattr(_observability, '_AsyncSuppressingClientInterceptor') to prevent failure against released PyPI core
…el fallback - Mark defensive except (AttributeError, TypeError) block in get_otel_interceptor with # pragma: NO COVER - Add unit test verifying graceful fallback when grpc.intercept_channel raises AttributeError or TypeError
| assert child_span.attributes.get("server.address") == "localhost" | ||
| assert child_span.attributes.get("server.port") == 7469 | ||
| assert child_span.attributes.get("url.domain") == "googleapis.com" | ||
| assert child_span.attributes.get("rpc.response.status_code") == "OK" |
There was a problem hiding this comment.
A good deal of the above is duplicative with test_auto_instrumentation_suppression_async()
It also happens to be duplicative with a number of the other existing tests.
I am putting together a PR with two standardized assertion functions to help reduce the duplication across the broader array of tests, but did not want to clutter up this PR.
| record_http_error = record_error | ||
|
|
||
| def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: | ||
| # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. |
There was a problem hiding this comment.
nit: can we change this to a TODO comment to remove the # type: ignore[misc] once libraries require google-api-core >= 2.42.0?
| # Track insertion index per interceptor list for OpenTelemetry interceptors. | ||
| # OpenTelemetry interceptors and downstream suppressors are inserted at the | ||
| # beginning of the channel's interceptor pipeline in order, so that our | ||
| # Google SDK spans wrap outside and downstream generic auto-instrumentation | ||
| # (e.g. GrpcAioInstrumentorClient) is suppressed. Other interceptors are appended. | ||
| otel_insert_indices: dict[str, int] = {} |
There was a problem hiding this comment.
nit: Gemini found a rare edge case where interceptors could be corrupted if the user provides a channel with an interceptor:
Re-applying interceptors to an existing channel wipes the insertion index tracking (otel_insert_indices starting back at 0), which pushes existing OpenTelemetry interceptors downstream and breaks execution order.
Gemini suggested this fix
| # Track insertion index per interceptor list for OpenTelemetry interceptors. | |
| # OpenTelemetry interceptors and downstream suppressors are inserted at the | |
| # beginning of the channel's interceptor pipeline in order, so that our | |
| # Google SDK spans wrap outside and downstream generic auto-instrumentation | |
| # (e.g. GrpcAioInstrumentorClient) is suppressed. Other interceptors are appended. | |
| otel_insert_indices: dict[str, int] = {} | |
| # Initialize insertion index trackers. | |
| # To support multiple sequential applications of this function on the same channel, | |
| # we inspect the existing pipeline instead of assuming a start index of 0. | |
| # This preserves the original order of previously inserted OTel interceptors. | |
| otel_insert_indices: dict[str, int] = {} | |
| for _, attr_name in mapping: | |
| target_list = getattr(channel, attr_name, []) | |
| if isinstance(target_list, list): | |
| otel_insert_indices[attr_name] = sum( | |
| 1 for i in target_list if getattr(i, "_is_otel_interceptor", False) is True | |
| ) |
Here is the test that fails without the fix
def test_apply_channel_interceptors_preserves_otel_order_on_multiple_calls():
"""Verifies that applying channel interceptors multiple times does not corrupt
the execution order of pre-existing OpenTelemetry interceptors.
This prevents pipeline regression when an end-user applies their own
OTel tracking interceptors first, and the GCP Client subsequently appends its own
internal suppressor interceptor during construction.
"""
# 1. Create mock interceptors representing the telemetry pipeline
user_otel_interceptor = mock.Mock(spec=["intercept_unary_unary"])
user_otel_interceptor._is_otel_interceptor = True
sdk_suppressor_interceptor = mock.Mock(spec=["intercept_unary_unary"])
sdk_suppressor_interceptor._is_otel_interceptor = True
user_logging_interceptor = mock.Mock(spec=["intercept_unary_unary"])
# 2. Set up the target mock channel
channel = mock.Mock()
channel._unary_unary_interceptors = []
# 3. Simulate FIRST call (User configures custom channel with their own telemetry)
grpc_helpers_async.apply_channel_interceptors(
channel, [user_logging_interceptor, user_otel_interceptor]
)
# The pipeline should have user OTel interceptor prepended and logging appended
assert channel._unary_unary_interceptors == [user_otel_interceptor, user_logging_interceptor]
# 4. Simulate SECOND call (GCP SDK client appends internal suppressor)
grpc_helpers_async.apply_channel_interceptors(
channel, [sdk_suppressor_interceptor]
)
# 5. Assertions: The parent span interceptor MUST execute BEFORE the suppressor.
# Expected order: [User OTel, Suppressor, User Logging]
# Without the fix, the suppressor is incorrectly prepended at index 0, altering order.
assert channel._unary_unary_interceptors == [
user_otel_interceptor,
sdk_suppressor_interceptor,
user_logging_interceptor,
]
| redundant generic downstream spans are bypassed. | ||
| """ | ||
|
|
||
| def __init__(self) -> None: |
There was a problem hiding this comment.
Gemini suggested we should add slots here to optimize memory, along with a regression test.
Here are the tests according to https://wiki.python.org/moin/UsingSlots : Why Not Use Slots?
- Do these interceptors need dynamic attribute creation?
No. The _SuppressingClientInterceptor and _AsyncSuppressingClientInterceptor are strictly designed to be stateless wrappers. Their only job is to intercept a call, enter a context manager, and yield to the continuation. The only attribute they store is self._is_otel_interceptor = True during initialization. Because this is hardcoded in the class logic, it can be handled by slots. They never need dynamic, on-the-fly attribute additions from external code.
- Do they need weak references (weakref)?
No. These interceptors are bound directly to the lifecycle of the gRPC channel (grpc.Channel). They are never held in global weak-reference tracking registries, nor do they require special non-blocking garbage collection patterns. They can be cleanly deallocated whenever the channel itself is garbage collected.
- Do they use descriptors or cached properties (like cached_property)?
No. These interceptors are incredibly lightweight. They do not have any properties, let alone lazy or cached properties. All of their logic is procedural and executes on every call.
| def __init__(self) -> None: | |
| __slots__ = () | |
| def __init__(self) -> None: |
Regression test
def test_suppressing_interceptors_define_slots():
"""Verifies that internal telemetry suppression interceptors explicitly define
empty __slots__ to prevent unnecessary duplicate instance dictionaries.
"""
# We verify that the classes themselves define the __slots__ property.
# This remains True regardless of whether base third-party packages introduce a __dict__.
assert hasattr(_observability._SuppressingClientInterceptor, "__slots__"), (
"Sync suppressing interceptor class is missing the __slots__ declaration."
)
assert hasattr(_observability._AsyncSuppressingClientInterceptor, "__slots__"), (
"Async suppressing interceptor class is missing the __slots__ declaration."
)
# Verify that the defined slot layout is empty to prevent dynamic allocations
assert _observability._SuppressingClientInterceptor.__slots__ == (), (
"Sync suppressing interceptor slots must be empty."
)
assert _observability._AsyncSuppressingClientInterceptor.__slots__ == (), (
"Async suppressing interceptor slots must be empty."
)
full diff for packages/google-api-core/google/api_core/_observability.py
git diff google/api_core/_observability.py
diff --git a/packages/google-api-core/google/api_core/_observability.py b/packages/google-api-core/google/api_core/_observability.py
index 89e33b25f28..268f11cf629 100644
--- a/packages/google-api-core/google/api_core/_observability.py
+++ b/packages/google-api-core/google/api_core/_observability.py
@@ -253,16 +253,16 @@ except ImportError: # pragma: NO COVER
# mypy: Fallback dummy classes when optional grpc is not installed.
# Four distinct empty classes avoid duplicate base class 'object' TypeError at runtime.
class _SyncUnaryUnaryClientInterceptor: # type: ignore[no-redef]
- pass
+ __slots__ = ()
class _SyncUnaryStreamClientInterceptor: # type: ignore[no-redef]
- pass
+ __slots__ = ()
class _SyncStreamUnaryClientInterceptor: # type: ignore[no-redef]
- pass
+ __slots__ = ()
class _SyncStreamStreamClientInterceptor: # type: ignore[no-redef]
- pass
+ __slots__ = ()
try:
@@ -282,16 +282,16 @@ except ImportError: # pragma: NO COVER
# mypy: Fallback dummy classes when optional grpc is not installed.
# Four distinct empty classes avoid duplicate base class 'object' TypeError at runtime.
class _AsyncUnaryUnaryClientInterceptor: # type: ignore[no-redef]
- pass
+ __slots__ = ()
class _AsyncUnaryStreamClientInterceptor: # type: ignore[no-redef]
- pass
+ __slots__ = ()
class _AsyncStreamUnaryClientInterceptor: # type: ignore[no-redef]
- pass
+ __slots__ = ()
class _AsyncStreamStreamClientInterceptor: # type: ignore[no-redef]
- pass
+ __slots__ = ()
class _SuppressingClientInterceptor(
@@ -309,6 +309,8 @@ class _SuppressingClientInterceptor(
redundant generic downstream spans are bypassed.
"""
+ __slots__ = ()
+
def __init__(self) -> None:
self._is_otel_interceptor = True
@@ -345,6 +347,8 @@ class _AsyncSuppressingClientInterceptor(
):
"""AsyncIO client interceptor that suppresses redundant downstream generic auto-instrumentation spans."""
+ __slots__ = ()
+
def __init__(self) -> None:
self._is_otel_interceptor = True
| ): | ||
| """Proves that otel_interceptor falls back gracefully if grpc.intercept_channel raises AttributeError or TypeError.""" | ||
| if _observability.grpc is None: | ||
| pytest.skip("grpc is not installed") |
There was a problem hiding this comment.
nit: minor clean up
| pytest.skip("grpc is not installed") | |
| @pytest.mark.skipif(_observability.grpc is None, reason="grpc is not installed") | |
| @pytest.mark.parametrize("error_cls", [AttributeError, TypeError]) | |
| def test_get_otel_interceptor_grpc_intercept_channel_fallback( | |
| monkeypatch, mock_otel_grpc, error_cls | |
| ): | |
| """Proves that otel_interceptor falls back gracefully if grpc.intercept_channel raises AttributeError or TypeError.""" |
| @pytest.mark.parametrize( | ||
| "method_name", | ||
| [ | ||
| "intercept_unary_unary", | ||
| "intercept_unary_stream", | ||
| "intercept_stream_unary", | ||
| "intercept_stream_stream", | ||
| ], | ||
| ids=["unary_unary", "unary_stream", "stream_unary", "stream_stream"], | ||
| ) | ||
| def test_async_suppressing_interceptor_invocations(method_name): | ||
| """Verifies that _AsyncSuppressingClientInterceptor wraps downstream async calls in suppress_instrumentation. | ||
|
|
||
| When global OpenTelemetry gRPC auto-instrumentation is enabled, it emits | ||
| generic, un-enriched attempt spans. We wrap downstream async calls in | ||
| suppress_instrumentation so only our enriched Google Cloud SDK T4 span is emitted. | ||
| """ | ||
| import asyncio | ||
|
|
||
| suppressor = _observability._AsyncSuppressingClientInterceptor() | ||
| assert getattr(suppressor, "_is_otel_interceptor", False) is True | ||
|
|
||
| async def _run_test(): | ||
| with mock.patch( | ||
| "google.api_core._observability._suppress_instrumentation" | ||
| ) as mock_suppress: | ||
| mock_continuation = mock.AsyncMock(return_value="expected_response") | ||
| interceptor_method = getattr(suppressor, method_name) | ||
| result = await interceptor_method( | ||
| mock_continuation, "call_details", "request_or_iterator" | ||
| ) | ||
|
|
||
| assert result == "expected_response" | ||
| mock_continuation.assert_called_once_with( | ||
| "call_details", "request_or_iterator" | ||
| ) | ||
| mock_suppress.assert_called_once() | ||
|
|
||
| asyncio.run(_run_test()) |
There was a problem hiding this comment.
nit: to improve readability, we can use the @pytest.mark.asyncio decorator.
nit: async tests should live in tests/asyncio/...
| @pytest.mark.parametrize( | |
| "method_name", | |
| [ | |
| "intercept_unary_unary", | |
| "intercept_unary_stream", | |
| "intercept_stream_unary", | |
| "intercept_stream_stream", | |
| ], | |
| ids=["unary_unary", "unary_stream", "stream_unary", "stream_stream"], | |
| ) | |
| def test_async_suppressing_interceptor_invocations(method_name): | |
| """Verifies that _AsyncSuppressingClientInterceptor wraps downstream async calls in suppress_instrumentation. | |
| When global OpenTelemetry gRPC auto-instrumentation is enabled, it emits | |
| generic, un-enriched attempt spans. We wrap downstream async calls in | |
| suppress_instrumentation so only our enriched Google Cloud SDK T4 span is emitted. | |
| """ | |
| import asyncio | |
| suppressor = _observability._AsyncSuppressingClientInterceptor() | |
| assert getattr(suppressor, "_is_otel_interceptor", False) is True | |
| async def _run_test(): | |
| with mock.patch( | |
| "google.api_core._observability._suppress_instrumentation" | |
| ) as mock_suppress: | |
| mock_continuation = mock.AsyncMock(return_value="expected_response") | |
| interceptor_method = getattr(suppressor, method_name) | |
| result = await interceptor_method( | |
| mock_continuation, "call_details", "request_or_iterator" | |
| ) | |
| assert result == "expected_response" | |
| mock_continuation.assert_called_once_with( | |
| "call_details", "request_or_iterator" | |
| ) | |
| mock_suppress.assert_called_once() | |
| asyncio.run(_run_test()) | |
| @pytest.mark.asyncio | |
| @pytest.mark.parametrize( | |
| "method_name", | |
| [ | |
| "intercept_unary_unary", | |
| "intercept_unary_stream", | |
| "intercept_stream_unary", | |
| "intercept_stream_stream", | |
| ], | |
| ids=["unary_unary", "unary_stream", "stream_unary", "stream_stream"], | |
| ) | |
| async def test_async_suppressing_interceptor_invocations(method_name): | |
| """Verifies that _AsyncSuppressingClientInterceptor wraps downstream async calls in suppress_instrumentation. | |
| When global OpenTelemetry gRPC auto-instrumentation is enabled, it emits | |
| generic, un-enriched attempt spans. We wrap downstream async calls in | |
| suppress_instrumentation so only our enriched Google Cloud SDK T4 span is emitted. | |
| """ | |
| suppressor = _observability._AsyncSuppressingClientInterceptor() | |
| assert getattr(suppressor, "_is_otel_interceptor", False) is True | |
| with mock.patch( | |
| "google.api_core._observability._suppress_instrumentation" | |
| ) as mock_suppress: | |
| mock_continuation = mock.AsyncMock(return_value="expected_response") | |
| interceptor_method = getattr(suppressor, method_name) | |
| result = await interceptor_method( | |
| mock_continuation, "call_details", "request_or_iterator" | |
| ) | |
| assert result == "expected_response" | |
| mock_continuation.assert_called_once_with( | |
| "call_details", "request_or_iterator" | |
| ) | |
| mock_suppress.assert_called_once() |
daniel-sanche
left a comment
There was a problem hiding this comment.
I don't have time to do a full in-depth review today, but I'd really suggest taking a step back and seeing if there's another way to solve this. This would add a lot of bloat, with the three layers of wrapers fighting each other, which I'm sure would show up in benchmarks. And the workarounds added to apply_channel_interceptors to identify and re-order interceptors feel like a big abstraction overreach
I added a comment about a possible approach to strip out unwanted interceptors before applying our custom otel ones. If that doesn't work, maybe spend some time considering other options, before committing to these new wrapper layers
| channel = grpc.intercept_channel(channel, suppressor) | ||
| except (AttributeError, TypeError): # pragma: NO COVER | ||
| # Fallback if grpc lacks intercept_channel or custom channel type rejects it. | ||
| pass |
There was a problem hiding this comment.
My unerstanding of this fix is that channel may already be wrapped with an interceptor, and we need to add another wrapper to supress it
Instead, can we unwrap the channel here?
def otel_interceptor(channel: grpc.Channel) -> grpc.Channel:
# Peel off any generic upstream auto-instrumentation wrapper once at setup time.
while (
isinstance(getattr(channel, "_interceptor", None), OpenTelemetryClientInterceptor)
and not getattr(channel._interceptor, "_is_otel_interceptor", False)
and hasattr(channel, "_channel")
):
channel = channel._channel
return otel_grpc.intercept_channel(channel, interceptor)
(We could do something similar on the async side, pruning existing the channel's interceptor list before adding our own otel one)
| # OpenTelemetry interceptors and downstream suppressors are inserted at the | ||
| # beginning of the channel's interceptor pipeline in order, so that our | ||
| # Google SDK spans wrap outside and downstream generic auto-instrumentation | ||
| # (e.g. GrpcAioInstrumentorClient) is suppressed. Other interceptors are appended. |
There was a problem hiding this comment.
I would really like to avoid this kind of special-case re-ordering if at all possible. Even using a different method for re-ordering, before calling apply_channel_interceptors if you have to
If you need to add it here, it should be called out in the docstring that otel is treated as a special case
| HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( | ||
| _observability is not None | ||
| and hasattr(_observability, "_AsyncSuppressingClientInterceptor") | ||
| ) |
There was a problem hiding this comment.
It seems like this is only used in a test file. Can this live there, instead of in each generated package?
Intent & Context
When upstream OpenTelemetry gRPC auto-instrumentation (
opentelemetry.instrumentation.grpc.GrpcInstrumentorClient().instrument()) is enabled globally in a customer application alongside Google Cloud SDK tracing, downstream channel invocations create redundant, un-enriched wire attempt spans.This PR introduces targeted downstream suppression interceptors for both sync (
grpc) and async (grpc.aio) transports:suppress_instrumentationcontext manager.feat/otel-tracing-universal-4path).