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

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

69 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 server-side classes.""" 

15 

16import abc 

17from typing import Generic, Iterable, Mapping, NoReturn, Optional, Sequence 

18 

19import grpc 

20 

21from ._typing import DoneCallbackType 

22from ._typing import MetadataType 

23from ._typing import RequestType 

24from ._typing import ResponseType 

25 

26 

27class Server(abc.ABC): 

28 """Serves RPCs.""" 

29 

30 @abc.abstractmethod 

31 def add_generic_rpc_handlers( 

32 self, generic_rpc_handlers: Sequence[grpc.GenericRpcHandler] 

33 ) -> None: 

34 """Registers GenericRpcHandlers with this Server. 

35 

36 This method is only safe to call before the server is started. 

37 

38 Args: 

39 generic_rpc_handlers: A sequence of GenericRpcHandlers that will be 

40 used to service RPCs. 

41 """ 

42 

43 @abc.abstractmethod 

44 def add_insecure_port(self, address: str) -> int: 

45 """Opens an insecure port for accepting RPCs. 

46 

47 A port is a communication endpoint that used by networking protocols, 

48 like TCP and UDP. To date, we only support TCP. 

49 

50 This method may only be called before starting the server. 

51 

52 Args: 

53 address: The address for which to open a port. If the port is 0, 

54 or not specified in the address, then the gRPC runtime will choose a port. 

55 

56 Returns: 

57 An integer port on which the server will accept RPC requests. 

58 """ 

59 

60 @abc.abstractmethod 

61 def add_secure_port( 

62 self, address: str, server_credentials: grpc.ServerCredentials 

63 ) -> int: 

64 """Opens a secure port for accepting RPCs. 

65 

66 A port is a communication endpoint that used by networking protocols, 

67 like TCP and UDP. To date, we only support TCP. 

68 

69 This method may only be called before starting the server. 

70 

71 Args: 

72 address: The address for which to open a port. 

73 if the port is 0, or not specified in the address, then the gRPC 

74 runtime will choose a port. 

75 server_credentials: A ServerCredentials object. 

76 

77 Returns: 

78 An integer port on which the server will accept RPC requests. 

79 """ 

80 

81 @abc.abstractmethod 

82 async def start(self) -> None: 

83 """Starts this Server. 

84 

85 This method may only be called once. (i.e. it is not idempotent). 

86 """ 

87 

88 @abc.abstractmethod 

89 async def stop(self, grace: Optional[float]) -> None: 

90 """Stops this Server. 

91 

92 This method immediately stops the server from servicing new RPCs in 

93 all cases. 

94 

95 If a grace period is specified, this method waits until all active 

96 RPCs are finished or until the grace period is reached. RPCs that haven't 

97 been terminated within the grace period are aborted. 

98 If a grace period is not specified (by passing None for grace), all 

99 existing RPCs are aborted immediately and this method blocks until 

100 the last RPC handler terminates. 

101 

102 This method is idempotent and may be called at any time. Passing a 

103 smaller grace value in a subsequent call will have the effect of 

104 stopping the Server sooner (passing None will have the effect of 

105 stopping the server immediately). Passing a larger grace value in a 

106 subsequent call will not have the effect of stopping the server later 

107 (i.e. the most restrictive grace value is used). 

108 

109 Args: 

110 grace: A duration of time in seconds or None. 

111 """ 

112 

113 @abc.abstractmethod 

114 async def wait_for_termination( 

115 self, timeout: Optional[float] = None 

116 ) -> bool: 

117 """Continues current coroutine once the server stops. 

118 

119 This is an EXPERIMENTAL API. 

120 

121 The wait will not consume computational resources during blocking, and 

122 it will block until one of the two following conditions are met: 

123 

124 1) The server is stopped or terminated; 

125 2) A timeout occurs if timeout is not `None`. 

126 

127 The timeout argument works in the same way as `threading.Event.wait()`. 

128 https://docs.python.org/3/library/threading.html#threading.Event.wait 

129 

130 Args: 

131 timeout: A floating point number specifying a timeout for the 

132 operation in seconds. 

133 

134 Returns: 

135 A bool indicates if the operation times out. 

136 """ 

137 

138 # Suppressing pyright[reportUnknownParameterType, reportMissingParameterType] 

139 # for type annotation of service_name and method_handlers as it will be 

140 # taken up along with the sync stack changes. 

141 def add_registered_method_handlers( # noqa: B027 

142 self, service_name, method_handlers # pyright: ignore # noqa: PGH003 

143 ) -> None: 

144 """Registers GenericRpcHandlers with this Server. 

