Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/aiohttp/client_proto.py: 19%

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

192 statements  

1import asyncio 

2from contextlib import suppress 

3from typing import Any, Callable 

4 

5from .base_protocol import BaseProtocol 

6from .client_exceptions import ( 

7 ClientConnectionError, 

8 ClientOSError, 

9 ClientPayloadError, 

10 ServerDisconnectedError, 

11 SocketTimeoutError, 

12) 

13from .helpers import ( 

14 _EXC_SENTINEL, 

15 DEFAULT_CHUNK_SIZE, 

16 EMPTY_BODY_STATUS_CODES, 

17 BaseTimerContext, 

18 set_exception, 

19 set_result, 

20) 

21from .http import HttpResponseParser, RawResponseMessage 

22from .http_exceptions import HttpProcessingError 

23from .streams import EMPTY_PAYLOAD, DataQueue, StreamReader 

24 

25 

26class ResponseHandler(BaseProtocol, DataQueue[tuple[RawResponseMessage, StreamReader]]): 

27 """Helper class to adapt between Protocol and StreamReader.""" 

28 

29 def __init__(self, loop: asyncio.AbstractEventLoop) -> None: 

30 BaseProtocol.__init__(self, loop=loop, parser=None) 

31 DataQueue.__init__(self, loop) 

32 

33 self._should_close = False 

34 

35 self._payload: StreamReader | None = None 

36 self._skip_payload = False 

37 self._payload_parser = None 

38 self._data_received_cb: Callable[[], None] | None = None 

39 

40 self._timer = None 

41 self._tail = b"" 

42 

43 self._read_timeout: float | None = None 

44 self._read_timeout_handle: asyncio.TimerHandle | None = None 

45 

46 self._timeout_ceil_threshold: float | None = 5 

47 

48 self._closed: None | asyncio.Future[None] = None 

49 self._connection_lost_called = False 

50 

51 @property 

52 def closed(self) -> None | asyncio.Future[None]: 

53 """Future that is set when the connection is closed. 

54 

55 This property returns a Future that will be completed when the connection 

56 is closed. The Future is created lazily on first access to avoid creating 

57 futures that will never be awaited. 

58 

59 Returns: 

60 - A Future[None] if the connection is still open or was closed after 

61 this property was accessed 

62 - None if connection_lost() was already called before this property 

63 was ever accessed (indicating no one is waiting for the closure) 

64 """ 

65 if self._closed is None and not self._connection_lost_called: 

66 self._closed = self._loop.create_future() 

67 return self._closed 

68 

69 @property 

70 def upgraded(self) -> bool: 

71 return self._upgraded 

72 

73 @property 

74 def should_close(self) -> bool: 

75 return bool( 

76 self._should_close 

77 or (self._payload is not None and not self._payload.is_eof()) 

78 or self._upgraded 

79 or self._exception is not None 

80 or self._payload_parser is not None 

81 or self._buffer 

82 or self._tail 

83 ) 

84 

85 def force_close(self) -> None: 

86 self._should_close = True 

87 

88 def close(self) -> None: 

89 self._exception = None # Break cyclic references 

90 transport = self.transport 

91 if transport is not None: 

92 transport.close() 

93 self.transport = None 

94 self._payload = None 

95 self._drop_timeout() 

96 

97 def abort(self) -> None: 

98 self._exception = None # Break cyclic references 

99 transport = self.transport 

100 if transport is not None: 

101 transport.abort() 

102 self.transport = None 

103 self._payload = None 

104 self._drop_timeout() 

105 

106 def is_connected(self) -> bool: 

107 return self.transport is not None and not self.transport.is_closing() 

108 

109 def connection_lost(self, exc: BaseException | None) -> None: 

110 self._connection_lost_called = True 

111 self._drop_timeout() 

112 

113 original_connection_error = exc 

114 reraised_exc = original_connection_error 

115 

116 connection_closed_cleanly = original_connection_error is None 

117 

118 if self._closed is not None: 

119 # If someone is waiting for the closed future, 

120 # we should set it to None or an exception. If 

121 # self._closed is None, it means that 

122 # connection_lost() was called already 

123 # or nobody is waiting for it. 

