codekingpro/portable-devtools
114k
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 