145 

146 This method is only safe to call before the server is started. 

147 

148 Args: 

149 service_name: The service name. 

150 method_handlers: A dictionary that maps method names to corresponding 

151 RpcMethodHandler. 

152 """ 

153 

154 

155# pylint: disable=too-many-public-methods 

156class ServicerContext(Generic[RequestType, ResponseType], abc.ABC): 

157 """A context object passed to method implementations.""" 

158 

159 @abc.abstractmethod 

160 async def read(self) -> RequestType: 

161 """Reads one message from the RPC. 

162 

163 Only one read operation is allowed simultaneously. 

164 

165 Returns: 

166 A response message of the RPC. 

167 

168 Raises: 

169 An RpcError exception if the read failed. 

170 """ 

171 

172 @abc.abstractmethod 

173 async def write(self, message: ResponseType) -> None: 

174 """Writes one message to the RPC. 

175 

176 Only one write operation is allowed simultaneously. 

177 

178 Raises: 

179 An RpcError exception if the write failed. 

180 """ 

181 

182 @abc.abstractmethod 

183 async def send_initial_metadata( 

184 self, initial_metadata: MetadataType 

185 ) -> None: 

186 """Sends the initial metadata value to the client. 

187 

188 This method need not be called by implementations if they have no 

189 metadata to add to what the gRPC runtime will transmit. 

190 

191 Args: 

192 initial_metadata: The initial :term:`metadata`. 

193 """ 

194 

195 @abc.abstractmethod 

196 async def abort( 

197 self, 

198 code: grpc.StatusCode, 

199 details: str = "", 

200 trailing_metadata: MetadataType = (), 

201 ) -> NoReturn: 

202 """Raises an exception to terminate the RPC with a non-OK status. 

203 

204 The code and details passed as arguments will supersede any existing 

205 ones. 

206 

207 Args: 

208 code: A StatusCode object to be sent to the client. 

209 It must not be StatusCode.OK. 

210 details: A UTF-8-encodable string to be sent to the client upon 

211 termination of the RPC. 

212 trailing_metadata: A sequence of tuple represents the trailing 

213 :term:`metadata`. 

214 

215 Raises: 

216 Exception: An exception is always raised to signal the abortion the 

217 RPC to the gRPC runtime. 

218 """ 

219 raise NotImplementedError() 

220 

221 @abc.abstractmethod 

222 async def abort_with_status(self, status: grpc.Status) -> NoReturn: 

223 """Raises an exception to terminate the RPC with a non-OK status. 

224 

225 The status passed as argument will supersede any existing status code, 

226 status message and trailing metadata. 

227 

228 This is an EXPERIMENTAL API. 

229 

230 Args: 

231 status: A grpc.Status object. The status code in it must not be 

232 StatusCode.OK. 

233 

234 Raises: 

235 Exception: An exception is always raised to signal the abortion of the 

236 RPC to the gRPC runtime. 

237 """ 

238 raise NotImplementedError() 

239 

240 @abc.abstractmethod 

241 def set_trailing_metadata(self, trailing_metadata: MetadataType) -> None: 

242 """Sends the trailing metadata for the RPC. 

243 

244 This method need not be called by implementations if they have no 

245 metadata to add to what the gRPC runtime will transmit. 

246 

247 Args: 

248 trailing_metadata: The trailing :term:`metadata`. 

249 """ 

250 

251 @abc.abstractmethod 

252 def invocation_metadata(self) -> Optional[MetadataType]: 

253 """Accesses the metadata sent by the client. 

254 

255 Returns: 

256 The invocation :term:`metadata`. 

257 """ 

258 

259 @abc.abstractmethod 

260 def set_code(self, code: grpc.StatusCode) -> None: 

261 """Sets the value to be used as status code upon RPC completion. 

262 

263 This method need not be called by method implementations if they wish 

264 the gRPC runtime to determine the status code of the RPC. 

265 

266 Args: 

267 code: A StatusCode object to be sent to the client. 

268 """ 

269 

270 @abc.abstractmethod 

271 def set_details(self, details: str) -> None: 

272 """Sets the value to be used the as detail string upon RPC completion. 

273 

274 This method need not be called by method implementations if they have 

275 no details to transmit. 

