Coverage for Lib/asyncio/selector_events.py: 88%

932 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-25 03:29 +0000

1"""Event loop using a selector and related classes. 

2 

3A selector is a "notify-when-ready" multiplexer. For a subclass which 

4also includes support for signal handling, see the unix_events sub-module. 

5""" 

6 

7__all__ = 'BaseSelectorEventLoop', 

8 

9import collections 

10import errno 

11import functools 

12import itertools 

13import os 

14import selectors 

15import socket 

16import warnings 

17import weakref 

18try: 

19 import ssl 

20except ImportError: # pragma: no cover 

21 ssl = None 

22 

23from . import base_events 

24from . import constants 

25from . import events 

26from . import futures 

27from . import protocols 

28from . import sslproto 

29from . import transports 

30from . import trsock 

31from .log import logger 

32 

33_HAS_SENDMSG = hasattr(socket.socket, 'sendmsg') 

34 

35if _HAS_SENDMSG: 35 ↛ 42line 35 didn't jump to line 42 because the condition on line 35 was always true

36 try: 

37 SC_IOV_MAX = os.sysconf('SC_IOV_MAX') 

38 except OSError: 

39 # Fallback to send 

40 _HAS_SENDMSG = False 

41 

42def _test_selector_event(selector, fd, event): 

43 # Test if the selector is monitoring 'event' events 

44 # for the file descriptor 'fd'. 

45 try: 

46 key = selector.get_key(fd) 

47 except KeyError: 

48 return False 

49 else: 

50 return bool(key.events & event) 

51 

52 

53class BaseSelectorEventLoop(base_events.BaseEventLoop): 

54 """Selector event loop. 

55 

56 See events.AbstractEventLoop for API specification. 

57 """ 

58 

59 def __init__(self, selector=None): 

60 super().__init__() 

61 

62 if selector is None: 

63 selector = selectors.DefaultSelector() 

64 logger.debug('Using selector: %s', selector.__class__.__name__) 

65 self._selector = selector 

66 self._make_self_pipe() 

67 self._transports = weakref.WeakValueDictionary() 

68 

69 def _make_socket_transport(self, sock, protocol, waiter=None, *, 

70 extra=None, server=None, context=None): 

71 self._ensure_fd_no_transport(sock) 

72 return _SelectorSocketTransport(self, sock, protocol, waiter, 

73 extra, server, context=context) 

74 

75 def _make_ssl_transport( 

76 self, rawsock, protocol, sslcontext, waiter=None, 

77 *, server_side=False, server_hostname=None, 

78 extra=None, server=None, 

79 ssl_handshake_timeout=constants.SSL_HANDSHAKE_TIMEOUT, 

80 ssl_shutdown_timeout=constants.SSL_SHUTDOWN_TIMEOUT, 

81 context=None, 

82 ): 

83 self._ensure_fd_no_transport(rawsock) 

84 ssl_protocol = sslproto.SSLProtocol( 

85 self, protocol, sslcontext, waiter, 

86 server_side, server_hostname, 

87 ssl_handshake_timeout=ssl_handshake_timeout, 

88 ssl_shutdown_timeout=ssl_shutdown_timeout, 

89 ) 

90 _SelectorSocketTransport(self, rawsock, ssl_protocol, 

91 extra=extra, server=server, context=context) 

92 return ssl_protocol._app_transport 

93 

94 def _make_datagram_transport(self, sock, protocol, 

95 address=None, waiter=None, extra=None): 

96 self._ensure_fd_no_transport(sock) 

97 return _SelectorDatagramTransport(self, sock, protocol, 

98 address, waiter, extra) 

99 

100 def close(self): 

101 if self.is_running(): 

102 raise RuntimeError("Cannot close a running event loop") 

103 if self.is_closed(): 

104 return 

105 self._close_self_pipe() 

106 super().close() 

107 if self._selector is not None: 

108 self._selector.close() 

109 self._selector = None 

110 

111 def _close_self_pipe(self): 

112 self._remove_reader(self._ssock.fileno()) 

113 self._ssock.close() 

114 self._ssock = None 

115 self._csock.close() 

116 self._csock = None 

117 self._internal_fds -= 1 

118 

119 def _make_self_pipe(self): 

120 # A self-socket, really. :-) 

121 self._ssock, self._csock = socket.socketpair() 

122 self._ssock.setblocking(False) 

123 self._csock.setblocking(False) 

124 self._internal_fds += 1 

125 self._add_reader(self._ssock.fileno(), self._read_from_self) 

126 

127 def _process_self_data(self, data): 

128 pass 

129 

130 def _read_from_self(self): 

131 while True: 

132 try: 

133 data = self._ssock.recv(4096) 

134 if not data: 134 ↛ 135line 134 didn't jump to line 135 because the condition on line 134 was never true

135 break 

136 self._process_self_data(data) 

137 except InterruptedError: 

138 continue 

139 except BlockingIOError: 

140 break 

141 

142 def _write_to_self(self): 

143 # This may be called from a different thread, possibly after 

144 # _close_self_pipe() has been called or even while it is 

145 # running. Guard for self._csock being None or closed. When 

146 # a socket is closed, send() raises OSError (with errno set to 

147 # EBADF, but let's not rely on the exact error code). 

148 csock = self._csock 

149 if csock is None: 149 ↛ 150line 149 didn't jump to line 150 because the condition on line 149 was never true

150 return 

151 

152 try: 

153 csock.send(b'\0') 

154 except OSError: 

155 if self._debug: 155 ↛ 156line 155 didn't jump to line 156 because the condition on line 155 was never true

156 logger.debug("Fail to write a null byte into the " 

157 "self-pipe socket", 

158 exc_info=True) 

159 

160 def _start_serving(self, protocol_factory, sock, 

161 sslcontext=None, server=None, backlog=100, 

162 ssl_handshake_timeout=constants.SSL_HANDSHAKE_TIMEOUT, 

163 ssl_shutdown_timeout=constants.SSL_SHUTDOWN_TIMEOUT, context=None): 

164 self._add_reader(sock.fileno(), self._accept_connection, 

165 protocol_factory, sock, sslcontext, server, backlog, 

166 ssl_handshake_timeout, ssl_shutdown_timeout, context) 

167 

168 def _accept_connection( 

169 self, protocol_factory, sock, 

170 sslcontext=None, server=None, backlog=100, 

171 ssl_handshake_timeout=constants.SSL_HANDSHAKE_TIMEOUT, 

172 ssl_shutdown_timeout=constants.SSL_SHUTDOWN_TIMEOUT, context=None): 

173 # This method is only called once for each event loop tick where the 

174 # listening socket has triggered an EVENT_READ. There may be multiple 

