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
« 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.
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"""
7__all__ = 'BaseSelectorEventLoop',
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
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
33_HAS_SENDMSG = hasattr(socket.socket, 'sendmsg')
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
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)
53class BaseSelectorEventLoop(base_events.BaseEventLoop):
54 """Selector event loop.
56 See events.AbstractEventLoop for API specification.
57 """
59 def __init__(self, selector=None):
60 super().__init__()
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()
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)
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
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)
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
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
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)
127 def _process_self_data(self, data):
128 pass
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
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
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)
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)
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)
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)
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.
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)
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}')
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
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))
312 if reader is not None:
313 reader.cancel()
314 return True
315 else:
316 return False
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
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))
348 if writer is not None:
349 writer.cancel()
350 return True
351 else:
352 return False
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)
359 def remove_reader(self, fd):
360 """Remove a reader callback."""
361 self._ensure_fd_no_transport(fd)
362 return self._remove_reader(fd)
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)
369 def remove_writer(self, fd):
370 """Remove a writer callback."""
371 self._ensure_fd_no_transport(fd)
372 return self._remove_writer(fd)
374 async def sock_recv(self, sock, n):
375 """Receive data from the socket.
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
396 def _sock_read_done(self, fd, fut, handle=None):
397 if handle is None or not handle.cancelled():
398 self.remove_reader(fd)
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)
416 async def sock_recv_into(self, sock, buf):
417 """Receive data from the socket.
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
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)
454 async def sock_recvfrom(self, sock, bufsize):
455 """Receive a datagram from a datagram socket.
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
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)
494 async def sock_recvfrom_into(self, sock, buf, nbytes=0):
495 """Receive data from the socket.
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)
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
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)
536 async def sock_sendall(self, sock, data):
537 """Send data to the socket.
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
554 if n == len(data):
555 # all data sent
556 return
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
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
583 start += n
585 if start == len(view):
586 fut.set_result(None)
587 else:
588 pos[0] = start
590 async def sock_sendto(self, sock, data, address):
591 """Send a datagram from sock to address.
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
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
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)
629 async def sock_connect(self, sock, address):
630 """Connect to a remote socket at address.
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")
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]
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
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
677 def _sock_write_done(self, fd, fut, handle=None):
678 if handle is None or not handle.cancelled():
679 self.remove_writer(fd)
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
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
702 async def sock_accept(self, sock):
703 """Accept a connection.
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
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))
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
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)
766 def _stop_serving(self, sock):
767 self._remove_reader(sock.fileno())
768 sock.close()
771class _SelectorTransport(transports._FlowControlMixin,
772 transports.Transport):
774 max_size = 256 * 1024 # Buffer size passed to recv().
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
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)
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
806 if self._server is not None:
807 self._server._attach(self)
808 loop._transports[self._sock_fd] = self
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')
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'
834 bufsize = self.get_write_buffer_size()
835 info.append(f'write=<{state}, bufsize={bufsize}>')
836 return '<{}>'.format(' '.join(info))
838 def abort(self):
839 self._force_close(None)
841 def set_protocol(self, protocol):
842 self._protocol = protocol
843 self._protocol_connected = True
845 def get_protocol(self):
846 return self._protocol
848 def is_closing(self):
849 return self._closing
851 def is_reading(self):
852 return not self.is_closing() and not self._paused
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)
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)
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)
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)
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)
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)
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
928 def get_write_buffer_size(self):
929 return self._buffer_size
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)
936 def _add_writer(self, fd, callback, *args):
937 self._loop._add_writer(fd, callback, *args, context=self._context)
939 def _call_soon(self, callback, *args):
940 self._loop.call_soon(callback, *args, context=self._context)
942class _SelectorSocketTransport(_SelectorTransport):
944 _start_tls_compatible = True
945 _sendfile_compatible = constants._SendfileMode.TRY_NATIVE
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)
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)
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
975 super().set_protocol(protocol)
977 def _read_ready(self):
978 self._read_ready_cb()
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
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
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
1005 if not nbytes:
1006 self._read_ready__on_eof()
1007 return
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.')
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
1030 if not data:
1031 self._read_ready__on_eof()
1032 return
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.')
1042 def _read_ready__on_eof(self):
1043 if self._loop.get_debug():
1044 logger.debug("%r received EOF", self)
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
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()
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
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
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)
1098 # Add it to the buffer.
1099 self._buffer.append(data)
1100 self._buffer_size += len(data)
1101 self._maybe_pause_protocol()
1103 def _get_sendmsg_buffer(self):
1104 return itertools.islice(self._buffer, SC_IOV_MAX)
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)
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
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)
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)
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
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
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()
1220 def can_write_eof(self):
1221 return True
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"))
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
1240 def _reset_empty_waiter(self):
1241 self._empty_waiter = None
1243 def close(self):
1244 self._read_ready_cb = None
1245 super().close()
1248class _SelectorDatagramTransport(_SelectorTransport, transports.DatagramTransport):
1250 _header_size = 8
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)
1264 def get_write_buffer_size(self):
1265 return self._buffer_size
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)
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}')
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
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
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
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()
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
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)