Skip to content

Commit ecc6b1f

Browse files
Merge branch 'googleapis:main' into feature/cert-rotation-streaming
2 parents e32e706 + 08f21a6 commit ecc6b1f

12 files changed

Lines changed: 536 additions & 76 deletions

File tree

packages/google-api-core/google/api_core/gapic_v1/method.py

Lines changed: 72 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020

2121
import enum
2222
import functools
23+
from typing import List, Tuple
2324

2425
from google.api_core import grpc_helpers
2526
from google.api_core.gapic_v1 import client_info
@@ -57,6 +58,52 @@ def _apply_decorators(func, decorators):
5758
return func
5859

5960

61+
def _deduplicate_metadata_tokens(*headers: str | None) -> str:
62+
"""
63+
Given one or more metadata payload strings, create a combined
64+
string with deduplicated tokens, while preserving token order.
65+
66+
Inputs are expected to contain a set of metadata tokens separated by spaces
67+
Example: `gl-python/3.14.0 grpc/1.76.0 gax/2.29.0 gapic/3.8.0 pb/6.33.4`
68+
69+
Args:
70+
*headers: one or more metadata payload strings
71+
72+
Returns:
73+
a single combined payload string
74+
"""
75+
# Split all non-empty headers into individual tokens
76+
token_list = " ".join(filter(None, headers)).split()
77+
# Deduplicate while preserving order
78+
return " ".join(dict.fromkeys(token_list))
79+
80+
81+
def _extract_metrics_header(metadata) -> Tuple[str, List[Tuple[str, str]]]:
82+
"""Extract x-google-api-client header from metadata list.
83+
84+
Args:
85+
metadata (Sequence[Tuple[str, str]]): The metadata to extract from.
86+
87+
Returns:
88+
A tuple containing:
89+
- a string representing the header value.
90+
- A sequence of remaining metadata tuples.
91+
"""
92+
if not metadata:
93+
return "", []
94+
95+
key_to_find = client_info.METRICS_METADATA_KEY
96+
97+
metric_str = _deduplicate_metadata_tokens(
98+
" ".join([v for k, v in metadata if k == key_to_find])
99+
)
100+
if not metric_str:
101+
return "", list(metadata)
102+
103+
arbitrary_metadata = [item for item in metadata if item[0] != key_to_find]
104+
return metric_str, arbitrary_metadata
105+
106+
60107
class _GapicCallable(object):
61108
"""Callable that applies retry, timeout, and metadata logic.
62109
@@ -90,7 +137,16 @@ def __init__(
90137
self._retry = retry
91138
self._timeout = timeout
92139
self._compression = compression
93-
self._metadata = metadata
140+
# Pre-extract the x-goog-api-client header from the initialized metadata.
141+
self._x_goog_api_client, remaining = _extract_metrics_header(metadata)
142+
self._static_metadata = tuple(remaining)
143+
if self._x_goog_api_client:
144+
self._default_metadata = (
145+
(client_info.METRICS_METADATA_KEY, self._x_goog_api_client),
146+
*self._static_metadata,
147+
)
148+
else:
149+
self._default_metadata = self._static_metadata
94150

95151
def __call__(
96152
self, *args, timeout=DEFAULT, retry=DEFAULT, compression=DEFAULT, **kwargs
@@ -112,16 +168,21 @@ def __call__(
112168
# Apply all applicable decorators.
113169
wrapped_func = _apply_decorators(self._target, [retry, timeout])
114170

115-
# Add the user agent metadata to the call.
116-
if self._metadata is not None:
117-
metadata = kwargs.get("metadata", [])
118-
# Due to the nature of invocation, None should be treated the same
119-
# as not specified.
120-
if metadata is None:
121-
metadata = []
122-
metadata = list(metadata)
123-
metadata.extend(self._metadata)
124-
kwargs["metadata"] = metadata
171+
if user_metadata := kwargs.get("metadata"):
172+
# Add the user agent metadata to the call.
173+
final_metadata = list(self._static_metadata)
174+
user_x_goog, remaining = _extract_metrics_header(user_metadata)
175+
176+
merged_header = _deduplicate_metadata_tokens(
177+
self._x_goog_api_client, user_x_goog
178+
)
179+
if merged_header:
180+
final_metadata.append((client_info.METRICS_METADATA_KEY, merged_header))
181+
final_metadata.extend(remaining)
182+
kwargs["metadata"] = final_metadata
183+
elif self._default_metadata:
184+
kwargs["metadata"] = self._default_metadata
185+
125186
if self._compression is not None:
126187
kwargs["compression"] = compression
127188

packages/google-api-core/google/api_core/grpc_helpers.py

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -254,11 +254,20 @@ def _create_composite_credentials(
254254
request = google.auth.transport.requests.Request()
255255

256256
# Create the metadata plugin for inserting the authorization header.
257-
metadata_plugin = google.auth.transport.grpc.AuthMetadataPlugin(
258-
credentials,
259-
request,
260-
default_host=default_host,
261-
)
257+
try:
258+
metadata_plugin = google.auth.transport.grpc.AuthMetadataPlugin(
259+
credentials,
260+
request,
261+
default_host=default_host,
262+
suppress_metrics_header=True,
263+
)
264+
except TypeError:
265+
# Support older versions of google-auth that do not accept suppress_metrics_header
266+
metadata_plugin = google.auth.transport.grpc.AuthMetadataPlugin(
267+
credentials,
268+
request,
269+
default_host=default_host,
270+
)
262271

263272
# Create a set of grpc.CallCredentials using the metadata plugin.
264273
google_auth_credentials = grpc.metadata_call_credentials(metadata_plugin)

packages/google-api-core/tests/unit/gapic/test_method.py

Lines changed: 98 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -119,6 +119,83 @@ def test_invoke_wrapped_method_with_metadata_as_none():
119119
assert len(metadata) == 1
120120

121121

122+
def test_invoke_wrapped_method_no_client_info_with_custom_metadata():
123+
method = mock.Mock(spec=["__call__"])
124+
125+
wrapped_method = google.api_core.gapic_v1.method.wrap_method(
126+
method, client_info=None
127+
)
128+
129+
wrapped_method(mock.sentinel.request, metadata=[("custom-header", "value")])
130+
131+
method.assert_called_once_with(
132+
mock.sentinel.request, metadata=[("custom-header", "value")]
133+
)
134+
135+
136+
def test_extract_metrics_header_duplicate_tokens():
137+
metadata = [
138+
("x-goog-api-client", "token1 token2"),
139+
("x-goog-api-client", "token2 token3 token1"),
140+
("other-header", "value"),
141+
("x-goog-api-client", "token4 token2"),
142+
]
143+
144+
metric_str, arbitrary_metadata = (
145+
google.api_core.gapic_v1.method._extract_metrics_header(metadata)
146+
)
147+
148+
# Should maintain order of first appearance and eliminate duplicates
149+
assert metric_str == "token1 token2 token3 token4"
150+
assert arbitrary_metadata == [("other-header", "value")]
151+
152+
153+
def test_invoke_wrapped_method_with_duplicate_x_goog_api_client_metadata():
154+
method = mock.Mock(spec=["__call__"])
155+
156+
# Create a custom ClientInfo with defined properties so we know exactly what is returned
157+
client_info = google.api_core.gapic_v1.client_info.ClientInfo(
158+
user_agent="custom-user-agent/1.0",
159+
python_version="3.14.0",
160+
grpc_version="1.76.0",
161+
api_core_version="2.29.0",
162+
)
163+
164+
wrapped_method = google.api_core.gapic_v1.method.wrap_method(
165+
method, client_info=client_info
166+
)
167+
168+
# Invoke the wrapped method with an explicit user-provided custom header that contains duplicates
169+
# both within its own items and overlapping with the default client_info
170+
wrapped_method(
171+
mock.sentinel.request,
172+
metadata=[
173+
("x-goog-api-client", "override-client/2.0"),
174+
(
175+
"x-goog-api-client",
176+
"override-client/2.0 grpc/1.76.0 custom-user-agent/1.0",
177+
),
178+
("other-header", "value"),
179+
],
180+
)
181+
182+
method.assert_called_once_with(mock.sentinel.request, metadata=mock.ANY)
183+
metadata = method.call_args[1]["metadata"]
184+
185+
# There should only be one "x-goog-api-client" header, containing both values joined by space,
186+
# plus the other-header.
187+
assert len(metadata) == 2
188+
metadata_dict = dict(metadata)
189+
assert "other-header" in metadata_dict
190+
assert metadata_dict["other-header"] == "value"
191+
assert "x-goog-api-client" in metadata_dict
192+
# Verify both the user-provided override value and the library system telemetry are merged explicitly
193+
assert (
194+
metadata_dict["x-goog-api-client"]
195+
== "custom-user-agent/1.0 gl-python/3.14.0 grpc/1.76.0 gax/2.29.0 override-client/2.0"
196+
)
197+
198+
122199
@mock.patch("time.sleep")
123200
def test_wrap_method_with_default_retry_and_timeout_and_compression(unused_sleep):
124201
method = mock.Mock(
@@ -248,3 +325,24 @@ def test_wrap_method_with_call_not_supported():
248325
with pytest.raises(ValueError) as exc_info:
249326
google.api_core.gapic_v1.method.wrap_method(method, with_call=True)
250327
assert "with_call=True is only supported for unary calls" in str(exc_info.value)
328+
329+
330+
@pytest.mark.parametrize(
331+
"headers,expected",
332+
[
333+
((), ""),
334+
(("",), ""),
335+
((None,), ""),
336+
(("", None, ""), ""),
337+
(("token1",), "token1"),
338+
(("token1 token1",), "token1"),
339+
(("token1", "token1"), "token1"),
340+
(("token1 token2 token1",), "token1 token2"),
341+
(("token1", "token2", "token1"), "token1 token2"),
342+
(("token1 token2", "token2 token3"), "token1 token2 token3"),
343+
(("token1", None, "token2", "", "token1"), "token1 token2"),
344+
],
345+
)
346+
def test__deduplicate_metadata_tokens(headers, expected):
347+
dedup = google.api_core.gapic_v1.method._deduplicate_metadata_tokens
348+
assert dedup(*headers) == expected

packages/google-auth/google/auth/transport/grpc.py

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,16 +50,21 @@ class AuthMetadataPlugin(grpc.AuthMetadataPlugin):
5050
default_host (Optional[str]): A host like "pubsub.googleapis.com".
5151
This is used when a self-signed JWT is created from service
5252
account credentials.
53+
suppress_metrics_header (bool): When enabled, ``x-goog-api-client``
54+
will be stripped from authorization headers.
5355
"""
5456