175 # connections waiting for an .accept() so it is called in a loop. 

176 # See https://bugs.python.org/issue27906 for more details. 

177 for _ in range(backlog + 1): 

178 try: 

179 conn, addr = sock.accept() 

180 if self._debug: 

181 logger.debug("%r got a new connection from %r: %r", 

182 server, addr, conn) 

183 conn.setblocking(False) 

184 except ConnectionAbortedError: 

185 # Discard connections that were aborted before accept(). 

186 continue 

187 except (BlockingIOError, InterruptedError): 

188 # Early exit because of a signal or 

189 # the socket accept buffer is empty. 

190 return 

191 except OSError as exc: 

192 # There's nowhere to send the error, so just log it. 

193 if exc.errno in (errno.EMFILE, errno.ENFILE, 193 ↛ 211line 193 didn't jump to line 211 because the condition on line 193 was always true

194 errno.ENOBUFS, errno.ENOMEM): 

195 # Some platforms (e.g. Linux keep reporting the FD as 

196 # ready, so we remove the read handler temporarily. 

197 # We'll try again in a while. 

198 self.call_exception_handler({ 

199 'message': 'socket.accept() out of system resource', 

200 'exception': exc, 

201 'socket': trsock.TransportSocket(sock), 

202 }) 

203 self._remove_reader(sock.fileno()) 

204 self.call_later(constants.ACCEPT_RETRY_DELAY, 

205 self._start_serving, 

206 protocol_factory, sock, sslcontext, server, 

207 backlog, ssl_handshake_timeout, 

208 ssl_shutdown_timeout, context) 

209 return 

210 else: 

211 raise # The event loop will catch, log and ignore it. 

212 else: 

213 extra = {'peername': addr} 

214 conn_context = context.copy() if context is not None else None 

215 accept = self._accept_connection2( 

216 protocol_factory, conn, extra, sslcontext, server, 

217 ssl_handshake_timeout, ssl_shutdown_timeout, context=conn_context) 

218 self.create_task(accept, context=conn_context) 

219 

220 async def _accept_connection2( 

221 self, protocol_factory, conn, extra, 

222 sslcontext=None, server=None, 

223 ssl_handshake_timeout=constants.SSL_HANDSHAKE_TIMEOUT, 

224 ssl_shutdown_timeout=constants.SSL_SHUTDOWN_TIMEOUT, context=None): 

225 protocol = None 

226 transport = None 

227 try: 

228 protocol = protocol_factory() 

229 waiter = self.create_future() 

230 if sslcontext: 

231 transport = self._make_ssl_transport( 

232 conn, protocol, sslcontext, waiter=waiter, 

233 server_side=True, extra=extra, server=server, 

234 ssl_handshake_timeout=ssl_handshake_timeout, 

235 ssl_shutdown_timeout=ssl_shutdown_timeout, 

236 context=context) 

237 else: 

238 transport = self._make_socket_transport( 

239 conn, protocol, waiter=waiter, extra=extra, 

240 server=server, context=context) 

241 

242 try: 

243 await waiter 

244 except BaseException: 

245 transport.close() 

246 # gh-109534: When an exception is raised by the SSLProtocol object the 

247 # exception set in this future can keep the protocol object alive and 

248 # cause a reference cycle. 

249 waiter = None 

250 raise 

251 # It's now up to the protocol to handle the connection. 

252 

253 except (SystemExit, KeyboardInterrupt): 

254 raise 

255 except BaseException as exc: 

256 if transport is None: 

257 conn.close() 

258 if transport is None or self._debug: 

259 context = { 

260 'message': 

261 'Error on transport creation for incoming connection', 

262 'exception': exc, 

263 } 

264 if protocol is not None: 

265 context['protocol'] = protocol 

266 if transport is not None: 266 ↛ 267line 266 didn't jump to line 267 because the condition on line 266 was never true

267 context['transport'] = transport 

268 self.call_exception_handler(context) 

269 

270 def _ensure_fd_no_transport(self, fd): 

271 fileno = fd 

272 if not isinstance(fileno, int): 

273 try: 

274 fileno = int(fileno.fileno()) 

275 except (AttributeError, TypeError, ValueError): 

276 # This code matches selectors._fileobj_to_fd function. 

277 raise ValueError(f"Invalid file object: {fd!r}") from None 

278 transport = self._transports.get(fileno) 

279 if transport and not transport.is_closing(): 

280 raise RuntimeError( 

281 f'File descriptor {fd!r} is used by transport ' 

282 f'{transport!r}') 

283 

284 def _add_reader(self, fd, callback, *args, context=None): 

285 self._check_closed() 

286 handle = events.Handle(callback, args, self, context=context) 

287 key = self._selector.get_map().get(fd) 

288 if key is None: 

289 self._selector.register(fd, selectors.EVENT_READ, 

290 (handle, None)) 

291 else: 

292 mask, (reader, writer) = key.events, key.data 

293 self._selector.modify(fd, mask | selectors.EVENT_READ, 

294 (handle, writer)) 

295 if reader is not None: 

296 reader.cancel() 

297 return handle 

298 

299 def _remove_reader(self, fd): 

300 if self.is_closed(): 

301 return False 

302 key = self._selector.get_map().get(fd) 

303 if key is None: 

304 return False 

305 mask, (reader, writer) = key.events, key.data 

306 mask &= ~selectors.EVENT_READ 

307 if not mask: 

308 self._selector.unregister(fd) 

309 else: 

310 self._selector.modify(fd, mask, (None, writer)) 

311 

312 if reader is not None: 

313 reader.cancel() 

314 return True 

315 else: 

316 return False 

317 

318 def _add_writer(self, fd, callback, *args, context=None): 

319 self._check_closed() 

320 handle = events.Handle(callback, args, self, context=context) 

321 key = self._selector.get_map().get(fd) 

322 if key is None: 

323 self._selector.register(fd, selectors.EVENT_WRITE, 

324 (None, handle)) 

325 else: 

326 mask, (reader, writer) = key.events, key.data 

327 self._selector.modify(fd, mask | selectors.EVENT_WRITE, 

328 (reader, handle)) 

329 if writer is not None: 

330 writer.cancel() 

331 return handle 

332 

333 def _remove_writer(self, fd): 

334 """Remove a writer callback.""" 

335 if self.is_closed(): 

336 return False 

337 key = self._selector.get_map().get(fd) 

338 if key is None: 

339 return False 

340 mask, (reader, writer) = key.events, key.data 

341 # Remove both writer and connector. 

342 mask &= ~selectors.EVENT_WRITE 

343 if not mask: 

344 self._selector.unregister(fd) 

345 else: 

346 self._selector.modify(fd, mask, (reader, None)) 

347 