124 if connection_closed_cleanly: 

125 set_result(self._closed, None) 

126 else: 

127 assert original_connection_error is not None 

128 set_exception( 

129 self._closed, 

130 ClientConnectionError( 

131 f"Connection lost: {original_connection_error !s}", 

132 ), 

133 original_connection_error, 

134 ) 

135 

136 if self._payload_parser is not None: 

137 with suppress(Exception): # FIXME: log this somehow? 

138 self._payload_parser.feed_eof() 

139 

140 uncompleted = None 

141 if self._parser is not None: 

142 try: 

143 uncompleted = self._parser.feed_eof() 

144 except Exception as underlying_exc: 

145 if self._payload is not None: 

146 client_payload_exc_msg = ( 

147 f"Response payload is not completed: {underlying_exc !r}" 

148 ) 

149 if not connection_closed_cleanly: 

150 client_payload_exc_msg = ( 

151 f"{client_payload_exc_msg !s}. " 

152 f"{original_connection_error !r}" 

153 ) 

154 set_exception( 

155 self._payload, 

156 ClientPayloadError(client_payload_exc_msg), 

157 underlying_exc, 

158 ) 

159 

160 if not self.is_eof(): 

161 if isinstance(original_connection_error, OSError): 

162 reraised_exc = ClientOSError(*original_connection_error.args) 

163 if connection_closed_cleanly: 

164 reraised_exc = ServerDisconnectedError(uncompleted) 

165 # assigns self._should_close to True as side effect, 

166 # we do it anyway below 

167 underlying_non_eof_exc = ( 

168 _EXC_SENTINEL 

169 if connection_closed_cleanly 

170 else original_connection_error 

171 ) 

172 assert underlying_non_eof_exc is not None 

173 assert reraised_exc is not None 

174 self.set_exception(reraised_exc, underlying_non_eof_exc) 

175 

176 self._should_close = True 

177 self._parser = None 

178 self._payload = None 

179 self._payload_parser = None 

180 self._reading_paused = False 

181 

182 super().connection_lost(reraised_exc) 

183 

184 def eof_received(self) -> None: 

185 # should call parser.feed_eof() most likely 

186 self._drop_timeout() 

187 

188 def pause_reading(self) -> None: 

189 super().pause_reading() 

190 self._drop_timeout() 

191 

192 def resume_reading(self, resume_parser: bool = True) -> None: 

193 was_paused = self._reading_paused 

194 super().resume_reading(resume_parser) 

195 if was_paused: 

196 self._reschedule_timeout() 

197 

198 def set_exception( 

199 self, 

200 exc: BaseException, 

201 exc_cause: BaseException = _EXC_SENTINEL, 

202 ) -> None: 

203 self._should_close = True 

204 self._drop_timeout() 

205 super().set_exception(exc, exc_cause) 

206 

207 def set_parser( 

208 self, 

209 parser: Any, 

210 payload: Any, 

211 data_received_cb: Callable[[], None] | None = None, 

212 ) -> None: 

213 # TODO: actual types are: 

214 # parser: WebSocketReader 

215 # payload: WebSocketDataQueue 

216 # but they are not generi enough 

217 # Need an ABC for both types 

218 self._payload = payload 

219 self._payload_parser = parser 

220 self._data_received_cb = data_received_cb 

221 

222 self._drop_timeout() 

223 

224 if self._tail: 

225 data, self._tail = self._tail, b"" 

226 self.data_received(data) 

227 

228 def set_response_params( 

229 self, 

230 *, 

231 timer: BaseTimerContext | None = None, 

232 skip_payload: bool = False, 

233 read_until_eof: bool = False, 

234 auto_decompress: bool = True, 

235 read_timeout: float | None = None, 

236 read_bufsize: int = DEFAULT_CHUNK_SIZE, 

237 timeout_ceil_threshold: float = 5, 

238 max_line_size: int = 8190, 

239 max_field_size: int = 8190, 

240 max_headers: int = 128, 

241 ) -> None: 

242 self._skip_payload = skip_payload 

243 

244 self._read_timeout = read_timeout 

245 

246 self._timeout_ceil_threshold = timeout_ceil_threshold 

247 