55-
def __init__(self, credentials, request, default_host=None):
57+
def __init__(
58+
self, credentials, request, default_host=None, *, suppress_metrics_header=False
59+
):
5660
# pylint: disable=no-value-for-parameter
5761
# pylint doesn't realize that the super method takes no arguments
5862
# because this class is the same name as the superclass.
5963
super(AuthMetadataPlugin, self).__init__()
6064
self._credentials = credentials
6165
self._request = request
6266
self._default_host = default_host
67+
self._suppress_metrics_header = suppress_metrics_header
6368

6469
def _get_authorization_headers(self, context):
6570
"""Gets the authorization headers for a request.
@@ -83,6 +88,9 @@ def _get_authorization_headers(self, context):
8388
self._request, context.method_name, context.service_url, headers
8489
)
8590

91+
if self._suppress_metrics_header and "x-goog-api-client" in headers:
92+
del headers["x-goog-api-client"]
93+
8694
return list(headers.items())
8795

8896
def __call__(self, context, callback):

packages/google-auth/tests/transport/test_grpc.py

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -142,6 +142,35 @@ def test__get_authorization_headers_with_service_account_and_default_host(self):
142142
"https://{}/".format(default_host)
143143
)
144144

145+
def test_suppress_metrics_header(self):
146+
credentials = mock.create_autospec(service_account.Credentials)
147+
148+
# Mock credentials before_request that adds metric and authorization
149+
def mock_before_request(request, method, url, headers):
150+
headers["x-goog-api-client"] = "foo"
151+
headers["authorization"] = "Bearer token"
152+
153+
credentials.before_request.side_effect = mock_before_request
154+
request = mock.create_autospec(transport.Request)
155+
156+
# By default, suppress_metrics_header=False
157+
plugin = google.auth.transport.grpc.AuthMetadataPlugin(credentials, request)
158+
context = mock.create_autospec(grpc.AuthMetadataContext, instance=True)
159+
context.method_name = "methodName"
160+
context.service_url = "https://pubsub.googleapis.com/methodName"
161+
162+
headers = dict(plugin._get_authorization_headers(context))
163+
assert "x-goog-api-client" in headers
164+
assert headers["x-goog-api-client"] == "foo"
165+
166+
# With suppress_metrics_header=True
167+
plugin_suppressed = google.auth.transport.grpc.AuthMetadataPlugin(
168+
credentials, request, suppress_metrics_header=True
169+
)
170+
headers_suppressed = dict(plugin_suppressed._get_authorization_headers(context))
171+
assert "x-goog-api-client" not in headers_suppressed
172+
assert headers_suppressed["authorization"] == "Bearer token"
173+
145174

146175
@mock.patch(
147176
"google.auth.transport._mtls_helper.get_client_ssl_credentials", autospec=True

packages/google-cloud-bigquery/google/cloud/bigquery/_versions_helpers.py

Lines changed: 50 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,17 +16,16 @@
1616
from typing import Any
1717

1818
import packaging.version
19-
2019
from google.cloud.bigquery import exceptions
2120

22-
2321
_MIN_PYARROW_VERSION = packaging.version.Version("3.0.0")
2422
_MIN_BQ_STORAGE_VERSION = packaging.version.Version("2.0.0")
2523
_BQ_STORAGE_OPTIONAL_READ_SESSION_VERSION = packaging.version.Version("2.6.0")
2624
_MIN_PANDAS_VERSION = packaging.version.Version("1.1.0")
2725

2826
_MIN_PANDAS_VERSION_RANGE = packaging.version.Version("1.5.0")
2927
_MIN_PYARROW_VERSION_RANGE = packaging.version.Version("10.0.1")
28+
_MIN_PANDAS_GBQ_DELEGATION_VERSION = packaging.version.Version("1.0.0")
3029

3130

3231
class PyarrowVersions:
@@ -247,3 +246,52 @@ def try_import(self, raise_if_error: bool = False) -> Any:
247246
and PYARROW_VERSIONS.try_import() is not None
248247
and PYARROW_VERSIONS.installed_version >= _MIN_PYARROW_VERSION_RANGE
249248
)
249+
250+
251+
class PandasGBQVersions:
252+
"""Version and delegation comparisons for pandas-gbq package."""
253+
254+
def __init__(self):
255+
self._installed_version = None
256+
self._delegation_api_version = None
257+
258+
@property
259+
def installed_version(self) -> packaging.version.Version:
260+
"""Return the parsed version of pandas-gbq."""
261+
if self._installed_version is not None:
262+
return self._installed_version
263+
264+
try:
265+
import pandas_gbq # type: ignore
266+
267+
self._installed_version = packaging.version.parse(
268+
getattr(pandas_gbq, "__version__", "0.0.0")
269+
)
270+
except Exception:
271+
self._installed_version = packaging.version.parse("0.0.0")
272+
return self._installed_version
273+
274+
@property
275+
def delegation_api_version(self) -> packaging.version.Version:
276+
"""Return the delegation API version of pandas-gbq if installed, otherwise 0.0.0."""
277+
if self._delegation_api_version is not None:
278+
return self._delegation_api_version
279+
280+
try:
281+
import pandas_gbq # type: ignore
282+
283+
raw_version = getattr(
284+
pandas_gbq, "_internal_delegation_api_version", "0.0.0"
285+
)
286+
self._delegation_api_version = packaging.version.parse(str(raw_version))
287+
except Exception:
288+
self._delegation_api_version = packaging.version.parse("0.0.0")
289+
return self._delegation_api_version
290+
291+
@property
292+
def is_delegation_supported(self) -> bool:
293+
"""True if the installed pandas-gbq version supports query delegation API."""
294+
return self.delegation_api_version >= _MIN_PANDAS_GBQ_DELEGATION_VERSION
295+
296+
297+
PANDAS_GBQ_VERSIONS = PandasGBQVersions()

0 commit comments

Comments
 (0)