348 if writer is not None: 

349 writer.cancel() 

350 return True 

351 else: 

352 return False 

353 

354 def add_reader(self, fd, callback, *args): 

355 """Add a reader callback.""" 

356 self._ensure_fd_no_transport(fd) 

357 self._add_reader(fd, callback, *args) 

358 

359 def remove_reader(self, fd): 

360 """Remove a reader callback.""" 

361 self._ensure_fd_no_transport(fd) 

362 return self._remove_reader(fd) 

363 

364 def add_writer(self, fd, callback, *args): 

365 """Add a writer callback..""" 

366 self._ensure_fd_no_transport(fd) 

367 self._add_writer(fd, callback, *args) 

368 

369 def remove_writer(self, fd): 

370 """Remove a writer callback.""" 

371 self._ensure_fd_no_transport(fd) 

372 return self._remove_writer(fd) 

373 

374 async def sock_recv(self, sock, n): 

375 """Receive data from the socket. 

376 

377 The return value is a bytes object representing the data received. 

378 The maximum amount of data to be received at once is specified by 

379 nbytes. 

380 """ 

381 base_events._check_ssl_socket(sock) 

382 if self._debug and sock.gettimeout() != 0: 

383 raise ValueError("the socket must be non-blocking") 

384 try: 

385 return sock.recv(n) 

386 except (BlockingIOError, InterruptedError): 

387 pass 

388 fut = self.create_future() 

389 fd = sock.fileno() 

390 self._ensure_fd_no_transport(fd) 

391 handle = self._add_reader(fd, self._sock_recv, fut, sock, n) 

392 fut.add_done_callback( 

393 functools.partial(self._sock_read_done, fd, handle=handle)) 

394 return await fut 

395 

396 def _sock_read_done(self, fd, fut, handle=None): 

397 if handle is None or not handle.cancelled(): 

398 self.remove_reader(fd) 

399 

400 def _sock_recv(self, fut, sock, n): 

401 # _sock_recv() can add itself as an I/O callback if the operation can't 

402 # be done immediately. Don't use it directly, call sock_recv(). 

403 if fut.done(): 403 ↛ 404line 403 didn't jump to line 404 because the condition on line 403 was never true

404 return 

405 try: 

406 data = sock.recv(n) 

407 except (BlockingIOError, InterruptedError): 

408 return # try again next time 

409 except (SystemExit, KeyboardInterrupt): 

410 raise 

411 except BaseException as exc: 

412 fut.set_exception(exc) 

413 else: 

414 fut.set_result(data) 

415 

416 async def sock_recv_into(self, sock, buf): 

417 """Receive data from the socket. 

418 

419 The received data is written into *buf* (a writable buffer). 

420 The return value is the number of bytes written. 

421 """ 

422 base_events._check_ssl_socket(sock) 

423 if self._debug and sock.gettimeout() != 0: 

424 raise ValueError("the socket must be non-blocking") 

425 try: 

426 return sock.recv_into(buf) 

427 except (BlockingIOError, InterruptedError): 

428 pass 

429 fut = self.create_future() 

430 fd = sock.fileno() 

431 self._ensure_fd_no_transport(fd) 

432 handle = self._add_reader(fd, self._sock_recv_into, fut, sock, buf) 

433 fut.add_done_callback( 

434 functools.partial(self._sock_read_done, fd, handle=handle)) 

435 return await fut 

436 

437 def _sock_recv_into(self, fut, sock, buf): 

438 # _sock_recv_into() can add itself as an I/O callback if the operation 

439 # can't be done immediately. Don't use it directly, call 

440 # sock_recv_into(). 

441 if fut.done(): 441 ↛ 442line 441 didn't jump to line 442 because the condition on line 441 was never true

442 return 

443 try: 

444 nbytes = sock.recv_into(buf) 

445 except (BlockingIOError, InterruptedError): 

446 return # try again next time 

447 except (SystemExit, KeyboardInterrupt): 

448 raise 

449 except BaseException as exc: 

450 fut.set_exception(exc) 

451 else: 

452 fut.set_result(nbytes) 

453 

454 async def sock_recvfrom(self, sock, bufsize): 

455 """Receive a datagram from a datagram socket. 

456 

457 The return value is a tuple of (bytes, address) representing the 

458 datagram received and the address it came from. 

459 The maximum amount of data to be received at once is specified by 

460 nbytes. 

461 """ 

462 base_events._check_ssl_socket(sock) 

463 if self._debug and sock.gettimeout() != 0: 463 ↛ 464line 463 didn't jump to line 464 because the condition on line 463 was never true

464 raise ValueError("the socket must be non-blocking") 

465 try: 

466 return sock.recvfrom(bufsize) 

467 except (BlockingIOError, InterruptedError): 

468 pass 

469 fut = self.create_future() 

470 fd = sock.fileno() 

471 self._ensure_fd_no_transport(fd) 

472 handle = self._add_reader(fd, self._sock_recvfrom, fut, sock, bufsize) 

473 fut.add_done_callback( 

474 functools.partial(self._sock_read_done, fd, handle=handle)) 

475 return await fut 

476 

477 def _sock_recvfrom(self, fut, sock, bufsize): 

478 # _sock_recvfrom() can add itself as an I/O callback if the operation 

479 # can't be done immediately. Don't use it directly, call 

480 # sock_recvfrom(). 

481 if fut.done(): 481 ↛ 482line 481 didn't jump to line 482 because the condition on line 481 was never true

482 return 

483 try: 

484 result = sock.recvfrom(bufsize) 

485 except (BlockingIOError, InterruptedError): 

486 return # try again next time 

487 except (SystemExit, KeyboardInterrupt): 

488 raise 

489 except BaseException as exc: 

490 fut.set_exception(exc) 

491 else: 

492 fut.set_result(result) 

493 

494 async def sock_recvfrom_into(self, sock, buf, nbytes=0): 

495 """Receive data from the socket. 

496 

497 The received data is written into *buf* (a writable buffer). 

498 The return value is a tuple of (number of bytes written, address). 

499 """ 

500 base_events._check_ssl_socket(sock) 

501 if self._debug and sock.gettimeout() != 0: 501 ↛ 502line 501 didn't jump to line 502 because the condition on line 501 was never true

502 raise ValueError("the socket must be non-blocking") 

503 if not nbytes: 

504 nbytes = len(buf) 

505 

506 try: 

507 return sock.recvfrom_into(buf, nbytes) 

508 except (BlockingIOError, InterruptedError): 

509 pass 

510 fut = self.create_future() 

511 fd = sock.fileno() 

512 self._ensure_fd_no_transport(fd) 

513 handle = self._add_reader(fd, self._sock_recvfrom_into, fut, sock, buf, 

514 nbytes) 

