Skip to content
This repository was archived by the owner on Apr 1, 2026. It is now read-only.

Commit 486068b

Browse files
committed
added try; generated sync
1 parent d5e012d commit 486068b

6 files changed

Lines changed: 68 additions & 36 deletions

File tree

google/cloud/bigtable/data/_async/metrics_interceptor.py

Lines changed: 13 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -44,18 +44,19 @@ def _with_operation_from_metadata(func):
4444
@wraps(func)
4545
def wrapper(self, continuation, client_call_details, request):
4646
found_operation_id: str | None = None
47-
new_metadata = client_call_details.metadata
48-
if client_call_details.metadata:
49-
# find operation key from metadata
50-
temp_metadata = []
51-
for k, v in client_call_details.metadata:
52-
if k == OPERATION_INTERCEPTOR_METADATA_KEY:
53-
found_operation_id = v
54-
else:
55-
temp_metadata.append((k, v))
56-
new_metadata = temp_metadata
57-
# update client_call_details to drop the operation key metadata
58-
client_call_details.metadata = new_metadata
47+
try:
48+
new_metadata: list[tuple[str, str]] = []
49+
if client_call_details.metadata:
50+
# find operation key from metadata
51+
for k, v in client_call_details.metadata:
52+
if k == OPERATION_INTERCEPTOR_METADATA_KEY:
53+
found_operation_id = v
54+
else:
55+
new_metadata.append((k, v))
56+
# update client_call_details to drop the operation key metadata
57+
client_call_details.metadata = new_metadata
58+
except Exception:
59+
pass
5960

6061
operation: "ActiveOperationMetric" = self.operation_map.get(found_operation_id)
6162
if operation:

google/cloud/bigtable/data/_metrics/data_model.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -166,7 +166,7 @@ class ActiveOperationMetric:
166166
@property
167167
def interceptor_metadata(self) -> tuple[str, str]:
168168
"""
169-
returns a tuple to attach to the grpc metadata.
169+
returns a tuple to attach to the grpc metadata.
170170
171171
This metadata field will be read by the BigtableMetricsInterceptor to associate a request with an operation
172172
"""

google/cloud/bigtable/data/_sync_autogen/metrics_interceptor.py

Lines changed: 13 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -33,15 +33,19 @@ def _with_operation_from_metadata(func):
3333

3434
@wraps(func)
3535
def wrapper(self, continuation, client_call_details, request):
36-
key = next(
37-
(
38-
m[1]
39-
for m in client_call_details.metadata
40-
if m[0] == OPERATION_INTERCEPTOR_METADATA_KEY
41-
),
42-
None,
43-
)
44-
operation: "ActiveOperationMetric" = self.operation_map.get(key)
36+
found_operation_id: str | None = None
37+
try:
38+
new_metadata: list[tuple[str, str]] = []
39+
if client_call_details.metadata:
40+
for k, v in client_call_details.metadata:
41+
if k == OPERATION_INTERCEPTOR_METADATA_KEY:
42+
found_operation_id = v
43+
else:
44+
new_metadata.append((k, v))
45+
client_call_details.metadata = new_metadata
46+
except Exception:
47+
pass
48+
operation: "ActiveOperationMetric" = self.operation_map.get(found_operation_id)
4549
if operation:
4650
if (
4751
operation.state == OperationState.CREATED

tests/unit/data/_async/test_metrics_interceptor.py

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -145,7 +145,10 @@ async def test_strip_operation_id_metadata(self):
145145
instance.operation_map[op.uuid] = op
146146
continuation = CrossSync.Mock()
147147
details = ClientCallDetails()
148-
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid), ("other_key", "other_value")]
148+
details.metadata = [
149+
(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid),
150+
("other_key", "other_value"),
151+
]
149152
await instance.intercept_unary_unary(continuation, details, mock.Mock())
150153
assert details.metadata == [("other_key", "other_value")]
151154
assert continuation.call_count == 1