248 self._parser = HttpResponseParser( 

249 self, 

250 self._loop, 

251 read_bufsize, 

252 timer=timer, 

253 payload_exception=ClientPayloadError, 

254 response_with_body=not skip_payload, 

255 read_until_eof=read_until_eof, 

256 auto_decompress=auto_decompress, 

257 max_line_size=max_line_size, 

258 max_field_size=max_field_size, 

259 max_headers=max_headers, 

260 ) 

261 

262 if self._tail: 

263 data, self._tail = self._tail, b"" 

264 self.data_received(data) 

265 

266 def _drop_timeout(self) -> None: 

267 if self._read_timeout_handle is not None: 

268 self._read_timeout_handle.cancel() 

269 self._read_timeout_handle = None 

270 

271 def _reschedule_timeout(self) -> None: 

272 timeout = self._read_timeout 

273 if self._read_timeout_handle is not None: 

274 self._read_timeout_handle.cancel() 

275 

276 if timeout: 

277 self._read_timeout_handle = self._loop.call_later( 

278 timeout, self._on_read_timeout 

279 ) 

280 else: 

281 self._read_timeout_handle = None 

282 

283 def start_timeout(self) -> None: 

284 self._reschedule_timeout() 

285 

286 @property 

287 def read_timeout(self) -> float | None: 

288 return self._read_timeout 

289 

290 @read_timeout.setter 

291 def read_timeout(self, read_timeout: float | None) -> None: 

292 self._read_timeout = read_timeout 

293 

294 def _on_read_timeout(self) -> None: 

295 exc = SocketTimeoutError("Timeout on reading data from socket") 

296 self.set_exception(exc) 

297 if self._payload is not None: 

298 set_exception(self._payload, exc) 

299 

300 def data_received(self, data: bytes) -> None: 

301 # If no data, then we are resuming decompression. We haven't received 

302 # data from the socket, so we can avoid the reschedule overhead. 

303 if data: 

304 self._reschedule_timeout() 

305 

306 # custom payload parser - currently always WebSocketReader 

307 if self._payload_parser is not None: 

308 if self._data_received_cb is not None: 

309 self._data_received_cb() 

310 eof, tail = self._payload_parser.feed_data(data) 

311 if eof: 

312 self._payload = None 

313 self._payload_parser = None 

314 

315 if tail: 

316 self.data_received(tail) 

317 return 

318 

319 if self._upgraded or self._parser is None: 

320 # i.e. websocket connection, websocket parser is not set yet 

321 self._tail += data 

322 return 

323 

324 # parse http messages 

325 try: 

326 messages, upgraded, tail = self._parser.feed_data(data) 

327 except BaseException as underlying_exc: 

328 if self.transport is not None: 

329 # connection.release() could be called BEFORE 

330 # data_received(), the transport is already 

331 # closed in this case 

332 self.transport.close() 

333 if not isinstance(underlying_exc, Exception): 

334 raise 

335 # should_close is True after the call 

336 if isinstance(underlying_exc, HttpProcessingError): 

337 exc = HttpProcessingError( 

338 code=underlying_exc.code, 

339 message=underlying_exc.message, 

340 headers=underlying_exc.headers, 

341 ) 

342 else: 

343 exc = HttpProcessingError() 

344 self.set_exception(exc, underlying_exc) 

345 return 

346 

347 self._upgraded = upgraded 

348 

349 payload: StreamReader | None = None 

350 for message, payload in messages: 

351 if message.should_close: 

352 self._should_close = True 

353 

354 self._payload = payload 

355 

356 if self._skip_payload or message.code in EMPTY_BODY_STATUS_CODES: 

357 self.feed_data((message, EMPTY_PAYLOAD), 0) 

358 else: 

359 self.feed_data((message, payload), 0) 

360 

361 if payload is not None: 

362 # new message(s) was processed 

363 # register timeout handler unsubscribing 

364 # either on end-of-stream or immediately for 

365 # EMPTY_PAYLOAD 

366 if payload is not EMPTY_PAYLOAD: 

367 payload.on_eof(self._drop_timeout) 

368 else: 

369 self._drop_timeout() 

370 

371 if upgraded and tail: 

372 self.data_received(tail)