515 fut.add_done_callback( 

516 functools.partial(self._sock_read_done, fd, handle=handle)) 

517 return await fut 

518 

519 def _sock_recvfrom_into(self, fut, sock, buf, bufsize): 

520 # _sock_recv_into() can add itself as an I/O callback if the operation 

521 # can't be done immediately. Don't use it directly, call 

522 # sock_recv_into(). 

523 if fut.done(): 523 ↛ 524line 523 didn't jump to line 524 because the condition on line 523 was never true

524 return 

525 try: 

526 result = sock.recvfrom_into(buf, bufsize) 

527 except (BlockingIOError, InterruptedError): 

528 return # try again next time 

529 except (SystemExit, KeyboardInterrupt): 

530 raise 

531 except BaseException as exc: 

532 fut.set_exception(exc) 

533 else: 

534 fut.set_result(result) 

535 

536 async def sock_sendall(self, sock, data): 

537 """Send data to the socket. 

538 

539 The socket must be connected to a remote socket. This method 

540 continues to send data from data until either all data has been 

541 sent or an error occurs. None is returned on success. On error, 

542 an exception is raised, and there is no way to determine how much 

543 data, if any, was successfully processed by the receiving end of 

544 the connection. 

545 """ 

546 base_events._check_ssl_socket(sock) 

547 if self._debug and sock.gettimeout() != 0: 

548 raise ValueError("the socket must be non-blocking") 

549 try: 

550 n = sock.send(data) 

551 except (BlockingIOError, InterruptedError): 

552 n = 0 

553 

554 if n == len(data): 

555 # all data sent 

556 return 

557 

558 fut = self.create_future() 

559 fd = sock.fileno() 

560 self._ensure_fd_no_transport(fd) 

561 # use a trick with a list in closure to store a mutable state 

562 handle = self._add_writer(fd, self._sock_sendall, fut, sock, 

563 memoryview(data), [n]) 

564 fut.add_done_callback( 

565 functools.partial(self._sock_write_done, fd, handle=handle)) 

566 return await fut 

567 

568 def _sock_sendall(self, fut, sock, view, pos): 

569 if fut.done(): 569 ↛ 571line 569 didn't jump to line 571 because the condition on line 569 was never true

570 # Future cancellation can be scheduled on previous loop iteration 

571 return 

572 start = pos[0] 

573 try: 

574 n = sock.send(view[start:]) 

575 except (BlockingIOError, InterruptedError): 

576 return 

577 except (SystemExit, KeyboardInterrupt): 

578 raise 

579 except BaseException as exc: 

580 fut.set_exception(exc) 

581 return 

582 

583 start += n 

584 

585 if start == len(view): 

586 fut.set_result(None) 

587 else: 

588 pos[0] = start 

589 

590 async def sock_sendto(self, sock, data, address): 

591 """Send a datagram from sock to address. 

592 

593 The socket does not have to be connected. This method sends the 

594 whole datagram in a single call. Return the number of bytes sent. 

595 """ 

596 base_events._check_ssl_socket(sock) 

597 if self._debug and sock.gettimeout() != 0: 597 ↛ 598line 597 didn't jump to line 598 because the condition on line 597 was never true

598 raise ValueError("the socket must be non-blocking") 

599 try: 

600 return sock.sendto(data, address) 

601 except (BlockingIOError, InterruptedError): 

602 pass 

603 

604 fut = self.create_future() 

605 fd = sock.fileno() 

606 self._ensure_fd_no_transport(fd) 

607 # use a trick with a list in closure to store a mutable state 

608 handle = self._add_writer(fd, self._sock_sendto, fut, sock, data, 

609 address) 

610 fut.add_done_callback( 

611 functools.partial(self._sock_write_done, fd, handle=handle)) 

612 return await fut 

613 

614 def _sock_sendto(self, fut, sock, data, address): 

615 if fut.done(): 615 ↛ 617line 615 didn't jump to line 617 because the condition on line 615 was never true

616 # Future cancellation can be scheduled on previous loop iteration 

617 return 

618 try: 

619 n = sock.sendto(data, 0, address) 

620 except (BlockingIOError, InterruptedError): 

621 return 

622 except (SystemExit, KeyboardInterrupt): 

623 raise 

624 except BaseException as exc: 

625 fut.set_exception(exc) 

626 else: 

627 fut.set_result(n) 

628 

629 async def sock_connect(self, sock, address): 

630 """Connect to a remote socket at address. 

631 

632 This method is a coroutine. 

633 """ 

634 base_events._check_ssl_socket(sock) 

635 if self._debug and sock.gettimeout() != 0: 

636 raise ValueError("the socket must be non-blocking") 

637 

638 if sock.family == socket.AF_INET or ( 

639 base_events._HAS_IPv6 and sock.family == socket.AF_INET6): 

640 resolved = await self._ensure_resolved( 

641 address, family=sock.family, type=sock.type, proto=sock.proto, 

642 loop=self, 

643 ) 

644 _, _, _, _, address = resolved[0] 

645 

646 fut = self.create_future() 

647 self._sock_connect(fut, sock, address) 

648 try: 

649 return await fut 

650 finally: 

651 # Needed to break cycles when an exception occurs. 

652 fut = None 

653 

654 def _sock_connect(self, fut, sock, address): 

655 fd = sock.fileno() 

656 try: 

657 sock.connect(address) 

658 except (BlockingIOError, InterruptedError): 

659 # Issue #23618: When the C function connect() fails with EINTR, the 

660 # connection runs in background. We have to wait until the socket 

661 # becomes writable to be notified when the connection succeed or 

662 # fails. 

663 self._ensure_fd_no_transport(fd) 

664 handle = self._add_writer( 

665 fd, self._sock_connect_cb, fut, sock, address) 

666 fut.add_done_callback( 

667 functools.partial(self._sock_write_done, fd, handle=handle)) 

668 except (SystemExit, KeyboardInterrupt): 

669 raise 

670 except BaseException as exc: 

671 fut.set_exception(exc) 

672 else: 

673 fut.set_result(None) 

674 finally: 

675 fut = None 

676 

677 def _sock_write_done(self, fd, fut, handle=None): 

678 if handle is None or not handle.cancelled(): 

679 self.remove_writer(fd) 

680 

681 def _sock_connect_cb(self, fut, sock, address): 

682 if fut.done(): 682 ↛ 683line 682 didn't jump to line 683 because the condition on line 682 was never true

683 return 

684 

685 try: 

686 err = sock.getsockopt(socket.SOL_SOCKET, socket.SO_ERROR) 

687 if err != 0: 

688 # Jump to any except clause below. 

689 raise OSError(err, f'Connect call failed {address}') 

