1# Copyright 2020 The gRPC Authors
2#
3# Licensed under the Apache License, Version 2.0 (the "License");
4# you may not use this file except in compliance with the License.
5# You may obtain a copy of the License at
6#
7# http://www.apache.org/licenses/LICENSE-2.0
8#
9# Unless required by applicable law or agreed to in writing, software
10# distributed under the License is distributed on an "AS IS" BASIS,
11# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12# See the License for the specific language governing permissions and
13# limitations under the License.
14"""Abstract base classes for Channel objects and Multicallable objects."""
15
16import abc
17from types import TracebackType
18from typing import Generic, Optional
19
20import grpc
21from typing_extensions import Self
22
23from . import _base_call # pyright: ignore[reportPrivateUsage]
24from ._typing import DeserializingFunction
25from ._typing import MetadataType
26from ._typing import RequestIterableType
27from ._typing import RequestType
28from ._typing import ResponseType
29from ._typing import SerializingFunction
30
31
32class UnaryUnaryMultiCallable(Generic[RequestType, ResponseType], abc.ABC):
33 """Enables asynchronous invocation of a unary-call RPC."""
34
35 @abc.abstractmethod
36 def __call__(
37 self,
38 request: RequestType,
39 *,
40 timeout: Optional[float] = None,
41 metadata: Optional[MetadataType] = None,
42 credentials: Optional[grpc.CallCredentials] = None,
43 wait_for_ready: Optional[bool] = None,
44 compression: Optional[grpc.Compression] = None,
45 ) -> _base_call.UnaryUnaryCall[RequestType, ResponseType]:
46 """Asynchronously invokes the underlying RPC.
47
48 Args:
49 request: The request value for the RPC.
50 timeout: An optional duration of time in seconds to allow
51 for the RPC.
52 metadata: Optional :term:`metadata` to be transmitted to the
53 service-side of the RPC.
54 credentials: An optional CallCredentials for the RPC. Only valid for
55 secure Channel.
56 wait_for_ready: An optional flag to enable :term:`wait_for_ready` mechanism.
57 compression: An element of grpc.Compression, e.g.
58 grpc.Compression.Gzip.
59
60 Returns:
61 A UnaryUnaryCall object.
62
63 Raises:
64 RpcError: Indicates that the RPC terminated with non-OK status. The
65 raised RpcError will also be a Call for the RPC affording the RPC's
66 metadata, status code, and details.
67 """
68
69
70class UnaryStreamMultiCallable(Generic[RequestType, ResponseType], abc.ABC):
71 """Enables asynchronous invocation of a server-streaming RPC."""
72
73 @abc.abstractmethod
74 def __call__(
75 self,
76 request: RequestType,
77 *,
78 timeout: Optional[float] = None,
79 metadata: Optional[MetadataType] = None,
80 credentials: Optional[grpc.CallCredentials] = None,
81 wait_for_ready: Optional[bool] = None,
82 compression: Optional[grpc.Compression] = None,
83 ) -> _base_call.UnaryStreamCall[RequestType, ResponseType]:
84 """Asynchronously invokes the underlying RPC.
85
86 Args:
87 request: The request value for the RPC.
88 timeout: An optional duration of time in seconds to allow
89 for the RPC.
90 metadata: Optional :term:`metadata` to be transmitted to the
91 service-side of the RPC.
92 credentials: An optional CallCredentials for the RPC. Only valid for
93 secure Channel.
94 wait_for_ready: An optional flag to enable :term:`wait_for_ready` mechanism.
95 compression: An element of grpc.Compression, e.g.
96 grpc.Compression.Gzip.
97
98 Returns:
99 A UnaryStreamCall object.
100
101 Raises:
102 RpcError: Indicates that the RPC terminated with non-OK status. The
103 raised RpcError will also be a Call for the RPC affording the RPC's
104 metadata, status code, and details.
105 """
106
107
108class StreamUnaryMultiCallable(Generic[RequestType, ResponseType], abc.ABC):
109 """Enables asynchronous invocation of a client-streaming RPC."""
110
111 @abc.abstractmethod
112 def __call__(
113 self,
114 request_iterator: Optional[RequestIterableType[RequestType]] = None,
115 timeout: Optional[float] = None,
116 metadata: Optional[MetadataType] = None,
117 credentials: Optional[grpc.CallCredentials] = None,
118 wait_for_ready: Optional[bool] = None,
119 compression: Optional[grpc.Compression] = None,
120 ) -> _base_call.StreamUnaryCall[RequestType, ResponseType]:
121 """Asynchronously invokes the underlying RPC.
122
123 Args:
124 request_iterator: An optional async iterable or iterable of request
125 messages for the RPC.
126 timeout: An optional duration of time in seconds to allow
127 for the RPC.
128 metadata: Optional :term:`metadata` to be transmitted to the
129 service-side of the RPC.
130 credentials: An optional CallCredentials for the RPC. Only valid for
131 secure Channel.
132 wait_for_ready: An optional flag to enable :term:`wait_for_ready` mechanism.
133 compression: An element of grpc.Compression, e.g.
134 grpc.Compression.Gzip.
135
136 Returns:
137 A StreamUnaryCall object.
138
139 Raises:
140 RpcError: Indicates that the RPC terminated with non-OK status. The
141 raised RpcError will also be a Call for the RPC affording the RPC's
142 metadata, status code, and details.
143 """
144
145
146class StreamStreamMultiCallable(Generic[RequestType, ResponseType], abc.ABC):
147 """Enables asynchronous invocation of a bidirectional-streaming RPC."""
148
149 @abc.abstractmethod
150 def __call__(
151 self,
152 request_iterator: Optional[RequestIterableType[RequestType]] = None,
153 timeout: Optional[float] = None,
154 metadata: Optional[MetadataType] = None,
155 credentials: Optional[grpc.CallCredentials] = None,
156 wait_for_ready: Optional[bool] = None,
157 compression: Optional[grpc.Compression] = None,
158 ) -> _base_call.StreamStreamCall[RequestType, ResponseType]:
159 """Asynchronously invokes the underlying RPC.
160
161 Args:
162 request_iterator: An optional async iterable or iterable of request
163 messages for the RPC.
164 timeout: An optional duration of time in seconds to allow
165 for the RPC.
166 metadata: Optional :term:`metadata` to be transmitted to the
167 service-side of the RPC.
168 credentials: An optional CallCredentials for the RPC. Only valid for
169 secure Channel.
170 wait_for_ready: An optional flag to enable :term:`wait_for_ready` mechanism.
171 compression: An element of grpc.Compression, e.g.
172 grpc.Compression.Gzip.
173
174 Returns:
175 A StreamStreamCall object.
176
177 Raises:
178 RpcError: Indicates that the RPC terminated with non-OK status. The
179 raised RpcError will also be a Call for the RPC affording the RPC's
180 metadata, status code, and details.
181 """
182
183
184class Channel(abc.ABC):
185 """Enables asynchronous RPC invocation as a client.
186
187 Channel objects implement the Asynchronous Context Manager (aka. async
188 with) type, although they are not supported to be entered and exited
189 multiple times.
190 """
191
192 @abc.abstractmethod
193 async def __aenter__(self) -> Self:
194 """Starts an asynchronous context manager.
195
196 Returns:
197 Channel the channel that was instantiated.
198 """
199
200 @abc.abstractmethod
201 async def __aexit__(
202 self,
203 exc_type: Optional[type[BaseException]],
204 exc_val: Optional[BaseException],
205 exc_tb: Optional[TracebackType],
206 ) -> Optional[bool]:
207 """Finishes the asynchronous context manager by closing the channel.
208
209 Still active RPCs will be cancelled.
210 """
211
212 @abc.abstractmethod
213 async def close(self, grace: Optional[float] = None):
214 """Closes this Channel and releases all resources held by it.
215
216 This method immediately stops the channel from executing new RPCs in
217 all cases.
218
219 If a grace period is specified, this method waits until all active
220 RPCs are finished or until the grace period is reached. RPCs that haven't
221 been terminated within the grace period are aborted.
222 If a grace period is not specified (by passing None for grace),
223 all existing RPCs are cancelled immediately.
224
225 This method is idempotent.
226 """
227
228 @abc.abstractmethod
229 def get_state(
230 self, try_to_connect: bool = False
231 ) -> grpc.ChannelConnectivity:
232 """Checks the connectivity state of a channel.
233
234 This is an EXPERIMENTAL API.
235
236 If the channel reaches a stable connectivity state, it is guaranteed
237 that the return value of this function will eventually converge to that
238 state.
239
240 Args:
241 try_to_connect: a bool indicate whether the Channel should try to
242 connect to peer or not.
243
244 Returns: A ChannelConnectivity object.
245 """
246
247 @abc.abstractmethod
248 async def wait_for_state_change(
249 self,
250 last_observed_state: grpc.ChannelConnectivity,
251 ) -> None:
252 """Waits for a change in connectivity state.
253
254 This is an EXPERIMENTAL API.
255
256 The function blocks until there is a change in the channel connectivity
257 state from the "last_observed_state". If the state is already
258 different, this function will return immediately.
259
260 There is an inherent race between the invocation of
261 "Channel.wait_for_state_change" and "Channel.get_state". The state can
262 change arbitrary many times during the race, so there is no way to
263 observe every state transition.
264
265 If there is a need to put a timeout for this function, please refer to
266 "asyncio.wait_for".
267
268 Args:
269 last_observed_state: A grpc.ChannelConnectivity object representing
270 the last known state.
271 """
272
273 @abc.abstractmethod
274 async def channel_ready(self) -> None:
275 """Creates a coroutine that blocks until the Channel is READY."""
276
277 @abc.abstractmethod
278 def unary_unary(
279 self,
280 method: str,
281 request_serializer: Optional[SerializingFunction[RequestType]] = None,
282 response_deserializer: Optional[
283 DeserializingFunction[ResponseType]
284 ] = None,
285 _registered_method: Optional[bool] = False,
286 ) -> UnaryUnaryMultiCallable[RequestType, ResponseType]:
287 """Creates a UnaryUnaryMultiCallable for a unary-unary method.
288
289 Args:
290 method: The name of the RPC method.
291 request_serializer: Optional :term:`serializer` for serializing the request
292 message. Request goes unserialized in case None is passed.
293 response_deserializer: Optional :term:`deserializer` for deserializing the
294 response message. Response goes undeserialized in case None
295 is passed.
296 _registered_method: Implementation Private. Optional: A bool representing
297 whether the method is registered.
298
299 Returns:
300 A UnaryUnaryMultiCallable value for the named unary-unary method.
301 """
302
303 @abc.abstractmethod
304 def unary_stream(
305 self,
306 method: str,
307 request_serializer: Optional[SerializingFunction[RequestType]] = None,
308 response_deserializer: Optional[
309 DeserializingFunction[ResponseType]
310 ] = None,
311 _registered_method: Optional[bool] = False,
312 ) -> UnaryStreamMultiCallable[RequestType, ResponseType]:
313 """Creates a UnaryStreamMultiCallable for a unary-stream method.
314
315 Args:
316 method: The name of the RPC method.
317 request_serializer: Optional :term:`serializer` for serializing the request
318 message. Request goes unserialized in case None is passed.
319 response_deserializer: Optional :term:`deserializer` for deserializing the
320 response message. Response goes undeserialized in case None
321 is passed.
322 _registered_method: Implementation Private. Optional: A bool representing
323 whether the method is registered.
324
325 Returns:
326 A UnaryStreamMultiCallable value for the named unary-stream method.
327 """
328
329 @abc.abstractmethod
330 def stream_unary(
331 self,
332 method: str,
333 request_serializer: Optional[SerializingFunction[RequestType]] = None,
334 response_deserializer: Optional[
335 DeserializingFunction[ResponseType]
336 ] = None,
337 _registered_method: Optional[bool] = False,
338 ) -> StreamUnaryMultiCallable[RequestType, ResponseType]:
339 """Creates a StreamUnaryMultiCallable for a stream-unary method.
340
341 Args:
342 method: The name of the RPC method.
343 request_serializer: Optional :term:`serializer` for serializing the request
344 message. Request goes unserialized in case None is passed.
345 response_deserializer: Optional :term:`deserializer` for deserializing the
346 response message. Response goes undeserialized in case None
347 is passed.
348 _registered_method: Implementation Private. Optional: A bool representing
349 whether the method is registered.
350
351 Returns:
352 A StreamUnaryMultiCallable value for the named stream-unary method.
353 """
354
355 @abc.abstractmethod
356 def stream_stream(
357 self,
358 method: str,
359 request_serializer: Optional[SerializingFunction[RequestType]] = None,
360 response_deserializer: Optional[
361 DeserializingFunction[ResponseType]
362 ] = None,
363 _registered_method: Optional[bool] = False,
364 ) -> StreamStreamMultiCallable[RequestType, ResponseType]:
365 """Creates a StreamStreamMultiCallable for a stream-stream method.
366
367 Args:
368 method: The name of the RPC method.
369 request_serializer: Optional :term:`serializer` for serializing the request
370 message. Request goes unserialized in case None is passed.
371 response_deserializer: Optional :term:`deserializer` for deserializing the
372 response message. Response goes undeserialized in case None
373 is passed.
374 _registered_method: Implementation Private. Optional: A bool representing
375 whether the method is registered.
376
377 Returns:
378 A StreamStreamMultiCallable value for the named stream-stream method.
379 """