/
usr
/
lib
/
python3
/
dist-packages
/
twisted
/
internet
/
/usr/lib/python3/dist-packages/twisted/internet
mkdir
upload
Name
Size
Mode
Actions
iocpreactor/
-
0755
rm
test/
-
0755
rm
__pycache__/
-
0755
rm
abstract.py
19295
0644
edit
dl
rm
address.py
5244
0644
edit
dl
rm
asyncioreactor.py
11131
0644
edit
dl
rm
base.py
47392
0644
edit
dl
rm
cfreactor.py
17499
0644
edit
dl
rm
default.py
1893
0644
edit
dl
rm
defer.py
85654
0644
edit
dl
rm
endpoints.py
77440
0644
edit
dl
rm
epollreactor.py
8941
0644
edit
dl
rm
error.py
13482
0644
edit
dl
rm
fdesc.py
3237
0644
edit
dl
rm
gireactor.py
4620
0644
edit
dl
rm
glib2reactor.py
1115
0644
edit
dl
rm
gtk2reactor.py
3640
0644
edit
dl
rm
gtk3reactor.py
1527
0644
edit
dl
rm
inotify.py
14396
0644
edit
dl
rm
interfaces.py
98048
0644
edit
dl
rm
kqreactor.py
10818
0644
edit
dl
rm
main.py
1006
0644
edit
dl
rm
pollreactor.py
5974
0644
edit
dl
rm
posixbase.py
27604
0644
edit
dl
rm
process.py
38516
0644
edit
dl
rm
protocol.py
27391
0644
edit
dl
rm
pyuisupport.py
843
0644
edit
dl
rm
reactor.py
1816
0644
edit
dl
rm
selectreactor.py
6102
0644
edit
dl
rm
serialport.py
2272
0644
edit
dl
rm
ssl.py
8643
0644
edit
dl
rm
stdio.py
1000
0644
edit
dl
rm
task.py
33608
0644
edit
dl
rm
tcp.py
54980
0644
edit
dl
rm
testing.py
29232
0644
edit
dl
rm
threads.py
3812
0644
edit
dl
rm
tksupport.py
1971
0644
edit
dl
rm
udp.py
18619
0644
edit
dl
rm
unix.py
22508
0644
edit
dl
rm
utils.py
8682
0644
edit
dl
rm
win32eventreactor.py
15266
0644
edit
dl
rm
wxreactor.py
5315
0644
edit
dl
rm
wxsupport.py
1305
0644
edit
dl
rm
_baseprocess.py
2003
0644
edit
dl
rm
_dumbwin32proc.py
12776
0644
edit
dl
rm
_glibbase.py
12706
0644
edit
dl
rm
_idna.py
1422
0644
edit
dl
rm
_newtls.py
9157
0644
edit
dl
rm
_pollingfile.py
8791
0644
edit
dl
rm
_posixserialport.py
2081
0644
edit
dl
rm
_posixstdio.py
4996
0644
edit
dl
rm
_producer_helpers.py
3909
0644
edit
dl
rm
_resolver.py
8465
0644
edit
dl
rm
_signals.py
2670
0644
edit
dl
rm
_sslverify.py
72796
0644
edit
dl
rm
_threadedselect.py
11582
0644
edit
dl
rm
_win32serialport.py
4914
0644
edit
dl
rm
_win32stdio.py
3140
0644
edit
dl
rm
__init__.py
521
0644
edit
dl
rm
Edit:
/usr/lib/python3/dist-packages/twisted/internet/asyncioreactor.py
(11131B)
# -*- test-case-name: twisted.test.test_internet -*- # Copyright (c) Twisted Matrix Laboratories. # See LICENSE for details. """ asyncio-based reactor implementation. """ import errno import sys from asyncio import AbstractEventLoop, get_event_loop from typing import Dict, Optional, Type from zope.interface import implementer from twisted.internet.abstract import FileDescriptor from twisted.internet.interfaces import IReactorFDSet from twisted.internet.posixbase import ( _NO_FILEDESC, PosixReactorBase, _ContinuousPolling, ) from twisted.logger import Logger from twisted.python.log import callWithLogger @implementer(IReactorFDSet) class AsyncioSelectorReactor(PosixReactorBase): """ Reactor running on top of L{asyncio.SelectorEventLoop}. On POSIX platforms, the default event loop is L{asyncio.SelectorEventLoop}. On Windows, the default event loop on Python 3.7 and older is C{asyncio.WindowsSelectorEventLoop}, but on Python 3.8 and newer the default event loop is C{asyncio.WindowsProactorEventLoop} which is incompatible with L{AsyncioSelectorReactor}. Applications that use L{AsyncioSelectorReactor} on Windows with Python 3.8+ must call C{asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy())} before instantiating and running L{AsyncioSelectorReactor}. """ _asyncClosed = False _log = Logger() def __init__(self, eventloop: Optional[AbstractEventLoop] = None): if eventloop is None: _eventloop: AbstractEventLoop = get_event_loop() else: _eventloop = eventloop # On Python 3.8+, asyncio.get_event_loop() on # Windows was changed to return a ProactorEventLoop # unless the loop policy has been changed. if sys.platform == "win32": from asyncio import ProactorEventLoop if isinstance(_eventloop, ProactorEventLoop): raise TypeError( f"ProactorEventLoop is not supported, got: {_eventloop}" ) self._asyncioEventloop: AbstractEventLoop = _eventloop self._writers: Dict[Type[FileDescriptor], int] = {} self._readers: Dict[Type[FileDescriptor], int] = {} self._continuousPolling = _ContinuousPolling(self) self._scheduledAt = None self._timerHandle = None super().__init__() def _unregisterFDInAsyncio(self, fd): """ Compensate for a bug in asyncio where it will not unregister a FD that it cannot handle in the epoll loop. It touches internal asyncio code. A description of the bug by markrwilliams: The C{add_writer} method of asyncio event loops isn't atomic because all the Selector classes in the selector module internally record a file object before passing it to the platform's selector implementation. If the platform's selector decides the file object isn't acceptable, the resulting exception doesn't cause the Selector to un-track the file object. The failing/hanging stdio test goes through the following sequence of events (roughly): * The first C{connection.write(intToByte(value))} call hits the asyncio reactor's C{addWriter} method. * C{addWriter} calls the asyncio loop's C{add_writer} method, which happens to live on C{_BaseSelectorEventLoop}. * The asyncio loop's C{add_writer} method checks if the file object has been registered before via the selector's C{get_key} method. * It hasn't, so the KeyError block runs and calls the selector's register method * Code examples that follow use EpollSelector, but the code flow holds true for any other selector implementation. The selector's register method first calls through to the next register method in the MRO * That next method is always C{_BaseSelectorImpl.register} which creates a C{SelectorKey} instance for the file object, stores it under the file object's file descriptor, and then returns it. * Control returns to the concrete selector implementation, which asks the operating system to track the file descriptor using the right API. * The operating system refuses! An exception is raised that, in this case, the asyncio reactor handles by creating a C{_ContinuousPolling} object to watch the file descriptor. * The second C{connection.write(intToByte(value))} call hits the asyncio reactor's C{addWriter} method, which hits the C{add_writer} method. But the loop's selector's get_key method now returns a C{SelectorKey}! Now the asyncio reactor's C{addWriter} method thinks the asyncio loop will watch the file descriptor, even though it won't. """ try: self._asyncioEventloop._selector.unregister(fd) except BaseException: pass def _readOrWrite(self, selectable, read): method = selectable.doRead if read else selectable.doWrite if selectable.fileno() == -1: self._disconnectSelectable(selectable, _NO_FILEDESC, read) return try: why = method() except Exception as e: why = e self._log.failure(None) if why: self._disconnectSelectable(selectable, why, read) def addReader(self, reader): if reader in self._readers.keys() or reader in self._continuousPolling._readers: return fd = reader.fileno() try: self._asyncioEventloop.add_reader( fd, callWithLogger, reader, self._readOrWrite, reader, True ) self._readers[reader] = fd except OSError as e: self._unregisterFDInAsyncio(fd) if e.errno == errno.EPERM: # epoll(7) doesn't support certain file descriptors, # e.g. filesystem files, so for those we just poll # continuously: self._continuousPolling.addReader(reader) else: raise def addWriter(self, writer): if writer in self._writers.keys() or writer in self._continuousPolling._writers: return fd = writer.fileno() try: self._asyncioEventloop.add_writer( fd, callWithLogger, writer, self._readOrWrite, writer, False ) self._writers[writer] = fd except PermissionError: self._unregisterFDInAsyncio(fd) # epoll(7) doesn't support certain file descriptors, # e.g. filesystem files, so for those we just poll # continuously: self._continuousPolling.addWriter(writer) except BrokenPipeError: # The kqueuereactor will raise this if there is a broken pipe self._unregisterFDInAsyncio(fd) except BaseException: self._unregisterFDInAsyncio(fd) raise def removeReader(self, reader): # First, see if they're trying to remove a reader that we don't have. if not ( reader in self._readers.keys() or self._continuousPolling.isReading(reader) ): # We don't have it, so just return OK. return # If it was a cont. polling reader, check there first. if self._continuousPolling.isReading(reader): self._continuousPolling.removeReader(reader) return fd = reader.fileno() if fd == -1: # If the FD is -1, we want to know what its original FD was, to # remove it. fd = self._readers.pop(reader) else: self._readers.pop(reader) self._asyncioEventloop.remove_reader(fd) def removeWriter(self, writer): # First, see if they're trying to remove a writer that we don't have. if not ( writer in self._writers.keys() or self._continuousPolling.isWriting(writer) ): # We don't have it, so just return OK. return # If it was a cont. polling writer, check there first. if self._continuousPolling.isWriting(writer): self._continuousPolling.removeWriter(writer) return fd = writer.fileno() if fd == -1: # If the FD is -1, we want to know what its original FD was, to # remove it. fd = self._writers.pop(writer) else: self._writers.pop(writer) self._asyncioEventloop.remove_writer(fd) def removeAll(self): return ( self._removeAll(self._readers.keys(), self._writers.keys()) + self._continuousPolling.removeAll() ) def getReaders(self): return list(self._readers.keys()) + self._continuousPolling.getReaders() def getWriters(self): return list(self._writers.keys()) + self._continuousPolling.getWriters() def iterate(self, timeout): self._asyncioEventloop.call_later(timeout + 0.01, self._asyncioEventloop.stop) self._asyncioEventloop.run_forever() def run(self, installSignalHandlers=True): self.startRunning(installSignalHandlers=installSignalHandlers) self._asyncioEventloop.run_forever() if self._justStopped: self._justStopped = False def stop(self): super().stop() # This will cause runUntilCurrent which in its turn # will call fireSystemEvent("shutdown") self.callLater(0, lambda: None) def crash(self): super().crash() self._asyncioEventloop.stop() def _onTimer(self): self._scheduledAt = None self.runUntilCurrent() self._reschedule() def _reschedule(self): timeout = self.timeout() if timeout is not None: abs_time = self._asyncioEventloop.time() + timeout self._scheduledAt = abs_time if self._timerHandle is not None: self._timerHandle.cancel() self._timerHandle = self._asyncioEventloop.call_at(abs_time, self._onTimer) def _moveCallLaterSooner(self, tple): PosixReactorBase._moveCallLaterSooner(self, tple) self._reschedule() def callLater(self, seconds, f, *args, **kwargs): dc = PosixReactorBase.callLater(self, seconds, f, *args, **kwargs) abs_time = self._asyncioEventloop.time() + self.timeout() if self._scheduledAt is None or abs_time < self._scheduledAt: self._reschedule() return dc def callFromThread(self, f, *args, **kwargs): g = lambda: self.callLater(0, f, *args, **kwargs) self._asyncioEventloop.call_soon_threadsafe(g) def install(eventloop=None): """ Install an asyncio-based reactor. @param eventloop: The asyncio eventloop to wrap. If default, the global one is selected. """ reactor = AsyncioSelectorReactor(eventloop) from twisted.internet.main import installReactor installReactor(reactor)
Save
cmd:
run