codekingpro/portable-devtools
114k
1"""Base implementation of event loop.2 3The event loop can be broken up into a multiplexer (the part4responsible for notifying us of I/O events) and the event loop proper,5which wraps a multiplexer with functionality for scheduling callbacks,6immediately or at a given time in the future.7 8Whenever a public API takes a callback, subsequent positional9arguments will be passed to the callback if/when it is called. This10avoids the proliferation of trivial lambdas implementing closures.11Keyword arguments for the callback are not supported; this is a12conscious design decision, leaving the door open for keyword arguments13to modify the meaning of the API call itself.14"""15 16import collections17import collections.abc18import concurrent.futures19import errno20import heapq21import itertools22import os23import socket24import stat25import subprocess26import threading27import time28import traceback29import sys30import warnings31import weakref32 33try:34 import ssl35except ImportError: # pragma: no cover36 ssl = None37 38from . import constants39from . import coroutines40from . import events41from . import exceptions42from . import futures43from . import protocols44from . import sslproto45from . import staggered46from . import tasks47from . import timeouts48from . import transports49from . import trsock50from .log import logger51 52 53__all__ = 'BaseEventLoop','Server',54 55 56# Minimum number of _scheduled timer handles before cleanup of57# cancelled handles is performed.58_MIN_SCHEDULED_TIMER_HANDLES = 10059 60# Minimum fraction of _scheduled timer handles that are cancelled61# before cleanup of cancelled handles is performed.62_MIN_CANCELLED_TIMER_HANDLES_FRACTION = 0.563 64 65_HAS_IPv6 = hasattr(socket, 'AF_INET6')66 67# Maximum timeout passed to select to avoid OS limitations68MAXIMUM_SELECT_TIMEOUT = 24 * 360069 70 71def _format_handle(handle):72 cb = handle._callback73 if isinstance(getattr(cb, '__self__', None), tasks.Task):74 # format the task75 return repr(cb.__self__)76 else:77 return str(handle)78 79 80def _format_pipe(fd):81 if fd == subprocess.PIPE:82 return '<pipe>'83 elif fd == subprocess.STDOUT:84 return '<stdout>'85 else:86 return repr(fd)87 88 89def _set_reuseport(sock):90 if not hasattr(socket, 'SO_REUSEPORT'):91 raise ValueError('reuse_port not supported by socket module')92 else:93 try:94 sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, 1)95 except OSError:96 raise ValueError('reuse_port not supported by socket module, '97 'SO_REUSEPORT defined but not implemented.')98 99 100def _ipaddr_info(host, port, family, type, proto, flowinfo=0, scopeid=0):101 # Try to skip getaddrinfo if "host" is already an IP. Users might have102 # handled name resolution in their own code and pass in resolved IPs.103 if not hasattr(socket, 'inet_pton'):104 return105 106 if proto not in {0, socket.IPPROTO_TCP, socket.IPPROTO_UDP} or \107 host is None:108 return None109 110 if type == socket.SOCK_STREAM:111 proto = socket.IPPROTO_TCP112 elif type == socket.SOCK_DGRAM:113 proto = socket.IPPROTO_UDP114 else:115 return None116 117 if port is None:118 port = 0119 elif isinstance(port, bytes) and port == b'':120 port = 0121 elif isinstance(port, str) and port == '':122 port = 0123 else:124 # If port's a service name like "http", don't skip getaddrinfo.125 try:126 port = int(port)127 except (TypeError, ValueError):128 return None129 130 if family == socket.AF_UNSPEC:131 afs = [socket.AF_INET]132 if _HAS_IPv6:133 afs.append(socket.AF_INET6)134 else:135 afs = [family]136 137 if isinstance(host, bytes):138 host = host.decode('idna')139 if '%' in host:140 # Linux's inet_pton doesn't accept an IPv6 zone index after host,141 # like '::1%lo0'.142 return None143 144 for af in afs:145 try:146 socket.inet_pton(af, host)147 # The host has already been resolved.148 if _HAS_IPv6 and af == socket.AF_INET6:149 return af, type, proto, '', (host, port, flowinfo, scopeid)150 else:151 return af, type, proto, '', (host, port)152 except OSError:153 pass154 155 # "host" is not an IP address.156 return None157 158 159def _interleave_addrinfos(addrinfos, first_address_family_count=1):160 """Interleave list of addrinfo tuples by family."""161 # Group addresses by family162 addrinfos_by_family = collections.OrderedDict()163 for addr in addrinfos:164 family = addr[0]165 if family not in addrinfos_by_family:166 addrinfos_by_family[family] = []167 addrinfos_by_family[family].append(addr)168 addrinfos_lists = list(addrinfos_by_family.values())169 170 reordered = []171 if first_address_family_count > 1:172 reordered.extend(addrinfos_lists[0][:first_address_family_count - 1])173 del addrinfos_lists[0][:first_address_family_count - 1]174 reordered.extend(175 a for a in itertools.chain.from_iterable(176 itertools.zip_longest(*addrinfos_lists)177 ) if a is not None)178 return reordered179 180 181def _run_until_complete_cb(fut):182 if not fut.cancelled():183 exc = fut.exception()184 if isinstance(exc, (SystemExit, KeyboardInterrupt)):185 # Issue #22429: run_forever() already finished, no need to186 # stop it.187 return188 futures._get_loop(fut).stop()189 190 191if hasattr(socket, 'TCP_NODELAY'):192 def _set_nodelay(sock):193 if (sock.family in {socket.AF_INET, socket.AF_INET6} and194 sock.type == socket.SOCK_STREAM and195 sock.proto == socket.IPPROTO_TCP):196 sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)197else:198 def _set_nodelay(sock):199 pass200 201 202def _check_ssl_socket(sock):203 if ssl is not None and isinstance(sock, ssl.SSLSocket):204 raise TypeError("Socket cannot be of type SSLSocket")205 206 207class _SendfileFallbackProtocol(protocols.Protocol):208 def __init__(self, transp):209 if not isinstance(transp, transports._FlowControlMixin):210 raise TypeError("transport should be _FlowControlMixin instance")211 self._transport = transp212 self._proto = transp.get_protocol()213 self._should_resume_reading = transp.is_reading()214 self._should_resume_writing = transp._protocol_paused215 transp.pause_reading()216 transp.set_protocol(self)217 if self._should_resume_writing:218 self._write_ready_fut = self._transport._loop.create_future()219 else:220 self._write_ready_fut = None221 222 async def drain(self):223 if self._transport.is_closing():224 raise ConnectionError("Connection closed by peer")225 fut = self._write_ready_fut226 if fut is None:227 return228 await fut229 230 def connection_made(self, transport):231 raise RuntimeError("Invalid state: "232 "connection should have been established already.")233 234 def connection_lost(self, exc):235 if self._write_ready_fut is not None:236 # Never happens if peer disconnects after sending the whole content237 # Thus disconnection is always an exception from user perspective238 if exc is None:239 self._write_ready_fut.set_exception(240 ConnectionError("Connection is closed by peer"))241 else:242 self._write_ready_fut.set_exception(exc)243 self._proto.connection_lost(exc)244 245 def pause_writing(self):246 if self._write_ready_fut is not None:247 return248 self._write_ready_fut = self._transport._loop.create_future()249 250 def resume_writing(self):251 if self._write_ready_fut is None:252 return253 self._write_ready_fut.set_result(False)254 self._write_ready_fut = None255 256 def data_received(self, data):257 raise RuntimeError("Invalid state: reading should be paused")258 259 def eof_received(self):260 raise RuntimeError("Invalid state: reading should be paused")261 262 async def restore(self):263 self._transport.set_protocol(self._proto)264 if self._should_resume_reading:265 self._transport.resume_reading()266 if self._write_ready_fut is not None:267 # Cancel the future.268 # Basically it has no effect because protocol is switched back,269 # no code should wait for it anymore.270 self._write_ready_fut.cancel()271 if self._should_resume_writing:272 self._proto.resume_writing()273 274 275class Server(events.AbstractServer):276 277 def __init__(self, loop, sockets, protocol_factory, ssl_context, backlog,278 ssl_handshake_timeout, ssl_shutdown_timeout=None):279 self._loop = loop280 self._sockets = sockets281 # Weak references so we don't break Transport's ability to282 # detect abandoned transports283 self._clients = weakref.WeakSet()284 self._waiters = []285 self._protocol_factory = protocol_factory286 self._backlog = backlog287 self._ssl_context = ssl_context288 self._ssl_handshake_timeout = ssl_handshake_timeout289 self._ssl_shutdown_timeout = ssl_shutdown_timeout290 self._serving = False291 self._serving_forever_fut = None292 293 def __repr__(self):294 return f'<{self.__class__.__name__} sockets={self.sockets!r}>'295 296 def _attach(self, transport):297 assert self._sockets is not None298 self._clients.add(transport)299 300 def _detach(self, transport):301 self._clients.discard(transport)302 if len(self._clients) == 0 and self._sockets is None:303 self._wakeup()304 305 def _wakeup(self):306 waiters = self._waiters307 self._waiters = None308 for waiter in waiters:309 if not waiter.done():310 waiter.set_result(None)311 312 def _start_serving(self):313 if self._serving:314 return315 self._serving = True316 for sock in self._sockets:317 sock.listen(self._backlog)318 self._loop._start_serving(319 self._protocol_factory, sock, self._ssl_context,320 self, self._backlog, self._ssl_handshake_timeout,321 self._ssl_shutdown_timeout)322 323 def get_loop(self):324 return self._loop325 326 def is_serving(self):327 return self._serving328 329 @property330 def sockets(self):331 if self._sockets is None:332 return ()333 return tuple(trsock.TransportSocket(s) for s in self._sockets)334 335 def close(self):336 sockets = self._sockets337 if sockets is None:338 return339 self._sockets = None340 341 for sock in sockets:342 self._loop._stop_serving(sock)343 344 self._serving = False345 346 if (self._serving_forever_fut is not None and347 not self._serving_forever_fut.done()):348 self._serving_forever_fut.cancel()349 self._serving_forever_fut = None350 351 if len(self._clients) == 0:352 self._wakeup()353 354 def close_clients(self):355 for transport in self._clients.copy():356 transport.close()357 358 def abort_clients(self):359 for transport in self._clients.copy():360 transport.abort()361 362 async def start_serving(self):363 self._start_serving()364 # Skip one loop iteration so that all 'loop.add_reader'365 # go through.366 await tasks.sleep(0)367 368 async def serve_forever(self):369 if self._serving_forever_fut is not None:370 raise RuntimeError(371 f'server {self!r} is already being awaited on serve_forever()')372 if self._sockets is None:373 raise RuntimeError(f'server {self!r} is closed')374 375 self._start_serving()376 self._serving_forever_fut = self._loop.create_future()377 378 try:379 await self._serving_forever_fut380 except exceptions.CancelledError:381 try:382 self.close()383 await self.wait_closed()384 finally:385 raise386 finally:387 self._serving_forever_fut = None388 389 async def wait_closed(self):390 """Wait until server is closed and all connections are dropped.391 392 - If the server is not closed, wait.393 - If it is closed, but there are still active connections, wait.394 395 Anyone waiting here will be unblocked once both conditions396 (server is closed and all connections have been dropped)397 have become true, in either order.398 399 Historical note: In 3.11 and before, this was broken, returning400 immediately if the server was already closed, even if there401 were still active connections. An attempted fix in 3.12.0 was402 still broken, returning immediately if the server was still403 open and there were no active connections. Hopefully in 3.12.1404 we have it right.405 """406 # Waiters are unblocked by self._wakeup(), which is called407 # from two places: self.close() and self._detach(), but only408 # when both conditions have become true. To signal that this409 # has happened, self._wakeup() sets self._waiters to None.410 if self._waiters is None:411 return412 waiter = self._loop.create_future()413 self._waiters.append(waiter)414 await waiter415 416 417class BaseEventLoop(events.AbstractEventLoop):418 419 def __init__(self):420 self._timer_cancelled_count = 0421 self._closed = False422 self._stopping = False423 self._ready = collections.deque()424 self._scheduled = []425 self._default_executor = None426 self._internal_fds = 0427 # Identifier of the thread running the event loop, or None if the428 # event loop is not running429 self._thread_id = None430 self._clock_resolution = time.get_clock_info('monotonic').resolution431 self._exception_handler = None432 self.set_debug(coroutines._is_debug_mode())433 # The preserved state of async generator hooks.434 self._old_agen_hooks = None435 # In debug mode, if the execution of a callback or a step of a task436 # exceed this duration in seconds, the slow callback/task is logged.437 self.slow_callback_duration = 0.1438 self._current_handle = None439 self._task_factory = None440 self._coroutine_origin_tracking_enabled = False441 self._coroutine_origin_tracking_saved_depth = None442 443 # A weak set of all asynchronous generators that are444 # being iterated by the loop.445 self._asyncgens = weakref.WeakSet()446 # Set to True when `loop.shutdown_asyncgens` is called.447 self._asyncgens_shutdown_called = False448 # Set to True when `loop.shutdown_default_executor` is called.449 self._executor_shutdown_called = False450 451 def __repr__(self):452 return (453 f'<{self.__class__.__name__} running={self.is_running()} '454 f'closed={self.is_closed()} debug={self.get_debug()}>'455 )456 457 def create_future(self):458 """Create a Future object attached to the loop."""459 return futures.Future(loop=self)460 461 def create_task(self, coro, **kwargs):462 """Schedule or begin executing a coroutine object.463 464 Return a task object.465 """466 self._check_closed()467 if self._task_factory is not None:468 return self._task_factory(self, coro, **kwargs)469 470 task = tasks.Task(coro, loop=self, **kwargs)471 if task._source_traceback:472 del task._source_traceback[-1]473 try:474 return task475 finally:476 # gh-128552: prevent a refcycle of477 # task.exception().__traceback__->BaseEventLoop.create_task->task478 del task479 480 def set_task_factory(self, factory):481 """Set a task factory that will be used by loop.create_task().482 483 If factory is None the default task factory will be set.484 485 If factory is a callable, it should have a signature matching486 '(loop, coro, **kwargs)', where 'loop' will be a reference to the active487 event loop, 'coro' will be a coroutine object, and **kwargs will be488 arbitrary keyword arguments that should be passed on to Task.489 The callable must return a Task.490 """491 if factory is not None and not callable(factory):492 raise TypeError('task factory must be a callable or None')493 self._task_factory = factory494 495 def get_task_factory(self):496 """Return a task factory, or None if the default one is in use."""497 return self._task_factory498 499 def _make_socket_transport(self, sock, protocol, waiter=None, *,500 extra=None, server=None):501 """Create socket transport."""502 raise NotImplementedError503 504 def _make_ssl_transport(505 self, rawsock, protocol, sslcontext, waiter=None,506 *, server_side=False, server_hostname=None,507 extra=None, server=None,508 ssl_handshake_timeout=None,509 ssl_shutdown_timeout=None,510 call_connection_made=True):511 """Create SSL transport."""512 raise NotImplementedError513 514 def _make_datagram_transport(self, sock, protocol,515 address=None, waiter=None, extra=None):516 """Create datagram transport."""517 raise NotImplementedError518 519 def _make_read_pipe_transport(self, pipe, protocol, waiter=None,520 extra=None):521 """Create read pipe transport."""522 raise NotImplementedError523 524 def _make_write_pipe_transport(self, pipe, protocol, waiter=None,525 extra=None):526 """Create write pipe transport."""527 raise NotImplementedError528 529 async def _make_subprocess_transport(self, protocol, args, shell,530 stdin, stdout, stderr, bufsize,531 extra=None, **kwargs):532 """Create subprocess transport."""533 raise NotImplementedError534 535 def _write_to_self(self):536 """Write a byte to self-pipe, to wake up the event loop.537 538 This may be called from a different thread.539 540 The subclass is responsible for implementing the self-pipe.541 """542 raise NotImplementedError543 544 def _process_events(self, event_list):545 """Process selector events."""546 raise NotImplementedError547 548 def _check_closed(self):549 if self._closed:550 raise RuntimeError('Event loop is closed')551 552 def _check_default_executor(self):553 if self._executor_shutdown_called:554 raise RuntimeError('Executor shutdown has been called')555 556 def _asyncgen_finalizer_hook(self, agen):557 self._asyncgens.discard(agen)558 if not self.is_closed():559 self.call_soon_threadsafe(self.create_task, agen.aclose())560 561 def _asyncgen_firstiter_hook(self, agen):562 if self._asyncgens_shutdown_called:563 warnings.warn(564 f"asynchronous generator {agen!r} was scheduled after "565 f"loop.shutdown_asyncgens() call",566 ResourceWarning, source=self)567 568 self._asyncgens.add(agen)569 570 async def shutdown_asyncgens(self):571 """Shutdown all active asynchronous generators."""572 self._asyncgens_shutdown_called = True573 574 if not len(self._asyncgens):575 # If Python version is <3.6 or we don't have any asynchronous576 # generators alive.577 return578 579 closing_agens = list(self._asyncgens)580 self._asyncgens.clear()581 582 results = await tasks.gather(583 *[ag.aclose() for ag in closing_agens],584 return_exceptions=True)585 586 for result, agen in zip(results, closing_agens):587 if isinstance(result, Exception):588 self.call_exception_handler({589 'message': f'an error occurred during closing of '590 f'asynchronous generator {agen!r}',591 'exception': result,592 'asyncgen': agen593 })594 595 async def shutdown_default_executor(self, timeout=None):596 """Schedule the shutdown of the default executor.597 598 The timeout parameter specifies the amount of time the executor will599 be given to finish joining. The default value is None, which means600 that the executor will be given an unlimited amount of time.601 """602 self._executor_shutdown_called = True603 if self._default_executor is None:604 return605 future = self.create_future()606 thread = threading.Thread(target=self._do_shutdown, args=(future,))607 thread.start()608 try:609 async with timeouts.timeout(timeout):610 await future611 except TimeoutError:612 warnings.warn("The executor did not finishing joining "613 f"its threads within {timeout} seconds.",614 RuntimeWarning, stacklevel=2)615 self._default_executor.shutdown(wait=False)616 else:617 thread.join()618 619 def _do_shutdown(self, future):620 try:621 self._default_executor.shutdown(wait=True)622 if not self.is_closed():623 self.call_soon_threadsafe(futures._set_result_unless_cancelled,624 future, None)625 except Exception as ex:626 if not self.is_closed() and not future.cancelled():627 self.call_soon_threadsafe(future.set_exception, ex)628 629 def _check_running(self):630 if self.is_running():631 raise RuntimeError('This event loop is already running')632 if events._get_running_loop() is not None:633 raise RuntimeError(634 'Cannot run the event loop while another loop is running')635 636 def _run_forever_setup(self):637 """Prepare the run loop to process events.638 639 This method exists so that custom event loop subclasses (e.g., event loops640 that integrate a GUI event loop with Python's event loop) have access to all the641 loop setup logic.642 """643 self._check_closed()644 self._check_running()645 self._set_coroutine_origin_tracking(self._debug)646 647 self._old_agen_hooks = sys.get_asyncgen_hooks()648 self._thread_id = threading.get_ident()649 sys.set_asyncgen_hooks(650 firstiter=self._asyncgen_firstiter_hook,651 finalizer=self._asyncgen_finalizer_hook652 )653 654 events._set_running_loop(self)655 656 def _run_forever_cleanup(self):657 """Clean up after an event loop finishes the looping over events.658 659 This method exists so that custom event loop subclasses (e.g., event loops660 that integrate a GUI event loop with Python's event loop) have access to all the661 loop cleanup logic.662 """663 self._stopping = False664 self._thread_id = None665 events._set_running_loop(None)666 self._set_coroutine_origin_tracking(False)667 # Restore any pre-existing async generator hooks.668 if self._old_agen_hooks is not None:669 sys.set_asyncgen_hooks(*self._old_agen_hooks)670 self._old_agen_hooks = None671 672 def run_forever(self):673 """Run until stop() is called."""674 self._run_forever_setup()675 try:676 while True:677 self._run_once()678 if self._stopping:679 break680 finally:681 self._run_forever_cleanup()682 683 def run_until_complete(self, future):684 """Run until the Future is done.685 686 If the argument is a coroutine, it is wrapped in a Task.687 688 WARNING: It would be disastrous to call run_until_complete()689 with the same coroutine twice -- it would wrap it in two690 different Tasks and that can't be good.691 692 Return the Future's result, or raise its exception.693 """694 self._check_closed()695 self._check_running()696 697 new_task = not futures.isfuture(future)698 future = tasks.ensure_future(future, loop=self)699 if new_task:700 # An exception is raised if the future didn't complete, so there701 # is no need to log the "destroy pending task" message702 future._log_destroy_pending = False703 704 future.add_done_callback(_run_until_complete_cb)705 try:706 self.run_forever()707 except:708 if new_task and future.done() and not future.cancelled():709 # The coroutine raised a BaseException. Consume the exception710 # to not log a warning, the caller doesn't have access to the711 # local task.712 future.exception()713 raise714 finally:715 future.remove_done_callback(_run_until_complete_cb)716 if not future.done():717 raise RuntimeError('Event loop stopped before Future completed.')718 719 return future.result()720 721 def stop(self):722 """Stop running the event loop.723 724 Every callback already scheduled will still run. This simply informs725 run_forever to stop looping after a complete iteration.726 """727 self._stopping = True728 729 def close(self):730 """Close the event loop.731 732 This clears the queues and shuts down the executor,733 but does not wait for the executor to finish.734 735 The event loop must not be running.736 """737 if self.is_running():738 raise RuntimeError("Cannot close a running event loop")739 if self._closed:740 return741 if self._debug:742 logger.debug("Close %r", self)743 self._closed = True744 self._ready.clear()745 self._scheduled.clear()746 self._executor_shutdown_called = True747 executor = self._default_executor748 if executor is not None:749 self._default_executor = None750 executor.shutdown(wait=False)751 752 def is_closed(self):753 """Returns True if the event loop was closed."""754 return self._closed755 756 def __del__(self, _warn=warnings.warn):757 if not self.is_closed():758 _warn(f"unclosed event loop {self!r}", ResourceWarning, source=self)759 if not self.is_running():760 self.close()761 762 def is_running(self):763 """Returns True if the event loop is running."""764 return (self._thread_id is not None)765 766 def time(self):767 """Return the time according to the event loop's clock.768 769 This is a float expressed in seconds since an epoch, but the770 epoch, precision, accuracy and drift are unspecified and may771 differ per event loop.772 """773 return time.monotonic()774 775 def call_later(self, delay, callback, *args, context=None):776 """Arrange for a callback to be called at a given time.777 778 Return a Handle: an opaque object with a cancel() method that779 can be used to cancel the call.780 781 The delay can be an int or float, expressed in seconds. It is782 always relative to the current time.783 784 Each callback will be called exactly once. If two callbacks785 are scheduled for exactly the same time, it is undefined which786 will be called first.787 788 Any positional arguments after the callback will be passed to789 the callback when it is called.790 """791 if delay is None:792 raise TypeError('delay must not be None')793 timer = self.call_at(self.time() + delay, callback, *args,794 context=context)795 if timer._source_traceback:796 del timer._source_traceback[-1]797 return timer798 799 def call_at(self, when, callback, *args, context=None):800 """Like call_later(), but uses an absolute time.801 802 Absolute time corresponds to the event loop's time() method.803 """804 if when is None:805 raise TypeError("when cannot be None")806 self._check_closed()807 if self._debug:808 self._check_thread()809 self._check_callback(callback, 'call_at')810 timer = events.TimerHandle(when, callback, args, self, context)811 if timer._source_traceback:812 del timer._source_traceback[-1]813 heapq.heappush(self._scheduled, timer)814 timer._scheduled = True815 return timer816 817 def call_soon(self, callback, *args, context=None):818 """Arrange for a callback to be called as soon as possible.819 820 This operates as a FIFO queue: callbacks are called in the821 order in which they are registered. Each callback will be822 called exactly once.823 824 Any positional arguments after the callback will be passed to825 the callback when it is called.826 """827 self._check_closed()828 if self._debug:829 self._check_thread()830 self._check_callback(callback, 'call_soon')831 handle = self._call_soon(callback, args, context)832 if handle._source_traceback:833 del handle._source_traceback[-1]834 return handle835 836 def _check_callback(self, callback, method):837 if (coroutines.iscoroutine(callback) or838 coroutines._iscoroutinefunction(callback)):839 raise TypeError(840 f"coroutines cannot be used with {method}()")841 if not callable(callback):842 raise TypeError(843 f'a callable object was expected by {method}(), '844 f'got {callback!r}')845 846 def _call_soon(self, callback, args, context):847 handle = events.Handle(callback, args, self, context)848 if handle._source_traceback:849 del handle._source_traceback[-1]850 self._ready.append(handle)851 return handle852 853 def _check_thread(self):854 """Check that the current thread is the thread running the event loop.855 856 Non-thread-safe methods of this class make this assumption and will857 likely behave incorrectly when the assumption is violated.858 859 Should only be called when (self._debug == True). The caller is860 responsible for checking this condition for performance reasons.861 """862 if self._thread_id is None:863 return864 thread_id = threading.get_ident()865 if thread_id != self._thread_id:866 raise RuntimeError(867 "Non-thread-safe operation invoked on an event loop other "868 "than the current one")869 870 def call_soon_threadsafe(self, callback, *args, context=None):871 """Like call_soon(), but thread-safe."""872 self._check_closed()873 if self._debug:874 self._check_callback(callback, 'call_soon_threadsafe')875 handle = events._ThreadSafeHandle(callback, args, self, context)876 self._ready.append(handle)877 if handle._source_traceback:878 del handle._source_traceback[-1]879 if handle._source_traceback:880 del handle._source_traceback[-1]881 self._write_to_self()882 return handle883 884 def run_in_executor(self, executor, func, *args):885 self._check_closed()886 if self._debug:887 self._check_callback(func, 'run_in_executor')888 if executor is None:889 executor = self._default_executor890 # Only check when the default executor is being used891 self._check_default_executor()892 if executor is None:893 executor = concurrent.futures.ThreadPoolExecutor(894 thread_name_prefix='asyncio'895 )896 self._default_executor = executor897 return futures.wrap_future(898 executor.submit(func, *args), loop=self)899 900 def set_default_executor(self, executor):901 if not isinstance(executor, concurrent.futures.ThreadPoolExecutor):902 raise TypeError('executor must be ThreadPoolExecutor instance')903 self._default_executor = executor904 905 def _getaddrinfo_debug(self, host, port, family, type, proto, flags):906 msg = [f"{host}:{port!r}"]907 if family:908 msg.append(f'family={family!r}')909 if type:910 msg.append(f'type={type!r}')911 if proto:912 msg.append(f'proto={proto!r}')913 if flags:914 msg.append(f'flags={flags!r}')915 msg = ', '.join(msg)916 logger.debug('Get address info %s', msg)917 918 t0 = self.time()919 addrinfo = socket.getaddrinfo(host, port, family, type, proto, flags)920 dt = self.time() - t0921 922 msg = f'Getting address info {msg} took {dt * 1e3:.3f}ms: {addrinfo!r}'923 if dt >= self.slow_callback_duration:924 logger.info(msg)925 else:926 logger.debug(msg)927 return addrinfo928 929 async def getaddrinfo(self, host, port, *,930 family=0, type=0, proto=0, flags=0):931 if self._debug:932 getaddr_func = self._getaddrinfo_debug933 else:934 getaddr_func = socket.getaddrinfo935 936 return await self.run_in_executor(937 None, getaddr_func, host, port, family, type, proto, flags)938 939 async def getnameinfo(self, sockaddr, flags=0):940 return await self.run_in_executor(941 None, socket.getnameinfo, sockaddr, flags)942 943 async def sock_sendfile(self, sock, file, offset=0, count=None,944 *, fallback=True):945 if self._debug and sock.gettimeout() != 0:946 raise ValueError("the socket must be non-blocking")947 _check_ssl_socket(sock)948 self._check_sendfile_params(sock, file, offset, count)949 try:950 return await self._sock_sendfile_native(sock, file,951 offset, count)952 except exceptions.SendfileNotAvailableError as exc:953 if not fallback:954 raise955 return await self._sock_sendfile_fallback(sock, file,956 offset, count)957 958 async def _sock_sendfile_native(self, sock, file, offset, count):959 # NB: sendfile syscall is not supported for SSL sockets and960 # non-mmap files even if sendfile is supported by OS961 raise exceptions.SendfileNotAvailableError(962 f"syscall sendfile is not available for socket {sock!r} "963 f"and file {file!r} combination")964 965 async def _sock_sendfile_fallback(self, sock, file, offset, count):966 if offset:967 file.seek(offset)968 blocksize = (969 min(count, constants.SENDFILE_FALLBACK_READBUFFER_SIZE)970 if count else constants.SENDFILE_FALLBACK_READBUFFER_SIZE971 )972 buf = bytearray(blocksize)973 total_sent = 0974 try:975 while True:976 if count:977 blocksize = min(count - total_sent, blocksize)978 if blocksize <= 0:979 break980 view = memoryview(buf)[:blocksize]981 read = await self.run_in_executor(None, file.readinto, view)982 if not read:983 break # EOF984 await self.sock_sendall(sock, view[:read])985 total_sent += read986 return total_sent987 finally:988 if total_sent > 0 and hasattr(file, 'seek'):989 file.seek(offset + total_sent)990 991 def _check_sendfile_params(self, sock, file, offset, count):992 if 'b' not in getattr(file, 'mode', 'b'):993 raise ValueError("file should be opened in binary mode")994 if not sock.type == socket.SOCK_STREAM:995 raise ValueError("only SOCK_STREAM type sockets are supported")996 if count is not None:997 if not isinstance(count, int):998 raise TypeError(999 "count must be a positive integer (got {!r})".format(count))1000 if count <= 0:1001 raise ValueError(1002 "count must be a positive integer (got {!r})".format(count))1003 if not isinstance(offset, int):1004 raise TypeError(1005 "offset must be a non-negative integer (got {!r})".format(1006 offset))1007 if offset < 0:1008 raise ValueError(1009 "offset must be a non-negative integer (got {!r})".format(1010 offset))1011 1012 async def _connect_sock(self, exceptions, addr_info, local_addr_infos=None):1013 """Create, bind and connect one socket."""1014 my_exceptions = []1015 exceptions.append(my_exceptions)1016 family, type_, proto, _, address = addr_info1017 sock = None1018 try:1019 try:1020 sock = socket.socket(family=family, type=type_, proto=proto)1021 sock.setblocking(False)1022 if local_addr_infos is not None:1023 for lfamily, _, _, _, laddr in local_addr_infos:1024 # skip local addresses of different family1025 if lfamily != family:1026 continue1027 try:1028 sock.bind(laddr)1029 break1030 except OSError as exc:1031 msg = (1032 f'error while attempting to bind on '1033 f'address {laddr!r}: {str(exc).lower()}'1034 )1035 exc = OSError(exc.errno, msg)1036 my_exceptions.append(exc)1037 else: # all bind attempts failed1038 if my_exceptions:1039 raise my_exceptions.pop()1040 else:1041 raise OSError(f"no matching local address with {family=} found")1042 await self.sock_connect(sock, address)1043 return sock1044 except OSError as exc:1045 my_exceptions.append(exc)1046 raise1047 except:1048 if sock is not None:1049 try:1050 sock.close()1051 except OSError:1052 # An error when closing a newly created socket is1053 # not important, but it can overwrite more important1054 # non-OSError error. So ignore it.1055 pass1056 raise1057 finally:1058 exceptions = my_exceptions = None1059 1060 async def create_connection(1061 self, protocol_factory, host=None, port=None,1062 *, ssl=None, family=0,1063 proto=0, flags=0, sock=None,1064 local_addr=None, server_hostname=None,1065 ssl_handshake_timeout=None,1066 ssl_shutdown_timeout=None,1067 happy_eyeballs_delay=None, interleave=None,1068 all_errors=False):1069 """Connect to a TCP server.1070 1071 Create a streaming transport connection to a given internet host and1072 port: socket family AF_INET or socket.AF_INET6 depending on host (or1073 family if specified), socket type SOCK_STREAM. protocol_factory must be1074 a callable returning a protocol instance.1075 1076 This method is a coroutine which will try to establish the connection1077 in the background. When successful, the coroutine returns a1078 (transport, protocol) pair.1079 """1080 if server_hostname is not None and not ssl:1081 raise ValueError('server_hostname is only meaningful with ssl')1082 1083 if server_hostname is None and ssl:1084 # Use host as default for server_hostname. It is an error1085 # if host is empty or not set, e.g. when an1086 # already-connected socket was passed or when only a port1087 # is given. To avoid this error, you can pass1088 # server_hostname='' -- this will bypass the hostname1089 # check. (This also means that if host is a numeric1090 # IP/IPv6 address, we will attempt to verify that exact1091 # address; this will probably fail, but it is possible to1092 # create a certificate for a specific IP address, so we1093 # don't judge it here.)1094 if not host:1095 raise ValueError('You must set server_hostname '1096 'when using ssl without a host')1097 server_hostname = host1098 1099 if ssl_handshake_timeout is not None and not ssl:1100 raise ValueError(1101 'ssl_handshake_timeout is only meaningful with ssl')1102 1103 if ssl_shutdown_timeout is not None and not ssl:1104 raise ValueError(1105 'ssl_shutdown_timeout is only meaningful with ssl')1106 1107 if sock is not None:1108 _check_ssl_socket(sock)1109 1110 if happy_eyeballs_delay is not None and interleave is None:1111 # If using happy eyeballs, default to interleave addresses by family1112 interleave = 11113 1114 if host is not None or port is not None:1115 if sock is not None:1116 raise ValueError(1117 'host/port and sock can not be specified at the same time')1118 1119 infos = await self._ensure_resolved(1120 (host, port), family=family,1121 type=socket.SOCK_STREAM, proto=proto, flags=flags, loop=self)1122 if not infos:1123 raise OSError('getaddrinfo() returned empty list')1124 1125 if local_addr is not None:1126 laddr_infos = await self._ensure_resolved(1127 local_addr, family=family,1128 type=socket.SOCK_STREAM, proto=proto,1129 flags=flags, loop=self)1130 if not laddr_infos:1131 raise OSError('getaddrinfo() returned empty list')1132 else:1133 laddr_infos = None1134 1135 if interleave:1136 infos = _interleave_addrinfos(infos, interleave)1137 1138 exceptions = []1139 if happy_eyeballs_delay is None:1140 # not using happy eyeballs1141 for addrinfo in infos:1142 try:1143 sock = await self._connect_sock(1144 exceptions, addrinfo, laddr_infos)1145 break1146 except OSError:1147 continue1148 else: # using happy eyeballs1149 sock = (await staggered.staggered_race(1150 (1151 # can't use functools.partial as it keeps a reference1152 # to exceptions1153 lambda addrinfo=addrinfo: self._connect_sock(1154 exceptions, addrinfo, laddr_infos1155 )1156 for addrinfo in infos1157 ),1158 happy_eyeballs_delay,1159 loop=self,1160 ))[0] # can't use sock, _, _ as it keeks a reference to exceptions1161 1162 if sock is None:1163 exceptions = [exc for sub in exceptions for exc in sub]1164 try:1165 if all_errors:1166 raise ExceptionGroup("create_connection failed", exceptions)1167 if len(exceptions) == 1:1168 raise exceptions[0]1169 elif exceptions:1170 # If they all have the same str(), raise one.1171 model = str(exceptions[0])1172 if all(str(exc) == model for exc in exceptions):1173 raise exceptions[0]1174 # Raise a combined exception so the user can see all1175 # the various error messages.1176 raise OSError('Multiple exceptions: {}'.format(1177 ', '.join(str(exc) for exc in exceptions)))1178 else:1179 # No exceptions were collected, raise a timeout error1180 raise TimeoutError('create_connection failed')1181 finally:1182 exceptions = None1183 1184 else:1185 if sock is None:1186 raise ValueError(1187 'host and port was not specified and no sock specified')1188 if sock.type != socket.SOCK_STREAM:1189 # We allow AF_INET, AF_INET6, AF_UNIX as long as they1190 # are SOCK_STREAM.1191 # We support passing AF_UNIX sockets even though we have1192 # a dedicated API for that: create_unix_connection.1193 # Disallowing AF_UNIX in this method, breaks backwards1194 # compatibility.1195 raise ValueError(1196 f'A Stream Socket was expected, got {sock!r}')1197 1198 transport, protocol = await self._create_connection_transport(1199 sock, protocol_factory, ssl, server_hostname,1200 ssl_handshake_timeout=ssl_handshake_timeout,