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

Commit 7d90a04

Browse files
committed
remove extra wrapper; added invalidate_stubs helper
1 parent 4e13783 commit 7d90a04

4 files changed

Lines changed: 34 additions & 153 deletions

File tree

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

Lines changed: 11 additions & 84 deletions
Original file line numberDiff line numberDiff line change
@@ -21,69 +21,12 @@
2121
from grpc import ChannelConnectivity
2222

2323
if CrossSync.is_async:
24-
from grpc.aio import Call
2524
from grpc.aio import Channel
26-
from grpc.aio import UnaryUnaryMultiCallable
27-
from grpc.aio import UnaryStreamMultiCallable
28-
from grpc.aio import StreamUnaryMultiCallable
29-
from grpc.aio import StreamStreamMultiCallable
3025
else:
31-
from grpc import Call
3226
from grpc import Channel
33-
from grpc import UnaryUnaryMultiCallable
34-
from grpc import UnaryStreamMultiCallable
35-
from grpc import StreamUnaryMultiCallable
36-
from grpc import StreamStreamMultiCallable
3727

3828
__CROSS_SYNC_OUTPUT__ = "google.cloud.bigtable.data._sync_autogen._replaceable_channel"
3929

40-
@CrossSync.convert_class
41-
class _WrappedMultiCallable:
42-
"""
43-
Wrapper class that implements the grpc MultiCallable interface.
44-
Allows generic functions that return calls to pass checks for
45-
MultiCallable objects.
46-
"""
47-
48-
def __init__(self, call_factory: Callable[..., Call]):
49-
self._call_factory = call_factory
50-
51-
def __call__(self, *args, **kwargs) -> Call:
52-
return self._call_factory(*args, **kwargs)
53-
54-
55-
class WrappedUnaryUnaryMultiCallable(
56-
_WrappedMultiCallable, UnaryUnaryMultiCallable
57-
):
58-
if not CrossSync.is_async:
59-
# add missing functions for sync unary callable
60-
61-
def with_call(self, *args, **kwargs):
62-
call = self.__call__(self, *args, **kwargs)
63-
return call(), call
64-
65-
def future(self, *args, **kwargs):
66-
raise NotImplementedError
67-
68-
69-
class WrappedUnaryStreamMultiCallable(
70-
_WrappedMultiCallable, UnaryStreamMultiCallable
71-
):
72-
pass
73-
74-
75-
class WrappedStreamUnaryMultiCallable(
76-
_WrappedMultiCallable, StreamUnaryMultiCallable
77-
):
78-
pass
79-
80-
81-
class WrappedStreamStreamMultiCallable(
82-
_WrappedMultiCallable, StreamStreamMultiCallable
83-
):
84-
pass
85-
86-
8730
@CrossSync.convert_class(sync_name="_WrappedChannel", rm_aio=True)
8831
class _AsyncWrappedChannel(Channel):
8932
"""
@@ -94,33 +37,17 @@ class _AsyncWrappedChannel(Channel):
9437
def __init__(self, channel: Channel):
9538
self._channel = channel
9639

97-
def unary_unary(self, *args, **kwargs) -> UnaryUnaryMultiCallable:
98-
return WrappedUnaryUnaryMultiCallable(
99-
lambda *call_args, **call_kwargs: self._channel.unary_unary(
100-
*args, **kwargs
101-
)(*call_args, **call_kwargs)
102-
)
103-
104-
def unary_stream(self, *args, **kwargs) -> UnaryStreamMultiCallable:
105-
return WrappedUnaryStreamMultiCallable(
106-
lambda *call_args, **call_kwargs: self._channel.unary_stream(
107-
*args, **kwargs
108-
)(*call_args, **call_kwargs)
109-
)
110-
111-
def stream_unary(self, *args, **kwargs) -> StreamUnaryMultiCallable:
112-
return WrappedStreamUnaryMultiCallable(
113-
lambda *call_args, **call_kwargs: self._channel.stream_unary(
114-
*args, **kwargs
115-
)(*call_args, **call_kwargs)
116-
)
117-
118-
def stream_stream(self, *args, **kwargs) -> StreamStreamMultiCallable:
119-
return WrappedStreamStreamMultiCallable(
120-
lambda *call_args, **call_kwargs: self._channel.stream_stream(
121-
*args, **kwargs
122-
)(*call_args, **call_kwargs)
123-
)
40+
def unary_unary(self, *args, **kwargs):
41+
return self._channel.unary_unary(*args, **kwargs)
42+
43+
def unary_stream(self, *args, **kwargs):
44+
return self._channel.unary_stream(*args, **kwargs)
45+
46+
def stream_unary(self, *args, **kwargs):
47+
return self._channel.stream_unary(*args, **kwargs)
48+
49+
def stream_stream(self, *args, **kwargs):
50+
return self._channel.stream_stream(*args, **kwargs)
12451

12552
# grace not supported by sync version
12653
@CrossSync.drop

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -346,6 +346,11 @@ async def _ping_and_warm_instances(
346346
)
347347
return [r or None for r in result_list]
348348

349+
def _invalidate_channel_stubs(self):
350+
"""Helper to reset the cached stubs. Needed when changing out the grpc channel"""
351+
self.transport._stubs = {}
352+
self.transport._prep_wrapped_messages(self.client_info)
353+
349354
@CrossSync.convert(replace_symbols={"_AsyncReplaceableChannel": "_ReplaceableChannel"})
350355
async def _manage_channel(
351356
self,
@@ -398,6 +403,7 @@ async def _manage_channel(
398403
await self._ping_and_warm_instances(channel=new_channel)
399404
# cycle channel out of use, with long grace window before closure
400405
old_channel = super_channel.replace_wrapped_channel(new_channel)
406+
self._invalidate_channel_stubs()
401407
# give old_channel a chance to complete existing rpcs
402408
if CrossSync.is_async:
403409
await old_channel.close(grace_period)

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

Lines changed: 11 additions & 69 deletions
Original file line numberDiff line numberDiff line change
@@ -18,49 +18,7 @@
1818
from __future__ import annotations
1919
from typing import Callable
2020
from grpc import ChannelConnectivity
21-
from grpc import Call
2221
from grpc import Channel
23-
from grpc import UnaryUnaryMultiCallable
24-
from grpc import UnaryStreamMultiCallable
25-
from grpc import StreamUnaryMultiCallable
26-
from grpc import StreamStreamMultiCallable
27-
28-
29-
class _WrappedMultiCallable:
30-
"""
31-
Wrapper class that implements the grpc MultiCallable interface.
32-
Allows generic functions that return calls to pass checks for
33-
MultiCallable objects.
34-
"""
35-
36-
def __init__(self, call_factory: Callable[..., Call]):
37-
self._call_factory = call_factory
38-
39-
def __call__(self, *args, **kwargs) -> Call:
40-
return self._call_factory(*args, **kwargs)
41-
42-
43-
class WrappedUnaryUnaryMultiCallable(_WrappedMultiCallable, UnaryUnaryMultiCallable):
44-
def with_call(self, *args, **kwargs):
45-
call = self.__call__(self, *args, **kwargs)
46-
return (call(), call)
47-
48-
def future(self, *args, **kwargs):
49-
raise NotImplementedError
50-
51-
52-
class WrappedUnaryStreamMultiCallable(_WrappedMultiCallable, UnaryStreamMultiCallable):
53-
pass
54-
55-
56-
class WrappedStreamUnaryMultiCallable(_WrappedMultiCallable, StreamUnaryMultiCallable):
57-
pass
58-
59-
60-
class WrappedStreamStreamMultiCallable(
61-
_WrappedMultiCallable, StreamStreamMultiCallable
62-
):
63-
pass
6422

6523

6624
class _WrappedChannel(Channel):
@@ -72,33 +30,17 @@ class _WrappedChannel(Channel):
7230
def __init__(self, channel: Channel):
7331
self._channel = channel
7432

75-
def unary_unary(self, *args, **kwargs) -> UnaryUnaryMultiCallable:
76-
return WrappedUnaryUnaryMultiCallable(
77-
lambda *call_args, **call_kwargs: self._channel.unary_unary(
78-
*args, **kwargs
79-
)(*call_args, **call_kwargs)
80-
)
81-
82-
def unary_stream(self, *args, **kwargs) -> UnaryStreamMultiCallable:
83-
return WrappedUnaryStreamMultiCallable(
84-
lambda *call_args, **call_kwargs: self._channel.unary_stream(
85-
*args, **kwargs
86-
)(*call_args, **call_kwargs)
87-
)
88-
89-
def stream_unary(self, *args, **kwargs) -> StreamUnaryMultiCallable:
90-
return WrappedStreamUnaryMultiCallable(
91-
lambda *call_args, **call_kwargs: self._channel.stream_unary(
92-
*args, **kwargs
93-
)(*call_args, **call_kwargs)
94-
)
95-
96-
def stream_stream(self, *args, **kwargs) -> StreamStreamMultiCallable:
97-
return WrappedStreamStreamMultiCallable(
98-
lambda *call_args, **call_kwargs: self._channel.stream_stream(
99-
*args, **kwargs
100-
)(*call_args, **call_kwargs)
101-
)
33+
def unary_unary(self, *args, **kwargs):
34+
return self._channel.unary_unary(*args, **kwargs)
35+
36+
def unary_stream(self, *args, **kwargs):
37+
return self._channel.unary_stream(*args, **kwargs)
38+
39+
def stream_unary(self, *args, **kwargs):
40+
return self._channel.stream_unary(*args, **kwargs)
41+
42+
def stream_stream(self, *args, **kwargs):
43+
return self._channel.stream_stream(*args, **kwargs)
10244

10345
def channel_ready(self):
10446
return self._channel.channel_ready()

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -262,6 +262,11 @@ def _ping_and_warm_instances(
262262
)
263263
return [r or None for r in result_list]
264264

265+
def _invalidate_channel_stubs(self):
266+
"""Helper to reset the cached stubs. Needed when changing out the grpc channel"""
267+
self.transport._stubs = {}
268+
self.transport._prep_wrapped_messages(self.client_info)
269+
265270
def _manage_channel(
266271
self,
267272
refresh_interval_min: float = 60 * 35,
@@ -304,6 +309,7 @@ def _manage_channel(
304309
new_channel = super_channel.create_channel()
305310
self._ping_and_warm_instances(channel=new_channel)
306311
old_channel = super_channel.replace_wrapped_channel(new_channel)
312+
self._invalidate_channel_stubs()
307313
if grace_period:
308314
self._is_closed.wait(grace_period)
309315
old_channel.close()

0 commit comments

Comments
 (0)