tests/unit/data/_metrics/test_data_model.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -577,6 +577,7 @@ def test_end_with_status_with_default_cluster_zone(self):
577577
called_with = h.on_operation_complete.call_args[0][0]
578578
assert called_with.cluster_id == DEFAULT_CLUSTER_ID
579579
assert called_with.zone == DEFAULT_ZONE
580+
580581
def test_end_with_success(self):
581582
"""
582583
end with success should be a pass-through helper for end_with_status

tests/unit/data/_sync_autogen/test_metrics_interceptor.py

Lines changed: 36 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717

1818
import pytest
1919
from grpc import RpcError
20+
from grpc import ClientCallDetails
2021
from google.cloud.bigtable.data._metrics.data_model import OperationState
2122
from google.cloud.bigtable.data._cross_sync import CrossSync
2223

@@ -103,11 +104,33 @@ def test_on_operation_cancelled(self):
103104
op.cancel()
104105
assert op.uuid not in instance.operation_map
105106

107+
def test_strip_operation_id_metadata(self):
108+
"""After operation id is detected in metadata, the field should be stripped out before calling continuation"""
109+
from google.cloud.bigtable.data._metrics.data_model import (
110+
OPERATION_INTERCEPTOR_METADATA_KEY,
111+
)
112+
113+
instance = self._make_one()
114+
op = mock.Mock()
115+
op.uuid = "test-uuid"
116+
op.state = OperationState.ACTIVE_ATTEMPT
117+
instance.operation_map[op.uuid] = op
118+
continuation = CrossSync._Sync_Impl.Mock()
119+
details = ClientCallDetails()
120+
details.metadata = [
121+
(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid),
122+
("other_key", "other_value"),
123+
]
124+
instance.intercept_unary_unary(continuation, details, mock.Mock())
125+
assert details.metadata == [("other_key", "other_value")]
126+
assert continuation.call_count == 1
127+
assert continuation.call_args[0][0].metadata == [("other_key", "other_value")]
128+
106129
def test_unary_unary_interceptor_op_not_found(self):
107130
"""Test that interceptor call cuntinuation if op is not found"""
108131
instance = self._make_one()
109132
continuation = CrossSync._Sync_Impl.Mock()
110-
details = mock.Mock()
133+
details = ClientCallDetails()
111134
details.metadata = []
112135
request = mock.Mock()
113136
instance.intercept_unary_unary(continuation, details, request)
@@ -128,7 +151,7 @@ def test_unary_unary_interceptor_success(self):
128151
call = continuation.return_value
129152
call.trailing_metadata = CrossSync._Sync_Impl.Mock(return_value=[("a", "b")])
130153
call.initial_metadata = CrossSync._Sync_Impl.Mock(return_value=[("c", "d")])
131-
details = mock.Mock()
154+
details = ClientCallDetails()
132155
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
133156
request = mock.Mock()
134157
result = instance.intercept_unary_unary(continuation, details, request)
@@ -152,7 +175,7 @@ def test_unary_unary_interceptor_failure(self):
152175
exc.trailing_metadata = CrossSync._Sync_Impl.Mock(return_value=[("a", "b")])
153176
exc.initial_metadata = CrossSync._Sync_Impl.Mock(return_value=[("c", "d")])
154177
continuation = CrossSync._Sync_Impl.Mock(side_effect=exc)
155-
details = mock.Mock()
178+
details = ClientCallDetails()
156179
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
157180
request = mock.Mock()
158181
with pytest.raises(RpcError) as e:
@@ -178,7 +201,7 @@ def test_unary_unary_interceptor_failure_no_metadata(self):
178201
call = continuation.return_value
179202
call.trailing_metadata = CrossSync._Sync_Impl.Mock(return_value=[("a", "b")])
180203
call.initial_metadata = CrossSync._Sync_Impl.Mock(return_value=[("c", "d")])
181-
details = mock.Mock()
204+
details = ClientCallDetails()
182205
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
183206
request = mock.Mock()
184207
with pytest.raises(RpcError) as e:
@@ -204,7 +227,7 @@ def test_unary_unary_interceptor_failure_generic(self):
204227
call = continuation.return_value
205228
call.trailing_metadata = CrossSync._Sync_Impl.Mock(return_value=[("a", "b")])
206229
call.initial_metadata = CrossSync._Sync_Impl.Mock(return_value=[("c", "d")])
207-
details = mock.Mock()
230+
details = ClientCallDetails()
208231
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
209232
request = mock.Mock()
210233
with pytest.raises(ValueError) as e:
@@ -218,7 +241,7 @@ def test_unary_stream_interceptor_op_not_found(self):
218241
"""Test that interceptor calls continuation if op is not found"""
219242
instance = self._make_one()
220243
continuation = CrossSync._Sync_Impl.Mock()
221-
details = mock.Mock()
244+
details = ClientCallDetails()
222245
details.metadata = []
223246
request = mock.Mock()
224247
instance.intercept_unary_stream(continuation, details, request)
@@ -243,7 +266,7 @@ def test_unary_stream_interceptor_success(self):
243266
call = continuation.return_value
244267
call.trailing_metadata = CrossSync._Sync_Impl.Mock(return_value=[("a", "b")])
245268
call.initial_metadata = CrossSync._Sync_Impl.Mock(return_value=[("c", "d")])
246-
details = mock.Mock()
269+
details = ClientCallDetails()
247270
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
248271
request = mock.Mock()
249272
wrapper = instance.intercept_unary_stream(continuation, details, request)
@@ -274,7 +297,7 @@ def test_unary_stream_interceptor_failure_mid_stream(self):
274297
call = continuation.return_value
275298
call.trailing_metadata = CrossSync._Sync_Impl.Mock(return_value=[("a", "b")])
276299
call.initial_metadata = CrossSync._Sync_Impl.Mock(return_value=[("c", "d")])
277-
details = mock.Mock()
300+
details = ClientCallDetails()
278301
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
279302
request = mock.Mock()
280303
wrapper = instance.intercept_unary_stream(continuation, details, request)
@@ -304,7 +327,7 @@ def test_unary_stream_interceptor_failure_start_stream(self):
304327
exc.initial_metadata = CrossSync._Sync_Impl.Mock(return_value=[("c", "d")])
305328
continuation = CrossSync._Sync_Impl.Mock()
306329
continuation.side_effect = exc
307-
details = mock.Mock()
330+
details = ClientCallDetails()
308331
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
309332
request = mock.Mock()
310333
with pytest.raises(RpcError) as e:
@@ -331,7 +354,7 @@ def test_unary_stream_interceptor_failure_start_stream_no_metadata(self):
331354
exc = RpcError("test")
332355
continuation = CrossSync._Sync_Impl.Mock()
333356
continuation.side_effect = exc
334-
details = mock.Mock()
357+
details = ClientCallDetails()
335358
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
336359
request = mock.Mock()
337360
with pytest.raises(RpcError) as e:
@@ -358,7 +381,7 @@ def test_unary_stream_interceptor_failure_start_stream_generic(self):
358381
exc = ValueError("test")
359382
continuation = CrossSync._Sync_Impl.Mock()
360383
continuation.side_effect = exc
361-
details = mock.Mock()
384+
details = ClientCallDetails()
362385
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
363386
request = mock.Mock()
364387
with pytest.raises(ValueError) as e:
@@ -387,7 +410,7 @@ def test_unary_unary_interceptor_start_operation(self, initial_state):
387410
call = continuation.return_value
388411
call.trailing_metadata = CrossSync._Sync_Impl.Mock(return_value=[])
389412
call.initial_metadata = CrossSync._Sync_Impl.Mock(return_value=[])
390-
details = mock.Mock()
413+
details = ClientCallDetails()
391414
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
392415
request = mock.Mock()
393416
instance.intercept_unary_unary(continuation, details, request)
@@ -413,7 +436,7 @@ def test_unary_stream_interceptor_start_operation(self, initial_state):
413436
call = continuation.return_value
414437
call.trailing_metadata = CrossSync._Sync_Impl.Mock(return_value=[])
415438
call.initial_metadata = CrossSync._Sync_Impl.Mock(return_value=[])
416-
details = mock.Mock()
439+
details = ClientCallDetails()
417440
details.metadata = [(OPERATION_INTERCEPTOR_METADATA_KEY, op.uuid)]
418441
request = mock.Mock()
419442
instance.intercept_unary_stream(continuation, details, request)

0 commit comments

Comments
 (0)