Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
asyncio.py754 linesDownload Raw Back to platform
1"""Bridges between the `asyncio` module and Tornado IOLoop.
2
3.. versionadded:: 3.2
4
5This module integrates Tornado with the ``asyncio`` module introduced
6in Python 3.4. This makes it possible to combine the two libraries on
7the same event loop.
8
9.. deprecated:: 5.0
10
11   While the code in this module is still used, it is now enabled
12   automatically when `asyncio` is available, so applications should
13   no longer need to refer to this module directly.
14
15.. note::
16
17   Tornado is designed to use a selector-based event loop. On Windows,
18   where a proactor-based event loop has been the default since Python 3.8,
19   a selector event loop is emulated by running ``select`` on a separate thread.
20   Configuring ``asyncio`` to use a selector event loop may improve performance
21   of Tornado (but may reduce performance of other ``asyncio``-based libraries
22   in the same process).
23"""
24
25import asyncio
26import atexit
27import concurrent.futures
28import contextvars
29import errno
30import functools
31import select
32import socket
33import sys
34import threading
35import typing
36import warnings
37from tornado.gen import convert_yielded
38from tornado.ioloop import IOLoop, _Selectable
39
40from typing import (
41    Any,
42    Callable,
43    Dict,
44    List,
45    Optional,
46    Protocol,
47    Set,
48    Tuple,
49    TypeVar,
50    Union,
51)
52
53if typing.TYPE_CHECKING:
54    from typing_extensions import TypeVarTuple, Unpack
55
56
57class _HasFileno(Protocol):
58    def fileno(self) -> int:
59        pass
60
61
62_FileDescriptorLike = Union[int, _HasFileno]
63
64_T = TypeVar("_T")
65
66if typing.TYPE_CHECKING:
67    _Ts = TypeVarTuple("_Ts")
68
69# Collection of selector thread event loops to shut down on exit.
70_selector_loops: Set["SelectorThread"] = set()
71
72
73def _atexit_callback() -> None:
74    for loop in _selector_loops:
75        with loop._select_cond:
76            loop._closing_selector = True
77            loop._select_cond.notify()
78        try:
79            loop._waker_w.send(b"a")
80        except BlockingIOError:
81            pass
82        if loop._thread is not None:
83            # If we don't join our (daemon) thread here, we may get a deadlock
84            # during interpreter shutdown. I don't really understand why. This
85            # deadlock happens every time in CI (both travis and appveyor) but
86            # I've never been able to reproduce locally.
87            loop._thread.join()
88    _selector_loops.clear()
89
90
91atexit.register(_atexit_callback)
92
93
94class BaseAsyncIOLoop(IOLoop):
95    def initialize(  # type: ignore
96        self, asyncio_loop: asyncio.AbstractEventLoop, **kwargs: Any
97    ) -> None:
98        # asyncio_loop is always the real underlying IOLoop. This is used in
99        # ioloop.py to maintain the asyncio-to-ioloop mappings.
100        self.asyncio_loop = asyncio_loop
101        # selector_loop is an event loop that implements the add_reader family of
102        # methods. Usually the same as asyncio_loop but differs on platforms such
103        # as windows where the default event loop does not implement these methods.
104        self.selector_loop = asyncio_loop
105        if hasattr(asyncio, "ProactorEventLoop") and isinstance(
106            asyncio_loop, asyncio.ProactorEventLoop
107        ):
108            # Ignore this line for mypy because the abstract method checker
109            # doesn't understand dynamic proxies.
110            self.selector_loop = AddThreadSelectorEventLoop(asyncio_loop)  # type: ignore
111        # Maps fd to (fileobj, handler function) pair (as in IOLoop.add_handler)
112        self.handlers: Dict[int, Tuple[Union[int, _Selectable], Callable]] = {}
113        # Set of fds listening for reads/writes
114        self.readers: Set[int] = set()
115        self.writers: Set[int] = set()
116        self.closing = False
117        # If an asyncio loop was closed through an asyncio interface
118        # instead of IOLoop.close(), we'd never hear about it and may
119        # have left a dangling reference in our map. In case an
120        # application (or, more likely, a test suite) creates and
121        # destroys a lot of event loops in this way, check here to
122        # ensure that we don't have a lot of dead loops building up in
123        # the map.
124        #
125        # TODO(bdarnell): consider making self.asyncio_loop a weakref
126        # for AsyncIOMainLoop and make _ioloop_for_asyncio a
127        # WeakKeyDictionary.
128        for loop in IOLoop._ioloop_for_asyncio.copy():
129            if loop.is_closed():
130                try:
131                    del IOLoop._ioloop_for_asyncio[loop]
132                except KeyError:
133                    pass
134
135        # Make sure we don't already have an IOLoop for this asyncio loop
136        existing_loop = IOLoop._ioloop_for_asyncio.setdefault(asyncio_loop, self)
137        if existing_loop is not self:
138            raise RuntimeError(
139                f"IOLoop {existing_loop} already associated with asyncio loop {asyncio_loop}"
140            )
141
142        super().initialize(**kwargs)
143
144    def close(self, all_fds: bool = False) -> None:
145        self.closing = True
146        for fd in list(self.handlers):
147            fileobj, handler_func = self.handlers[fd]
148            self.remove_handler(fd)
149            if all_fds:
150                self.close_fd(fileobj)
151        # Remove the mapping before closing the asyncio loop. If this
152        # happened in the other order, we could race against another
153        # initialize() call which would see the closed asyncio loop,
154        # assume it was closed from the asyncio side, and do this
155        # cleanup for us, leading to a KeyError.
156        del IOLoop._ioloop_for_asyncio[self.asyncio_loop]
157        if self.selector_loop is not self.asyncio_loop:
158            self.selector_loop.close()
159        self.asyncio_loop.close()
160
161    def add_handler(
162        self, fd: Union[int, _Selectable], handler: Callable[..., None], events: int
163    ) -> None:
164        fd, fileobj = self.split_fd(fd)
165        if fd in self.handlers:
166            raise ValueError("fd %s added twice" % fd)
167        self.handlers[fd] = (fileobj, handler)
168        if events & IOLoop.READ:
169            self.selector_loop.add_reader(fd, self._handle_events, fd, IOLoop.READ)
170            self.readers.add(fd)
171        if events & IOLoop.WRITE:
172            self.selector_loop.add_writer(fd, self._handle_events, fd, IOLoop.WRITE)
173            self.writers.add(fd)
174
175    def update_handler(self, fd: Union[int, _Selectable], events: int) -> None:
176        fd, fileobj = self.split_fd(fd)
177        if events & IOLoop.READ:
178            if fd not in self.readers:
179                self.selector_loop.add_reader(fd, self._handle_events, fd, IOLoop.READ)
180                self.readers.add(fd)
181        else:
182            if fd in self.readers:
183                self.selector_loop.remove_reader(fd)
184                self.readers.remove(fd)
185        if events & IOLoop.WRITE:
186            if fd not in self.writers:
187                self.selector_loop.add_writer(fd, self._handle_events, fd, IOLoop.WRITE)
188                self.writers.add(fd)
189        else:
190            if fd in self.writers:
191                self.selector_loop.remove_writer(fd)
192                self.writers.remove(fd)
193
194    def remove_handler(self, fd: Union[int, _Selectable]) -> None:
195        fd, fileobj = self.split_fd(fd)
196        if fd not in self.handlers:
197            return
198        if fd in self.readers:
199            self.selector_loop.remove_reader(fd)
200            self.readers.remove(fd)
201        if fd in self.writers:
202            self.selector_loop.remove_writer(fd)
203            self.writers.remove(fd)
204        del self.handlers[fd]
205
206    def _handle_events(self, fd: int, events: int) -> None:
207        fileobj, handler_func = self.handlers[fd]
208        handler_func(fileobj, events)
209
210    def start(self) -> None:
211        self.asyncio_loop.run_forever()
212
213    def stop(self) -> None:
214        self.asyncio_loop.stop()
215
216    def call_at(
217        self, when: float, callback: Callable, *args: Any, **kwargs: Any
218    ) -> object:
219        # asyncio.call_at supports *args but not **kwargs, so bind them here.
220        # We do not synchronize self.time and asyncio_loop.time, so
221        # convert from absolute to relative.
222        return self.asyncio_loop.call_later(
223            max(0, when - self.time()),
224            self._run_callback,
225            functools.partial(callback, *args, **kwargs),
226        )
227
228    def remove_timeout(self, timeout: object) -> None:
229        timeout.cancel()  # type: ignore
230
231    def add_callback(self, callback: Callable, *args: Any, **kwargs: Any) -> None:
232        try:
233            if asyncio.get_running_loop() is self.asyncio_loop:
234                call_soon = self.asyncio_loop.call_soon
235            else:
236                call_soon = self.asyncio_loop.call_soon_threadsafe
237        except RuntimeError:
238            call_soon = self.asyncio_loop.call_soon_threadsafe
239
240        try:
241            call_soon(self._run_callback, functools.partial(callback, *args, **kwargs))
242        except RuntimeError:
243            # "Event loop is closed". Swallow the exception for
244            # consistency with PollIOLoop (and logical consistency
245            # with the fact that we can't guarantee that an
246            # add_callback that completes without error will
247            # eventually execute).
248            pass
249        except AttributeError:
250            # ProactorEventLoop may raise this instead of RuntimeError
251            # if call_soon_threadsafe races with a call to close().
252            # Swallow it too for consistency.
253            pass
254
255    def add_callback_from_signal(
256        self, callback: Callable, *args: Any, **kwargs: Any
257    ) -> None:
258        warnings.warn("add_callback_from_signal is deprecated", DeprecationWarning)
259        try:
260            self.asyncio_loop.call_soon_threadsafe(
261                self._run_callback, functools.partial(callback, *args, **kwargs)
262            )
263        except RuntimeError:
264            pass
265
266    def run_in_executor(
267        self,
268        executor: Optional[concurrent.futures.Executor],
269        func: Callable[..., _T],
270        *args: Any,
271    ) -> "asyncio.Future[_T]":
272        return self.asyncio_loop.run_in_executor(executor, func, *args)
273
274    def set_default_executor(self, executor: concurrent.futures.Executor) -> None:
275        return self.asyncio_loop.set_default_executor(executor)
276
277
278class AsyncIOMainLoop(BaseAsyncIOLoop):
279    """``AsyncIOMainLoop`` creates an `.IOLoop` that corresponds to the
280    current ``asyncio`` event loop (i.e. the one returned by
281    ``asyncio.get_event_loop()``).
282
283    .. deprecated:: 5.0
284
285       Now used automatically when appropriate; it is no longer necessary
286       to refer to this class directly.
287
288    .. versionchanged:: 5.0
289
290       Closing an `AsyncIOMainLoop` now closes the underlying asyncio loop.
291    """
292
293    def initialize(self, **kwargs: Any) -> None:  # type: ignore
294        super().initialize(asyncio.get_event_loop(), **kwargs)
295
296    def _make_current(self) -> None:
297        # AsyncIOMainLoop already refers to the current asyncio loop so
298        # nothing to do here.
299        pass
300
301
302class AsyncIOLoop(BaseAsyncIOLoop):
303    """``AsyncIOLoop`` is an `.IOLoop` that runs on an ``asyncio`` event loop.
304    This class follows the usual Tornado semantics for creating new
305    ``IOLoops``; these loops are not necessarily related to the
306    ``asyncio`` default event loop.
307
308    Each ``AsyncIOLoop`` creates a new ``asyncio.EventLoop``; this object
309    can be accessed with the ``asyncio_loop`` attribute.
310
311    .. versionchanged:: 6.2
312
313       Support explicit ``asyncio_loop`` argument
314       for specifying the asyncio loop to attach to,
315       rather than always creating a new one with the default policy.
316
317    .. versionchanged:: 5.0
318
319       When an ``AsyncIOLoop`` becomes the current `.IOLoop`, it also sets
320       the current `asyncio` event loop.
321
322    .. deprecated:: 5.0
323
324       Now used automatically when appropriate; it is no longer necessary
325       to refer to this class directly.
326    """
327
328    def initialize(self, **kwargs: Any) -> None:  # type: ignore
329        self.is_current = False
330        loop = None
331        if "asyncio_loop" not in kwargs:
332            kwargs["asyncio_loop"] = loop = asyncio.new_event_loop()
333        try:
334            super().initialize(**kwargs)
335        except Exception:
336            # If initialize() does not succeed (taking ownership of the loop),
337            # we have to close it.
338            if loop is not None:
339                loop.close()
340            raise
341
342    def close(self, all_fds: bool = False) -> None:
343        if self.is_current:
344            self._clear_current()
345        super().close(all_fds=all_fds)
346
347    def _make_current(self) -> None:
348        if not self.is_current:
349            try:
350                self.old_asyncio = asyncio.get_event_loop()
351            except (RuntimeError, AssertionError):
352                self.old_asyncio = None  # type: ignore
353            self.is_current = True
354        asyncio.set_event_loop(self.asyncio_loop)
355
356    def _clear_current_hook(self) -> None:
357        if self.is_current:
358            asyncio.set_event_loop(self.old_asyncio)
359            self.is_current = False
360
361
362def to_tornado_future(asyncio_future: asyncio.Future) -> asyncio.Future:
363    """Convert an `asyncio.Future` to a `tornado.concurrent.Future`.
364
365    .. versionadded:: 4.1
366
367    .. deprecated:: 5.0
368       Tornado ``Futures`` have been merged with `asyncio.Future`,
369       so this method is now a no-op.
370    """
371    return asyncio_future
372
373
374def to_asyncio_future(tornado_future: asyncio.Future) -> asyncio.Future:
375    """Convert a Tornado yieldable object to an `asyncio.Future`.
376
377    .. versionadded:: 4.1
378
379    .. versionchanged:: 4.3
380       Now accepts any yieldable object, not just
381       `tornado.concurrent.Future`.
382
383    .. deprecated:: 5.0
384       Tornado ``Futures`` have been merged with `asyncio.Future`,
385       so this method is now equivalent to `tornado.gen.convert_yielded`.
386    """
387    return convert_yielded(tornado_future)
388
389
390_AnyThreadEventLoopPolicy = None
391
392
393def __getattr__(name: str) -> typing.Any:
394    # The event loop policy system is deprecated in Python 3.14; simply accessing
395    # the name asyncio.DefaultEventLoopPolicy will raise a warning. Lazily create
396    # the AnyThreadEventLoopPolicy class so that the warning is only raised if
397    # the policy is used.
398    if name != "AnyThreadEventLoopPolicy":
399        raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
400
401    global _AnyThreadEventLoopPolicy
402    if _AnyThreadEventLoopPolicy is None:
403        if sys.platform == "win32" and hasattr(
404            asyncio, "WindowsSelectorEventLoopPolicy"
405        ):
406            # "Any thread" and "selector" should be orthogonal, but there's not a clean
407            # interface for composing policies so pick the right base.
408            _BasePolicy = asyncio.WindowsSelectorEventLoopPolicy  # type: ignore
409        else:
410            _BasePolicy = asyncio.DefaultEventLoopPolicy
411
412        class AnyThreadEventLoopPolicy(_BasePolicy):  # type: ignore
413            """Event loop policy that allows loop creation on any thread.
414
415            The default `asyncio` event loop policy only automatically creates
416            event loops in the main threads. Other threads must create event
417            loops explicitly or `asyncio.get_event_loop` (and therefore
418            `.IOLoop.current`) will fail. Installing this policy allows event
419            loops to be created automatically on any thread, matching the
420            behavior of Tornado versions prior to 5.0 (or 5.0 on Python 2).
421
422            Usage::
423
424                asyncio.set_event_loop_policy(AnyThreadEventLoopPolicy())
425
426            .. versionadded:: 5.0
427
428            .. deprecated:: 6.2
429
430                ``AnyThreadEventLoopPolicy`` affects the implicit creation
431                of an event loop, which is deprecated in Python 3.10 and
432                will be removed in a future version of Python. At that time
433                ``AnyThreadEventLoopPolicy`` will no longer be useful.
434                If you are relying on it, use `asyncio.new_event_loop`
435                or `asyncio.run` explicitly in any non-main threads that
436                need event loops.
437            """
438
439            def __init__(self) -> None:
440                super().__init__()
441                warnings.warn(
442                    "AnyThreadEventLoopPolicy is deprecated, use asyncio.run "
443                    "or asyncio.new_event_loop instead",
444                    DeprecationWarning,
445                    stacklevel=2,
446                )
447
448            def get_event_loop(self) -> asyncio.AbstractEventLoop:
449                try:
450                    return super().get_event_loop()
451                except RuntimeError:
452                    # "There is no current event loop in thread %r"
453                    loop = self.new_event_loop()
454                    self.set_event_loop(loop)
455                    return loop
456
457        _AnyThreadEventLoopPolicy = AnyThreadEventLoopPolicy
458
459    return _AnyThreadEventLoopPolicy
460
461
462class SelectorThread:
463    """Define ``add_reader`` methods to be called in a background select thread.
464
465    Instances of this class start a second thread to run a selector.
466    This thread is completely hidden from the user;
467    all callbacks are run on the wrapped event loop's thread.
468
469    Typically used via ``AddThreadSelectorEventLoop``,
470    but can be attached to a running asyncio loop.
471    """
472
473    _closed = False
474
475    def __init__(self, real_loop: asyncio.AbstractEventLoop) -> None:
476        self._main_thread_ctx = contextvars.copy_context()
477
478        self._real_loop = real_loop
479
480        self._select_cond = threading.Condition()
481        self._select_args: Optional[
482            Tuple[List[_FileDescriptorLike], List[_FileDescriptorLike]]
483        ] = None
484        self._closing_selector = False
485        self._thread: Optional[threading.Thread] = None
486        self._thread_manager_handle = self._thread_manager()
487
488        async def thread_manager_anext() -> None:
489            # the anext builtin wasn't added until 3.10. We just need to iterate
490            # this generator one step.
491            await self._thread_manager_handle.__anext__()
492
493        # When the loop starts, start the thread. Not too soon because we can't
494        # clean up if we get to this point but the event loop is closed without
495        # starting.
496        self._real_loop.call_soon(
497            lambda: self._real_loop.create_task(thread_manager_anext()),
498            context=self._main_thread_ctx,
499        )
500
501        self._readers: Dict[_FileDescriptorLike, Callable] = {}
502        self._writers: Dict[_FileDescriptorLike, Callable] = {}
503
504        # Writing to _waker_w will wake up the selector thread, which
505        # watches for _waker_r to be readable.
506        self._waker_r, self._waker_w = socket.socketpair()
507        self._waker_r.setblocking(False)
508        self._waker_w.setblocking(False)
509        _selector_loops.add(self)
510        self.add_reader(self._waker_r, self._consume_waker)
511
512    def close(self) -> None:
513        if self._closed:
514            return
515        with self._select_cond:
516            self._closing_selector = True
517            self._select_cond.notify()
518        self._wake_selector()
519        if self._thread is not None:
520            self._thread.join()
521        _selector_loops.discard(self)
522        self.remove_reader(self._waker_r)
523        self._waker_r.close()
524        self._waker_w.close()
525        self._closed = True
526
527    async def _thread_manager(self) -> typing.AsyncGenerator[None, None]:
528        # Create a thread to run the select system call. We manage this thread
529        # manually so we can trigger a clean shutdown from an atexit hook. Note
530        # that due to the order of operations at shutdown, only daemon threads
531        # can be shut down in this way (non-daemon threads would require the
532        # introduction of a new hook: https://bugs.python.org/issue41962)
533        self._thread = threading.Thread(
534            name="Tornado selector",
535            daemon=True,
536            target=self._run_select,
537        )
538        self._thread.start()
539        self._start_select()
540        try:
541            # The presense of this yield statement means that this coroutine
542            # is actually an asynchronous generator, which has a special
543            # shutdown protocol. We wait at this yield point until the
544            # event loop's shutdown_asyncgens method is called, at which point
545            # we will get a GeneratorExit exception and can shut down the
546            # selector thread.
547            yield
548        except GeneratorExit:
549            self.close()
550            raise
551
552    def _wake_selector(self) -> None:
553        if self._closed:
554            return
555        try:
556            self._waker_w.send(b"a")
557        except BlockingIOError:
558            pass
559
560    def _consume_waker(self) -> None:
561        try:
562            self._waker_r.recv(1024)
563        except BlockingIOError:
564            pass
565
566    def _start_select(self) -> None:
567        # Capture reader and writer sets here in the event loop
568        # thread to avoid any problems with concurrent
569        # modification while the select loop uses them.
570        with self._select_cond:
571            assert self._select_args is None
572            self._select_args = (list(self._readers.keys()), list(self._writers.keys()))
573            self._select_cond.notify()
574
575    def _run_select(self) -> None:
576        while True:
577            with self._select_cond:
578                while self._select_args is None and not self._closing_selector:
579                    self._select_cond.wait()
580                if self._closing_selector:
581                    return
582                assert self._select_args is not None
583                to_read, to_write = self._select_args
584                self._select_args = None
585
586            # We use the simpler interface of the select module instead of
587            # the more stateful interface in the selectors module because
588            # this class is only intended for use on windows, where
589            # select.select is the only option. The selector interface
590            # does not have well-documented thread-safety semantics that
591            # we can rely on so ensuring proper synchronization would be
592            # tricky.
593            try:
594                # On windows, selecting on a socket for write will not
595                # return the socket when there is an error (but selecting
596                # for reads works). Also select for errors when selecting
597                # for writes, and merge the results.
598                #
599                # This pattern is also used in
600                # https://github.com/python/cpython/blob/v3.8.0/Lib/selectors.py#L312-L317
601                rs, ws, xs = select.select(to_read, to_write, to_write)
602                ws = ws + xs
603            except OSError as e:
604                # After remove_reader or remove_writer is called, the file
605                # descriptor may subsequently be closed on the event loop
606                # thread. It's possible that this select thread hasn't
607                # gotten into the select system call by the time that
608                # happens in which case (at least on macOS), select may
609                # raise a "bad file descriptor" error. If we get that
610                # error, check and see if we're also being woken up by
611                # polling the waker alone. If we are, just return to the
612                # event loop and we'll get the updated set of file
613                # descriptors on the next iteration. Otherwise, raise the
614                # original error.
615                if e.errno == getattr(errno, "WSAENOTSOCK", errno.EBADF):
616                    rs, _, _ = select.select([self._waker_r.fileno()], [], [], 0)
617                    if rs:
618                        ws = []
619                    else:
620                        raise
621                else:
622                    raise
623
624            try:
625                self._real_loop.call_soon_threadsafe(
626                    self._handle_select, rs, ws, context=self._main_thread_ctx
627                )
628            except RuntimeError:
629                # "Event loop is closed". Swallow the exception for
630                # consistency with PollIOLoop (and logical consistency
631                # with the fact that we can't guarantee that an
632                # add_callback that completes without error will
633                # eventually execute).
634                pass
635            except AttributeError:
636                # ProactorEventLoop may raise this instead of RuntimeError
637                # if call_soon_threadsafe races with a call to close().
638                # Swallow it too for consistency.
639                pass
640
641    def _handle_select(
642        self, rs: List[_FileDescriptorLike], ws: List[_FileDescriptorLike]
643    ) -> None:
644        for r in rs:
645            self._handle_event(r, self._readers)
646        for w in ws:
647            self._handle_event(w, self._writers)
648        self._start_select()
649
650    def _handle_event(
651        self,
652        fd: _FileDescriptorLike,
653        cb_map: Dict[_FileDescriptorLike, Callable],
654    ) -> None:
655        try:
656            callback = cb_map[fd]
657        except KeyError:
658            return
659        callback()
660
661    def add_reader(
662        self, fd: _FileDescriptorLike, callback: Callable[..., None], *args: Any
663    ) -> None:
664        self._readers[fd] = functools.partial(callback, *args)
665        self._wake_selector()
666
667    def add_writer(
668        self, fd: _FileDescriptorLike, callback: Callable[..., None], *args: Any
669    ) -> None:
670        self._writers[fd] = functools.partial(callback, *args)
671        self._wake_selector()
672
673    def remove_reader(self, fd: _FileDescriptorLike) -> bool:
674        try:
675            del self._readers[fd]
676        except KeyError:
677            return False
678        self._wake_selector()
679        return True
680
681    def remove_writer(self, fd: _FileDescriptorLike) -> bool:
682        try:
683            del self._writers[fd]
684        except KeyError:
685            return False
686        self._wake_selector()
687        return True
688
689
690class AddThreadSelectorEventLoop(asyncio.AbstractEventLoop):
691    """Wrap an event loop to add implementations of the ``add_reader`` method family.
692
693    Instances of this class start a second thread to run a selector.
694    This thread is completely hidden from the user; all callbacks are
695    run on the wrapped event loop's thread.
696
697    This class is used automatically by Tornado; applications should not need
698    to refer to it directly.
699
700    It is safe to wrap any event loop with this class, although it only makes sense
701    for event loops that do not implement the ``add_reader`` family of methods
702    themselves (i.e. ``WindowsProactorEventLoop``)
703
704    Closing the ``AddThreadSelectorEventLoop`` also closes the wrapped event loop.
705
706    """
707
708    # This class is a __getattribute__-based proxy. All attributes other than those
709    # in this set are proxied through to the underlying loop.
710    MY_ATTRIBUTES = {
711        "_real_loop",
712        "_selector",
713        "add_reader",
714        "add_writer",
715        "close",
716        "remove_reader",
717        "remove_writer",
718    }
719
720    def __getattribute__(self, name: str) -> Any:
721        if name in AddThreadSelectorEventLoop.MY_ATTRIBUTES:
722            return super().__getattribute__(name)
723        return getattr(self._real_loop, name)
724
725    def __init__(self, real_loop: asyncio.AbstractEventLoop) -> None:
726        self._real_loop = real_loop
727        self._selector = SelectorThread(real_loop)
728
729    def close(self) -> None:
730        self._selector.close()
731        self._real_loop.close()
732
733    def add_reader(
734        self,
735        fd: "_FileDescriptorLike",
736        callback: Callable[..., None],
737        *args: "Unpack[_Ts]",
738    ) -> None:
739        return self._selector.add_reader(fd, callback, *args)
740
741    def add_writer(
742        self,
743        fd: "_FileDescriptorLike",
744        callback: Callable[..., None],
745        *args: "Unpack[_Ts]",
746    ) -> None:
747        return self._selector.add_writer(fd, callback, *args)
748
749    def remove_reader(self, fd: "_FileDescriptorLike") -> bool:
750        return self._selector.remove_reader(fd)
751
752    def remove_writer(self, fd: "_FileDescriptorLike") -> bool:
753        return self._selector.remove_writer(fd)
754 
codekingpro/portable-devtools · Team Ai