690 except (BlockingIOError, InterruptedError): 

691 # socket is still registered, the callback will be retried later 

692 pass 

693 except (SystemExit, KeyboardInterrupt): 

694 raise 

695 except BaseException as exc: 

696 fut.set_exception(exc) 

697 else: 

698 fut.set_result(None) 

699 finally: 

700 fut = None 

701 

702 async def sock_accept(self, sock): 

703 """Accept a connection. 

704 

705 The socket must be bound to an address and listening for 

706 connections. The return value is a pair (conn, address) where 

707 conn is a new socket object usable to send and receive data on the 

708 connection, and address is the address bound to the socket on the 

709 other end of the connection. 

710 """ 

711 base_events._check_ssl_socket(sock) 

712 if self._debug and sock.gettimeout() != 0: 

713 raise ValueError("the socket must be non-blocking") 

714 fut = self.create_future() 

715 self._sock_accept(fut, sock) 

716 return await fut 

717 

718 def _sock_accept(self, fut, sock): 

719 # gh-153761: _sock_accept must not scheduled with already cancelled future 

720 if fut.done(): 

721 return 

722 fd = sock.fileno() 

723 try: 

724 conn, address = sock.accept() 

725 conn.setblocking(False) 

726 except (BlockingIOError, InterruptedError): 

727 self._ensure_fd_no_transport(fd) 

728 handle = self._add_reader(fd, self._sock_accept, fut, sock) 

729 fut.add_done_callback( 

730 functools.partial(self._sock_read_done, fd, handle=handle)) 

731 except (SystemExit, KeyboardInterrupt): 

732 raise 

733 except BaseException as exc: 

734 fut.set_exception(exc) 

735 else: 

736 fut.set_result((conn, address)) 

737 

738 async def _sendfile_native(self, transp, file, offset, count): 

739 del self._transports[transp._sock_fd] 

740 resume_reading = transp.is_reading() 

741 transp.pause_reading() 

742 try: 

743 await transp._make_empty_waiter() 

744 return await self.sock_sendfile(transp._sock, file, offset, count, 

745 fallback=False) 

746 finally: 

747 transp._reset_empty_waiter() 

748 if resume_reading: 

749 transp.resume_reading() 

750 self._transports[transp._sock_fd] = transp 

751 

752 def _process_events(self, event_list): 

753 for key, mask in event_list: 

754 fileobj, (reader, writer) = key.fileobj, key.data 

755 if mask & selectors.EVENT_READ and reader is not None: 

756 if reader._cancelled: 

757 self._remove_reader(fileobj) 

758 else: 

759 self._add_callback(reader) 

760 if mask & selectors.EVENT_WRITE and writer is not None: 

761 if writer._cancelled: 

762 self._remove_writer(fileobj) 

763 else: 

764 self._add_callback(writer) 

765 

766 def _stop_serving(self, sock): 

767 self._remove_reader(sock.fileno()) 

768 sock.close() 

769 

770 

771class _SelectorTransport(transports._FlowControlMixin, 

772 transports.Transport): 

773 

774 max_size = 256 * 1024 # Buffer size passed to recv(). 

775 

776 # Attribute used in the destructor: it must be set even if the constructor 

777 # is not called (see _SelectorSslTransport which may start by raising an 

778 # exception) 

779 _sock = None 

780 

781 def __init__(self, loop, sock, protocol, extra=None, server=None, context=None): 

782 super().__init__(extra, loop) 

783 self._extra['socket'] = trsock.TransportSocket(sock) 

784 try: 

785 self._extra['sockname'] = sock.getsockname() 

786 except OSError: 

787 self._extra['sockname'] = None 

788 if 'peername' not in self._extra: 

789 try: 

790 self._extra['peername'] = sock.getpeername() 

791 except socket.error: 

792 self._extra['peername'] = None 

793 self._sock = sock 

794 self._sock_fd = sock.fileno() 

795 self._context = context 

796 self._protocol_connected = False 

797 self.set_protocol(protocol) 

798 

799 self._server = server 

800 self._buffer = collections.deque() 

801 self._buffer_size = 0 

802 self._conn_lost = 0 # Set when call to connection_lost scheduled. 

803 self._closing = False # Set when close() called. 

804 self._paused = False # Set when pause_reading() called 

805 

806 if self._server is not None: 

807 self._server._attach(self) 

808 loop._transports[self._sock_fd] = self 

809 

810 def __repr__(self): 

811 info = [self.__class__.__name__] 

812 if self._sock is None: 812 ↛ 813line 812 didn't jump to line 813 because the condition on line 812 was never true

813 info.append('closed') 

814 elif self._closing: 

815 info.append('closing') 

816 info.append(f'fd={self._sock_fd}') 

817 # test if the transport was closed 

818 if self._loop is not None and not self._loop.is_closed(): 

819 polling = _test_selector_event(self._loop._selector, 

820 self._sock_fd, selectors.EVENT_READ) 

821 if polling: 

822 info.append('read=polling') 

823 else: 

824 info.append('read=idle') 

825 

826 polling = _test_selector_event(self._loop._selector, 

827 self._sock_fd, 

828 selectors.EVENT_WRITE) 

829 if polling: 829 ↛ 830line 829 didn't jump to line 830 because the condition on line 829 was never true

830 state = 'polling' 

831 else: 

832 state = 'idle' 

833 

834 bufsize = self.get_write_buffer_size() 

835 info.append(f'write=<{state}, bufsize={bufsize}>') 

836 return '<{}>'.format(' '.join(info)) 

837 

838 def abort(self): 

839 self._force_close(None) 

840 

841 def set_protocol(self, protocol): 

842 self._protocol = protocol 

843 self._protocol_connected = True 

844 

845 def get_protocol(self): 

846 return self._protocol 

847 

848 def is_closing(self): 

849 return self._closing 

850 

851 def is_reading(self): 

852 return not self.is_closing() and not self._paused 

853 

854 def pause_reading(self): 

855 if not self.is_reading(): 

856 return 

857 self._paused = True 

858 self._loop._remove_reader(self._sock_fd) 

859 if self._loop.get_debug(): 

860 logger.debug("%r pauses reading", self) 

861 

862 def resume_reading(self): 

863 if self._closing or not self._paused: 

864 return 

865 self._paused = False 

866 self._add_reader(self._sock_fd, self._read_ready) 

867 if self._loop.get_debug(): 867 ↛ 868line 867 didn't jump to line 868 because the condition on line 867 was never true

868 logger.debug("%r resumes reading", self) 

869 

870 def close(self): 

871 if self._closing: 

872 return 

873 self._closing = True 

874 self._loop._remove_reader(self._sock_fd) 

875 if not self._buffer: 