276 

277 Args: 

278 details: A UTF-8-encodable string to be sent to the client upon 

279 termination of the RPC. 

280 """ 

281 

282 @abc.abstractmethod 

283 def set_compression(self, compression: grpc.Compression) -> None: 

284 """Set the compression algorithm to be used for the entire call. 

285 

286 Args: 

287 compression: An element of grpc.compression, e.g. 

288 grpc.compression.Gzip. 

289 """ 

290 

291 @abc.abstractmethod 

292 def disable_next_message_compression(self) -> None: 

293 """Disables compression for the next response message. 

294 

295 This method will override any compression configuration set during 

296 server creation or set on the call. 

297 """ 

298 

299 @abc.abstractmethod 

300 def peer(self) -> str: 

301 """Identifies the peer that invoked the RPC being serviced. 

302 

303 Returns: 

304 A string identifying the peer that invoked the RPC being serviced. 

305 The string format is determined by gRPC runtime. 

306 """ 

307 

308 @abc.abstractmethod 

309 def peer_identities(self) -> Optional[Iterable[bytes]]: 

310 """Gets one or more peer identity(s). 

311 

312 Equivalent to 

313 servicer_context.auth_context().get(servicer_context.peer_identity_key()) 

314 

315 Returns: 

316 An iterable of the identities, or None if the call is not 

317 authenticated. Each identity is returned as a raw bytes type. 

318 """ 

319 

320 @abc.abstractmethod 

321 def peer_identity_key(self) -> Optional[str]: 

322 """The auth property used to identify the peer. 

323 

324 For example, "x509_common_name" or "x509_subject_alternative_name" are 

325 used to identify an SSL peer. 

326 

327 Returns: 

328 The auth property (string) that indicates the 

329 peer identity, or None if the call is not authenticated. 

330 """ 

331 

332 @abc.abstractmethod 

333 def auth_context(self) -> Mapping[str, Iterable[bytes]]: 

334 """Gets the auth context for the call. 

335 

336 Returns: 

337 A map of strings to an iterable of bytes for each auth property. 

338 """ 

339 

340 def time_remaining(self) -> float: 

341 """Describes the length of allowed time remaining for the RPC. 

342 

343 Returns: 

344 A nonnegative float indicating the length of allowed time in seconds 

345 remaining for the RPC to complete before it is considered to have 

346 timed out, or None if no deadline was specified for the RPC. 

347 """ 

348 raise NotImplementedError() 

349 

350 def trailing_metadata(self) -> MetadataType: 

351 """Access value to be used as trailing metadata upon RPC completion. 

352 

353 This is an EXPERIMENTAL API. 

354 

355 Returns: 

356 The trailing :term:`metadata` for the RPC. 

357 """ 

358 raise NotImplementedError() 

359 

360 def code(self) -> grpc.StatusCode: 

361 """Accesses the value to be used as status code upon RPC completion. 

362 

363 This is an EXPERIMENTAL API. 

364 

365 Returns: 

366 The StatusCode value for the RPC. 

367 """ 

368 raise NotImplementedError() 

369 

370 def details(self): # pyright: ignore[reportUnknownParameterType] 

371 """Accesses the value to be used as detail string upon RPC completion. 

372 

373 This is an EXPERIMENTAL API. 

374 

375 Returns: 

376 The details string of the RPC. 

377 """ 

378 raise NotImplementedError() 

379 

380 def add_done_callback(self, callback: DoneCallbackType) -> None: 

381 """Registers a callback to be called on RPC termination. 

382 

383 This is an EXPERIMENTAL API. 

384 

385 Args: 

386 callback: A callable object will be called with the servicer context 

387 object as its only argument. 

388 """ 

389 

390 def cancelled(self) -> bool: 

391 """Return True if the RPC is cancelled. 

392 

393 The RPC is cancelled when the cancellation was requested with cancel(). 

394 

395 This is an EXPERIMENTAL API. 

396 

397 Returns: 

398 A bool indicates whether the RPC is cancelled or not. 

399 """ 

400 raise NotImplementedError() 

401 

402 def done(self) -> bool: 

403 """Return True if the RPC is done. 

404 

405 An RPC is done if the RPC is completed, cancelled or aborted. 

406 

407 This is an EXPERIMENTAL API. 

408 

409 Returns: 

410 A bool indicates if the RPC is done. 

411 """ 

412 raise NotImplementedError()