Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/grpc/aio/_base_channel.py: 98%

Shortcuts on this page

r m x   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

47 statements  

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 """