876 self._conn_lost += 1 

877 self._loop._remove_writer(self._sock_fd) 

878 self._call_soon(self._call_connection_lost, None) 

879 

880 def __del__(self, _warn=warnings.warn): 

881 if self._sock is not None: 881 ↛ 882line 881 didn't jump to line 882 because the condition on line 881 was never true

882 _warn(f"unclosed transport {self!r}", ResourceWarning, source=self) 

883 self._sock.close() 

884 if self._server is not None: 

885 self._server._detach(self) 

886 

887 def _fatal_error(self, exc, message='Fatal error on transport'): 

888 # Should be called from exception handler only. 

889 if isinstance(exc, OSError): 

890 if self._loop.get_debug(): 890 ↛ 891line 890 didn't jump to line 891 because the condition on line 890 was never true

891 logger.debug("%r: %s", self, message, exc_info=True) 

892 else: 

893 self._loop.call_exception_handler({ 

894 'message': message, 

895 'exception': exc, 

896 'transport': self, 

897 'protocol': self._protocol, 

898 }) 

899 self._force_close(exc) 

900 

901 def _force_close(self, exc): 

902 if self._conn_lost: 

903 return 

904 if self._buffer: 

905 self._buffer.clear() 

906 self._buffer_size = 0 

907 self._loop._remove_writer(self._sock_fd) 

908 if not self._closing: 

909 self._closing = True 

910 self._loop._remove_reader(self._sock_fd) 

911 self._conn_lost += 1 

912 self._call_soon(self._call_connection_lost, exc) 

913 

914 def _call_connection_lost(self, exc): 

915 try: 

916 if self._protocol_connected: 916 ↛ 919line 916 didn't jump to line 919 because the condition on line 916 was always true

917 self._protocol.connection_lost(exc) 

918 finally: 

919 self._sock.close() 

920 self._sock = None 

921 self._protocol = None 

922 self._loop = None 

923 server = self._server 

924 if server is not None: 

925 server._detach(self) 

926 self._server = None 

927 

928 def get_write_buffer_size(self): 

929 return self._buffer_size 

930 

931 def _add_reader(self, fd, callback, *args): 

932 if not self.is_reading(): 

933 return 

934 self._loop._add_reader(fd, callback, *args, context=self._context) 

935 

936 def _add_writer(self, fd, callback, *args): 

937 self._loop._add_writer(fd, callback, *args, context=self._context) 

938 

939 def _call_soon(self, callback, *args): 

940 self._loop.call_soon(callback, *args, context=self._context) 

941 

942class _SelectorSocketTransport(_SelectorTransport): 

943 

944 _start_tls_compatible = True 

945 _sendfile_compatible = constants._SendfileMode.TRY_NATIVE 

946 

947 def __init__(self, loop, sock, protocol, waiter=None, 

948 extra=None, server=None, context=None): 

949 self._read_ready_cb = None 

950 super().__init__(loop, sock, protocol, extra, server, context) 

951 self._eof = False 

952 self._empty_waiter = None 

953 if _HAS_SENDMSG: 953 ↛ 956line 953 didn't jump to line 956 because the condition on line 953 was always true

954 self._write_ready = self._write_sendmsg 

955 else: 

956 self._write_ready = self._write_send 

957 # Disable the Nagle algorithm -- small writes will be 

958 # sent without waiting for the TCP ACK. This generally 

959 # decreases the latency (in some cases significantly.) 

960 base_events._set_nodelay(self._sock) 

961 

962 self._call_soon(self._protocol.connection_made, self) 

963 # only start reading when connection_made() has been called 

964 self._call_soon(self._add_reader, self._sock_fd, self._read_ready) 

965 if waiter is not None: 

966 # only wake up the waiter when connection_made() has been called 

967 self._call_soon(futures._set_result_unless_cancelled, waiter, None) 

968 

969 def set_protocol(self, protocol): 

970 if isinstance(protocol, protocols.BufferedProtocol): 

971 self._read_ready_cb = self._read_ready__get_buffer 

972 else: 

973 self._read_ready_cb = self._read_ready__data_received 

974 

975 super().set_protocol(protocol) 

976 

977 def _read_ready(self): 

978 self._read_ready_cb() 

979 

980 def _read_ready__get_buffer(self): 

981 if self._conn_lost: 981 ↛ 982line 981 didn't jump to line 982 because the condition on line 981 was never true

982 return 

983 

984 try: 

985 buf = self._protocol.get_buffer(-1) 

986 if not len(buf): 

987 raise RuntimeError('get_buffer() returned an empty buffer') 

988 except (SystemExit, KeyboardInterrupt): 

989 raise 

990 except BaseException as exc: 

991 self._fatal_error( 

992 exc, 'Fatal error: protocol.get_buffer() call failed.') 

993 return 

994 

995 try: 

996 nbytes = self._sock.recv_into(buf) 

997 except (BlockingIOError, InterruptedError): 

998 return 

999 except (SystemExit, KeyboardInterrupt): 

1000 raise 

1001 except BaseException as exc: 

1002 self._fatal_error(exc, 'Fatal read error on socket transport') 

1003 return 

1004 

1005 if not nbytes: 

1006 self._read_ready__on_eof() 

1007 return 

1008 

1009 try: 

1010 self._protocol.buffer_updated(nbytes) 

1011 except (SystemExit, KeyboardInterrupt): 

1012 raise 

1013 except BaseException as exc: 

1014 self._fatal_error( 

1015 exc, 'Fatal error: protocol.buffer_updated() call failed.') 

1016 

1017 def _read_ready__data_received(self): 

1018 if self._conn_lost: 1018 ↛ 1019line 1018 didn't jump to line 1019 because the condition on line 1018 was never true

1019 return 

1020 try: 

1021 data = self._sock.recv(self.max_size) 

1022 except (BlockingIOError, InterruptedError): 

1023 return 

1024 except (SystemExit, KeyboardInterrupt): 

1025 raise 

1026 except BaseException as exc: 

1027 self._fatal_error(exc, 'Fatal read error on socket transport') 

1028 return 

1029 

1030 if not data: 

1031 self._read_ready__on_eof() 

1032 return 

1033 

1034 try: 

1035 self._protocol.data_received(data) 

1036 except (SystemExit, KeyboardInterrupt): 

1037 raise 

1038 except BaseException as exc: 

1039 self._fatal_error( 

1040 exc, 'Fatal error: protocol.data_received() call failed.') 

1041 

1042 def _read_ready__on_eof(self): 

1043 if self._loop.get_debug(): 

1044 logger.debug("%r received EOF", self) 

1045 

1046 try: 

1047 keep_open = self._protocol.eof_received() 

1048 except (SystemExit, KeyboardInterrupt): 

