Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
base_events.py2083 linesDownload Raw Back to asyncio
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,

Showing the first 1,200 of 2083 lines. Download the file for the rest.

codekingpro/portable-devtools · Team Ai