Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
tpool.py344 linesDownload Raw Back to eventlet
1# Copyright (c) 2007-2009, Linden Research, Inc.2# Copyright (c) 2007, IBM Corp.3#4# Licensed under the Apache License, Version 2.0 (the "License");5# you may not use this file except in compliance with the License.6# You may obtain a copy of the License at7#8#   http://www.apache.org/licenses/LICENSE-2.09#10# Unless required by applicable law or agreed to in writing, software11# distributed under the License is distributed on an "AS IS" BASIS,12# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.13# See the License for the specific language governing permissions and14# limitations under the License.15 16import atexit17try:18    import _imp as imp19except ImportError:20    import imp21import os22import sys23import traceback24 25import eventlet26from eventlet import event, greenio, greenthread, patcher, timeout27import six28 29__all__ = ['execute', 'Proxy', 'killall', 'set_num_threads']30 31 32EXC_CLASSES = (Exception, timeout.Timeout)33SYS_EXCS = (GeneratorExit, KeyboardInterrupt, SystemExit)34 35QUIET = True36 37socket = patcher.original('socket')38threading = patcher.original('threading')39if six.PY2:40    Queue_module = patcher.original('Queue')41if six.PY3:42    Queue_module = patcher.original('queue')43 44Empty = Queue_module.Empty45Queue = Queue_module.Queue46 47_bytetosend = b' '48_coro = None49_nthreads = int(os.environ.get('EVENTLET_THREADPOOL_SIZE', 20))50_reqq = _rspq = None51_rsock = _wsock = None52_setup_already = False53_threads = []54 55 56def tpool_trampoline():57    global _rspq58    while True:59        try:60            _c = _rsock.recv(1)61            assert _c62        # FIXME: this is probably redundant since using sockets instead of pipe now63        except ValueError:64            break  # will be raised when pipe is closed65        while not _rspq.empty():66            try:67                (e, rv) = _rspq.get(block=False)68                e.send(rv)69                e = rv = None70            except Empty:71                pass72 73 74def tworker():75    global _rspq76    while True:77        try:78            msg = _reqq.get()79        except AttributeError:80            return  # can't get anything off of a dud queue81        if msg is None:82            return83        (e, meth, args, kwargs) = msg84        rv = None85        try:86            rv = meth(*args, **kwargs)87        except SYS_EXCS:88            raise89        except EXC_CLASSES:90            rv = sys.exc_info()91            if sys.version_info >= (3, 4):92                traceback.clear_frames(rv[1].__traceback__)93        if six.PY2:94            sys.exc_clear()95        # test_leakage_from_tracebacks verifies that the use of96        # exc_info does not lead to memory leaks97        _rspq.put((e, rv))98        msg = meth = args = kwargs = e = rv = None99        _wsock.sendall(_bytetosend)100 101 102def execute(meth, *args, **kwargs):103    """104    Execute *meth* in a Python thread, blocking the current coroutine/105    greenthread until the method completes.106 107    The primary use case for this is to wrap an object or module that is not108    amenable to monkeypatching or any of the other tricks that Eventlet uses109    to achieve cooperative yielding.  With tpool, you can force such objects to110    cooperate with green threads by sticking them in native threads, at the cost111    of some overhead.112    """113    setup()114    # if already in tpool, don't recurse into the tpool115    # also, call functions directly if we're inside an import lock, because116    # if meth does any importing (sadly common), it will hang117    my_thread = threading.current_thread()118    if my_thread in _threads or imp.lock_held() or _nthreads == 0:119        return meth(*args, **kwargs)120 121    e = event.Event()122    _reqq.put((e, meth, args, kwargs))123 124    rv = e.wait()125    if isinstance(rv, tuple) \126            and len(rv) == 3 \127            and isinstance(rv[1], EXC_CLASSES):128        (c, e, tb) = rv129        if not QUIET:130            traceback.print_exception(c, e, tb)131            traceback.print_stack()132        six.reraise(c, e, tb)133    return rv134 135 136def proxy_call(autowrap, f, *args, **kwargs):137    """138    Call a function *f* and returns the value.  If the type of the return value139    is in the *autowrap* collection, then it is wrapped in a :class:`Proxy`140    object before return.141 142    Normally *f* will be called in the threadpool with :func:`execute`; if the143    keyword argument "nonblocking" is set to ``True``, it will simply be144    executed directly.  This is useful if you have an object which has methods145    that don't need to be called in a separate thread, but which return objects146    that should be Proxy wrapped.147    """148    if kwargs.pop('nonblocking', False):149        rv = f(*args, **kwargs)150    else:151        rv = execute(f, *args, **kwargs)152    if isinstance(rv, autowrap):153        return Proxy(rv, autowrap)154    else:155        return rv156 157 158class Proxy(object):159    """160    a simple proxy-wrapper of any object that comes with a161    methods-only interface, in order to forward every method162    invocation onto a thread in the native-thread pool.  A key163    restriction is that the object's methods should not switch164    greenlets or use Eventlet primitives, since they are in a165    different thread from the main hub, and therefore might behave166    unexpectedly.  This is for running native-threaded code167    only.168 169    It's common to want to have some of the attributes or return170    values also wrapped in Proxy objects (for example, database171    connection objects produce cursor objects which also should be172    wrapped in Proxy objects to remain nonblocking).  *autowrap*, if173    supplied, is a collection of types; if an attribute or return174    value matches one of those types (via isinstance), it will be175    wrapped in a Proxy.  *autowrap_names* is a collection176    of strings, which represent the names of attributes that should be177    wrapped in Proxy objects when accessed.178    """179 180    def __init__(self, obj, autowrap=(), autowrap_names=()):181        self._obj = obj182        self._autowrap = autowrap183        self._autowrap_names = autowrap_names184 185    def __getattr__(self, attr_name):186        f = getattr(self._obj, attr_name)187        if not hasattr(f, '__call__'):188            if isinstance(f, self._autowrap) or attr_name in self._autowrap_names:189                return Proxy(f, self._autowrap)190            return f191 192        def doit(*args, **kwargs):193            result = proxy_call(self._autowrap, f, *args, **kwargs)194            if attr_name in self._autowrap_names and not isinstance(result, Proxy):195                return Proxy(result)196            return result197        return doit198 199    # the following are a buncha methods that the python interpeter200    # doesn't use getattr to retrieve and therefore have to be defined201    # explicitly202    def __getitem__(self, key):203        return proxy_call(self._autowrap, self._obj.__getitem__, key)204 205    def __setitem__(self, key, value):206        return proxy_call(self._autowrap, self._obj.__setitem__, key, value)207 208    def __deepcopy__(self, memo=None):209        return proxy_call(self._autowrap, self._obj.__deepcopy__, memo)210 211    def __copy__(self, memo=None):212        return proxy_call(self._autowrap, self._obj.__copy__, memo)213 214    def __call__(self, *a, **kw):215        if '__call__' in self._autowrap_names:216            return Proxy(proxy_call(self._autowrap, self._obj, *a, **kw))217        else:218            return proxy_call(self._autowrap, self._obj, *a, **kw)219 220    def __enter__(self):221        return proxy_call(self._autowrap, self._obj.__enter__)222 223    def __exit__(self, *exc):224        return proxy_call(self._autowrap, self._obj.__exit__, *exc)225 226    # these don't go through a proxy call, because they're likely to227    # be called often, and are unlikely to be implemented on the228    # wrapped object in such a way that they would block229    def __eq__(self, rhs):230        return self._obj == rhs231 232    def __hash__(self):233        return self._obj.__hash__()234 235    def __repr__(self):236        return self._obj.__repr__()237 238    def __str__(self):239        return self._obj.__str__()240 241    def __len__(self):242        return len(self._obj)243 244    def __nonzero__(self):245        return bool(self._obj)246    # Python3247    __bool__ = __nonzero__248 249    def __iter__(self):250        it = iter(self._obj)251        if it == self._obj:252            return self253        else:254            return Proxy(it)255 256    def next(self):257        return proxy_call(self._autowrap, next, self._obj)258    # Python3259    __next__ = next260 261 262def setup():263    global _rsock, _wsock, _coro, _setup_already, _rspq, _reqq264    if _setup_already:265        return266    else:267        _setup_already = True268 269    assert _nthreads >= 0, "Can't specify negative number of threads"270    if _nthreads == 0:271        import warnings272        warnings.warn("Zero threads in tpool.  All tpool.execute calls will\273            execute in main thread.  Check the value of the environment \274            variable EVENTLET_THREADPOOL_SIZE.", RuntimeWarning)275    _reqq = Queue(maxsize=-1)276    _rspq = Queue(maxsize=-1)277 278    # connected socket pair279    sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)280    sock.bind(('127.0.0.1', 0))281    sock.listen(1)282    csock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)283    csock.connect(sock.getsockname())284    csock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, True)285    _wsock, _addr = sock.accept()286    _wsock.settimeout(None)287    _wsock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, True)288    sock.close()289    _rsock = greenio.GreenSocket(csock)290    _rsock.settimeout(None)291 292    for i in six.moves.range(_nthreads):293        t = threading.Thread(target=tworker,294                             name="tpool_thread_%s" % i)295        t.daemon = True296        t.start()297        _threads.append(t)298 299    _coro = greenthread.spawn_n(tpool_trampoline)300    # This yield fixes subtle error with GreenSocket.__del__301    eventlet.sleep(0)302 303 304# Avoid ResourceWarning unclosed socket on Python3.2+305@atexit.register306def killall():307    global _setup_already, _rspq, _rsock, _wsock308    if not _setup_already:309        return310 311    # This yield fixes freeze in some scenarios312    eventlet.sleep(0)313 314    for thr in _threads:315        _reqq.put(None)316    for thr in _threads:317        thr.join()318    del _threads[:]319 320    # return any remaining results321    while (_rspq is not None) and not _rspq.empty():322        try:323            (e, rv) = _rspq.get(block=False)324            e.send(rv)325            e = rv = None326        except Empty:327            pass328 329    if _coro is not None:330        greenthread.kill(_coro)331    if _rsock is not None:332        _rsock.close()333        _rsock = None334    if _wsock is not None:335        _wsock.close()336        _wsock = None337    _rspq = None338    _setup_already = False339 340 341def set_num_threads(nthreads):342    global _nthreads343    _nthreads = nthreads344 
codekingpro/portable-devtools · Team Ai