1049 raise 

1050 except BaseException as exc: 

1051 self._fatal_error( 

1052 exc, 'Fatal error: protocol.eof_received() call failed.') 

1053 return 

1054 

1055 if keep_open: 

1056 # We're keeping the connection open so the 

1057 # protocol can write more, but we still can't 

1058 # receive more, so remove the reader callback. 

1059 self._loop._remove_reader(self._sock_fd) 

1060 else: 

1061 self.close() 

1062 

1063 def write(self, data): 

1064 if not isinstance(data, (bytes, bytearray, memoryview)): 

1065 raise TypeError(f'data argument must be a bytes, bytearray, or memoryview ' 

1066 f'object, not {type(data).__name__!r}') 

1067 if self._eof: 1067 ↛ 1068line 1067 didn't jump to line 1068 because the condition on line 1067 was never true

1068 raise RuntimeError('Cannot call write() after write_eof()') 

1069 if self._empty_waiter is not None: 

1070 raise RuntimeError('unable to write; sendfile is in progress') 

1071 if not data: 

1072 return 

1073 

1074 if self._conn_lost: 

1075 if self._conn_lost >= constants.LOG_THRESHOLD_FOR_CONNLOST_WRITES: 

1076 logger.warning('socket.send() raised exception.') 

1077 self._conn_lost += 1 

1078 return 

1079 

1080 if not self._buffer: 

1081 # Optimization: try to send now. 

1082 try: 

1083 n = self._sock.send(data) 

1084 except (BlockingIOError, InterruptedError): 

1085 pass 

1086 except (SystemExit, KeyboardInterrupt): 

1087 raise 

1088 except BaseException as exc: 

1089 self._fatal_error(exc, 'Fatal write error on socket transport') 

1090 return 

1091 else: 

1092 data = memoryview(data)[n:] 

1093 if not data: 

1094 return 

1095 # Not all was written; register write handler. 

1096 self._add_writer(self._sock_fd, self._write_ready) 

1097 

1098 # Add it to the buffer. 

1099 self._buffer.append(data) 

1100 self._buffer_size += len(data) 

1101 self._maybe_pause_protocol() 

1102 

1103 def _get_sendmsg_buffer(self): 

1104 return itertools.islice(self._buffer, SC_IOV_MAX) 

1105 

1106 def _write_sendmsg(self): 

1107 assert self._buffer, 'Data should not be empty' 

1108 if self._conn_lost: 1108 ↛ 1109line 1108 didn't jump to line 1109 because the condition on line 1108 was never true

1109 return 

1110 try: 

1111 nbytes = self._sock.sendmsg(self._get_sendmsg_buffer()) 

1112 self._adjust_leftover_buffer(nbytes) 

1113 except (BlockingIOError, InterruptedError): 

1114 pass 

1115 except (SystemExit, KeyboardInterrupt): 

1116 raise 

1117 except BaseException as exc: 

1118 self._loop._remove_writer(self._sock_fd) 

1119 self._buffer.clear() 

1120 self._buffer_size = 0 

1121 self._fatal_error(exc, 'Fatal write error on socket transport') 

1122 if self._empty_waiter is not None: 1122 ↛ 1123line 1122 didn't jump to line 1123 because the condition on line 1122 was never true

1123 self._empty_waiter.set_exception(exc) 

1124 else: 

1125 self._maybe_resume_protocol() # May append to buffer. 

1126 if not self._buffer: 

1127 self._loop._remove_writer(self._sock_fd) 

1128 if self._empty_waiter is not None: 

1129 self._empty_waiter.set_result(None) 

1130 # gh-156512: don't let _call_connection_lost be called twice 

1131 if self._closing and not self._conn_lost: 

1132 self._conn_lost += 1 

1133 self._call_connection_lost(None) 

1134 elif self._eof: 1134 ↛ 1135line 1134 didn't jump to line 1135 because the condition on line 1134 was never true

1135 self._sock.shutdown(socket.SHUT_WR) 

1136 

1137 def _adjust_leftover_buffer(self, nbytes: int) -> None: 

1138 self._buffer_size -= nbytes 

1139 buffer = self._buffer 

1140 while nbytes: 

1141 b = buffer.popleft() 

1142 b_len = len(b) 

1143 if b_len <= nbytes: 

1144 nbytes -= b_len 

1145 else: 

1146 buffer.appendleft(b[nbytes:]) 

1147 break 

1148 

1149 def _write_send(self): 

1150 assert self._buffer, 'Data should not be empty' 

1151 if self._conn_lost: 1151 ↛ 1152line 1151 didn't jump to line 1152 because the condition on line 1151 was never true

1152 return 

1153 try: 

1154 buffer = self._buffer.popleft() 

1155 n = self._sock.send(buffer) 

1156 if n != len(buffer): 

1157 # Not all data was written 

1158 self._buffer.appendleft(buffer[n:]) 

1159 self._buffer_size -= n 

1160 except (BlockingIOError, InterruptedError): 

1161 self._buffer.appendleft(buffer) 

1162 return 

1163 except (SystemExit, KeyboardInterrupt): 

1164 raise 

1165 except BaseException as exc: 

1166 self._loop._remove_writer(self._sock_fd) 

1167 self._buffer.clear() 

1168 self._buffer_size = 0 

1169 self._fatal_error(exc, 'Fatal write error on socket transport') 

1170 if self._empty_waiter is not None: 1170 ↛ 1171line 1170 didn't jump to line 1171 because the condition on line 1170 was never true

1171 self._empty_waiter.set_exception(exc) 

1172 else: 

1173 self._maybe_resume_protocol() # May append to buffer. 

1174 if not self._buffer: 

1175 self._loop._remove_writer(self._sock_fd) 

1176 if self._empty_waiter is not None: 1176 ↛ 1177line 1176 didn't jump to line 1177 because the condition on line 1176 was never true

1177 self._empty_waiter.set_result(None) 

1178 # gh-156512: don't let _call_connection_lost be called twice 

1179 if self._closing and not self._conn_lost: 

1180 self._conn_lost += 1 

1181 self._call_connection_lost(None) 

1182 elif self._eof: 

1183 self._sock.shutdown(socket.SHUT_WR) 

1184 

1185 def write_eof(self): 

1186 if self._closing or self._eof: 

1187 return 

1188 self._eof = True 

1189 if not self._buffer: 

1190 self._sock.shutdown(socket.SHUT_WR) 

1191 

1192 def writelines(self, list_of_data): 

1193 if self._eof: 1193 ↛ 1194line 1193 didn't jump to line 1194 because the condition on line 1193 was never true

1194 raise RuntimeError('Cannot call writelines() after write_eof()') 

1195 if self._empty_waiter is not None: 1195 ↛ 1196line 1195 didn't jump to line 1196 because the condition on line 1195 was never true

1196 raise RuntimeError('unable to writelines; sendfile is in progress') 

1197 if not list_of_data: 1197 ↛ 1198line 1197 didn't jump to line 1198 because the condition on line 1197 was never true

1198 return 

1199 

1200 if self._conn_lost: 

1201 if self._conn_lost >= constants.LOG_THRESHOLD_FOR_CONNLOST_WRITES: 1201 ↛ 1202line 1201 didn't jump to line 1202 because the condition on line 1201 was never true

1202 logger.warning('socket.send() raised exception.') 

1203 self._conn_lost += 1 

1204 return 

1205 

1206 for data in list_of_data: 

1207 # gh-155888: an empty chunk can never be drained, so never buffer it 

1208 if not data: 

1209 continue 

1210 self._buffer.append(memoryview(data)) 

1211 self._buffer_size += len(data) 

1212 if not self._buffer: 1212 ↛ 1213line 1212 didn't jump to line 1213 because the condition on line 1212 was never true

1213 return 

1214 self._write_ready() 

1215 # If the entire buffer couldn't be written, register a write handler 

1216 if self._buffer: 

1217 self._add_writer(self._sock_fd, self._write_ready) 

1218 self._maybe_pause_protocol() 

1219 

1220 def can_write_eof(self): 

1221 return True 

1222 

1223 def _call_connection_lost(self, exc): 

1224 try: 

1225 super()._call_connection_lost(exc) 

1226 finally: 

1227 self._write_ready = None 

1228 if self._empty_waiter is not None: 1228 ↛ 1229line 1228 didn't jump to line 1229 because the condition on line 1228 was never true

1229 self._empty_waiter.set_exception( 

1230 ConnectionError("Connection is closed by peer")) 

1231 

1232 def _make_empty_waiter(self): 

1233 if self._empty_waiter is not None: 1233 ↛ 1234line 1233 didn't jump to line 1234 because the condition on line 1233 was never true

1234 raise RuntimeError("Empty waiter is already set") 

1235 self._empty_waiter = self._loop.create_future() 

1236 if not self._buffer: 

1237 self._empty_waiter.set_result(None) 

1238 return self._empty_waiter 

1239 

1240 def _reset_empty_waiter(self): 

1241 self._empty_waiter = None 

1242 

1243 def close(self): 

1244 self._read_ready_cb = None 

1245 super().close() 

1246 

1247 

1248class _SelectorDatagramTransport(_SelectorTransport, transports.DatagramTransport): 

1249 

1250 _header_size = 8 

1251 

1252 def __init__(self, loop, sock, protocol, address=None, 

1253 waiter=None, extra=None): 

1254 super().__init__(loop, sock, protocol, extra) 

1255 self._address = address 

1256 self._buffer_size = 0 

1257 self._call_soon(self._protocol.connection_made, self) 

1258 # only start reading when connection_made() has been called 

1259 self._call_soon(self._add_reader, self._sock_fd, self._read_ready) 

1260 if waiter is not None: 

1261 # only wake up the waiter when connection_made() has been called 

1262 self._call_soon(futures._set_result_unless_cancelled, waiter, None) 

1263 

1264 def get_write_buffer_size(self): 

1265 return self._buffer_size 

1266 

1267 def _read_ready(self): 

1268 if self._conn_lost: 1268 ↛ 1269line 1268 didn't jump to line 1269 because the condition on line 1268 was never true

1269 return 

1270 try: 

1271 data, addr = self._sock.recvfrom(self.max_size) 

1272 except (BlockingIOError, InterruptedError): 

1273 pass 

1274 except OSError as exc: 

1275 self._protocol.error_received(exc) 

1276 except (SystemExit, KeyboardInterrupt): 

1277 raise 

1278 except BaseException as exc: 

1279 self._fatal_error(exc, 'Fatal read error on datagram transport') 

1280 else: 

1281 self._protocol.datagram_received(data, addr) 

1282 

1283 def sendto(self, data, addr=None): 

1284 if not isinstance(data, (bytes, bytearray, memoryview)): 

1285 raise TypeError(f'data argument must be a bytes-like object, ' 

1286 f'not {type(data).__name__!r}') 

1287 

1288 if self._address: 

1289 if addr not in (None, self._address): 

1290 raise ValueError( 

1291 f'Invalid address: must be None or {self._address}') 

1292 addr = self._address 

1293 

1294 if self._conn_lost and self._address: 

1295 if self._conn_lost >= constants.LOG_THRESHOLD_FOR_CONNLOST_WRITES: 

1296 logger.warning('socket.send() raised exception.') 

1297 self._conn_lost += 1 

1298 return 

1299 

1300 if not self._buffer: 

1301 # Attempt to send it right away first. 

1302 try: 

1303 if self._extra['peername']: 

1304 self._sock.send(data) 

1305 else: 

1306 self._sock.sendto(data, addr) 

1307 return 

1308 except (BlockingIOError, InterruptedError): 

1309 self._add_writer(self._sock_fd, self._sendto_ready) 

1310 except OSError as exc: 

1311 self._protocol.error_received(exc) 

1312 return 

1313 except (SystemExit, KeyboardInterrupt): 

1314 raise 

1315 except BaseException as exc: 

1316 self._fatal_error( 

1317 exc, 'Fatal write error on datagram transport') 

1318 return 

1319 

1320 # Ensure that what we buffer is immutable. 

1321 self._buffer.append((bytes(data), addr)) 

1322 self._buffer_size += len(data) + self._header_size 

1323 self._maybe_pause_protocol() 

1324 

1325 def _sendto_ready(self): 

1326 while self._buffer: 

1327 data, addr = self._buffer.popleft() 

1328 self._buffer_size -= len(data) + self._header_size 

1329 try: 

1330 if self._extra['peername']: 

1331 self._sock.send(data) 

1332 else: 

1333 self._sock.sendto(data, addr) 

1334 except (BlockingIOError, InterruptedError): 

1335 self._buffer.appendleft((data, addr)) # Try again later. 

1336 self._buffer_size += len(data) + self._header_size 

1337 break 

1338 except OSError as exc: 

1339 self._protocol.error_received(exc) 

1340 return 

1341 except (SystemExit, KeyboardInterrupt): 

1342 raise 

1343 except BaseException as exc: 

1344 self._fatal_error( 

1345 exc, 'Fatal write error on datagram transport') 

1346 return 

1347 

1348 self._maybe_resume_protocol() # May append to buffer. 

1349 if not self._buffer: 

1350 self._loop._remove_writer(self._sock_fd) 

1351 if self._closing: 

1352 self._call_connection_lost(None)