Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
queue.py493 linesDownload Raw Back to eventlet
1# Copyright (c) 2009 Denis Bilenko, denis.bilenko at gmail com2# Copyright (c) 2010 Eventlet Contributors (see AUTHORS)3# and licensed under the MIT license:4#5# Permission is hereby granted, free of charge, to any person obtaining a copy6# of this software and associated documentation files (the "Software"), to deal7# in the Software without restriction, including without limitation the rights8# to use, copy, modify, merge, publish, distribute, sublicense, and/or sell9# copies of the Software, and to permit persons to whom the Software is10# furnished to do so, subject to the following conditions:11#12# The above copyright notice and this permission notice shall be included in13# all copies or substantial portions of the Software.14#15# THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR16# IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,17# FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE18# AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER19# LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,20# OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN21# THE SOFTWARE.22 23"""Synchronized queues.24 25The :mod:`eventlet.queue` module implements multi-producer, multi-consumer26queues that work across greenlets, with the API similar to the classes found in27the standard :mod:`Queue` and :class:`multiprocessing <multiprocessing.Queue>`28modules.29 30A major difference is that queues in this module operate as channels when31initialized with *maxsize* of zero. In such case, both :meth:`Queue.empty`32and :meth:`Queue.full` return ``True`` and :meth:`Queue.put` always blocks until33a call to :meth:`Queue.get` retrieves the item.34 35An interesting difference, made possible because of greenthreads, is36that :meth:`Queue.qsize`, :meth:`Queue.empty`, and :meth:`Queue.full` *can* be37used as indicators of whether the subsequent :meth:`Queue.get`38or :meth:`Queue.put` will not block.  The new methods :meth:`Queue.getting`39and :meth:`Queue.putting` report on the number of greenthreads blocking40in :meth:`put <Queue.put>` or :meth:`get <Queue.get>` respectively.41"""42from __future__ import print_function43 44import sys45import heapq46import collections47import traceback48 49from eventlet.event import Event50from eventlet.greenthread import getcurrent51from eventlet.hubs import get_hub52import six53from six.moves import queue as Stdlib_Queue54from eventlet.timeout import Timeout55 56 57__all__ = ['Queue', 'PriorityQueue', 'LifoQueue', 'LightQueue', 'Full', 'Empty']58 59_NONE = object()60Full = six.moves.queue.Full61Empty = six.moves.queue.Empty62 63 64class Waiter(object):65    """A low level synchronization class.66 67    Wrapper around greenlet's ``switch()`` and ``throw()`` calls that makes them safe:68 69    * switching will occur only if the waiting greenlet is executing :meth:`wait`70      method currently. Otherwise, :meth:`switch` and :meth:`throw` are no-ops.71    * any error raised in the greenlet is handled inside :meth:`switch` and :meth:`throw`72 73    The :meth:`switch` and :meth:`throw` methods must only be called from the :class:`Hub` greenlet.74    The :meth:`wait` method must be called from a greenlet other than :class:`Hub`.75    """76    __slots__ = ['greenlet']77 78    def __init__(self):79        self.greenlet = None80 81    def __repr__(self):82        if self.waiting:83            waiting = ' waiting'84        else:85            waiting = ''86        return '<%s at %s%s greenlet=%r>' % (87            type(self).__name__, hex(id(self)), waiting, self.greenlet,88        )89 90    def __str__(self):91        """92        >>> print(Waiter())93        <Waiter greenlet=None>94        """95        if self.waiting:96            waiting = ' waiting'97        else:98            waiting = ''99        return '<%s%s greenlet=%s>' % (type(self).__name__, waiting, self.greenlet)100 101    def __nonzero__(self):102        return self.greenlet is not None103 104    __bool__ = __nonzero__105 106    @property107    def waiting(self):108        return self.greenlet is not None109 110    def switch(self, value=None):111        """Wake up the greenlet that is calling wait() currently (if there is one).112        Can only be called from Hub's greenlet.113        """114        assert getcurrent() is get_hub(115        ).greenlet, "Can only use Waiter.switch method from the mainloop"116        if self.greenlet is not None:117            try:118                self.greenlet.switch(value)119            except Exception:120                traceback.print_exc()121 122    def throw(self, *throw_args):123        """Make greenlet calling wait() wake up (if there is a wait()).124        Can only be called from Hub's greenlet.125        """126        assert getcurrent() is get_hub(127        ).greenlet, "Can only use Waiter.switch method from the mainloop"128        if self.greenlet is not None:129            try:130                self.greenlet.throw(*throw_args)131            except Exception:132                traceback.print_exc()133 134    # XXX should be renamed to get() ? and the whole class is called Receiver?135    def wait(self):136        """Wait until switch() or throw() is called.137        """138        assert self.greenlet is None, 'This Waiter is already used by %r' % (self.greenlet, )139        self.greenlet = getcurrent()140        try:141            return get_hub().switch()142        finally:143            self.greenlet = None144 145 146class LightQueue(object):147    """148    This is a variant of Queue that behaves mostly like the standard149    :class:`Stdlib_Queue`.  It differs by not supporting the150    :meth:`task_done <Stdlib_Queue.task_done>` or151    :meth:`join <Stdlib_Queue.join>` methods, and is a little faster for152    not having that overhead.153    """154 155    def __init__(self, maxsize=None):156        if maxsize is None or maxsize < 0:  # None is not comparable in 3.x157            self.maxsize = None158        else:159            self.maxsize = maxsize160        self.getters = set()161        self.putters = set()162        self._event_unlock = None163        self._init(maxsize)164 165    # QQQ make maxsize into a property with setter that schedules unlock if necessary166 167    def _init(self, maxsize):168        self.queue = collections.deque()169 170    def _get(self):171        return self.queue.popleft()172 173    def _put(self, item):174        self.queue.append(item)175 176    def __repr__(self):177        return '<%s at %s %s>' % (type(self).__name__, hex(id(self)), self._format())178 179    def __str__(self):180        return '<%s %s>' % (type(self).__name__, self._format())181 182    def _format(self):183        result = 'maxsize=%r' % (self.maxsize, )184        if getattr(self, 'queue', None):185            result += ' queue=%r' % self.queue186        if self.getters:187            result += ' getters[%s]' % len(self.getters)188        if self.putters:189            result += ' putters[%s]' % len(self.putters)190        if self._event_unlock is not None:191            result += ' unlocking'192        return result193 194    def qsize(self):195        """Return the size of the queue."""196        return len(self.queue)197 198    def resize(self, size):199        """Resizes the queue's maximum size.200 201        If the size is increased, and there are putters waiting, they may be woken up."""202        # None is not comparable in 3.x203        if self.maxsize is not None and (size is None or size > self.maxsize):204            # Maybe wake some stuff up205            self._schedule_unlock()206        self.maxsize = size207 208    def putting(self):209        """Returns the number of greenthreads that are blocked waiting to put210        items into the queue."""211        return len(self.putters)212 213    def getting(self):214        """Returns the number of greenthreads that are blocked waiting on an215        empty queue."""216        return len(self.getters)217 218    def empty(self):219        """Return ``True`` if the queue is empty, ``False`` otherwise."""220        return not self.qsize()221 222    def full(self):223        """Return ``True`` if the queue is full, ``False`` otherwise.224 225        ``Queue(None)`` is never full.226        """227        # None is not comparable in 3.x228        return self.maxsize is not None and self.qsize() >= self.maxsize229 230    def put(self, item, block=True, timeout=None):231        """Put an item into the queue.232 233        If optional arg *block* is true and *timeout* is ``None`` (the default),234        block if necessary until a free slot is available. If *timeout* is235        a positive number, it blocks at most *timeout* seconds and raises236        the :class:`Full` exception if no free slot was available within that time.237        Otherwise (*block* is false), put an item on the queue if a free slot238        is immediately available, else raise the :class:`Full` exception (*timeout*239        is ignored in that case).240        """241        if self.maxsize is None or self.qsize() < self.maxsize:242            # there's a free slot, put an item right away243            self._put(item)244            if self.getters:245                self._schedule_unlock()246        elif not block and get_hub().greenlet is getcurrent():247            # we're in the mainloop, so we cannot wait; we can switch() to other greenlets though248            # find a getter and deliver an item to it249            while self.getters:250                getter = self.getters.pop()251                if getter:252                    self._put(item)253                    item = self._get()254                    getter.switch(item)255                    return256            raise Full257        elif block:258            waiter = ItemWaiter(item, block)259            self.putters.add(waiter)260            timeout = Timeout(timeout, Full)261            try:262                if self.getters:263                    self._schedule_unlock()264                result = waiter.wait()265                assert result is waiter, "Invalid switch into Queue.put: %r" % (result, )266                if waiter.item is not _NONE:267                    self._put(item)268            finally:269                timeout.cancel()270                self.putters.discard(waiter)271        elif self.getters:272            waiter = ItemWaiter(item, block)273            self.putters.add(waiter)274            self._schedule_unlock()275            result = waiter.wait()276            assert result is waiter, "Invalid switch into Queue.put: %r" % (result, )277            if waiter.item is not _NONE:278                raise Full279        else:280            raise Full281 282    def put_nowait(self, item):283        """Put an item into the queue without blocking.284 285        Only enqueue the item if a free slot is immediately available.286        Otherwise raise the :class:`Full` exception.287        """288        self.put(item, False)289 290    def get(self, block=True, timeout=None):291        """Remove and return an item from the queue.292 293        If optional args *block* is true and *timeout* is ``None`` (the default),294        block if necessary until an item is available. If *timeout* is a positive number,295        it blocks at most *timeout* seconds and raises the :class:`Empty` exception296        if no item was available within that time. Otherwise (*block* is false), return297        an item if one is immediately available, else raise the :class:`Empty` exception298        (*timeout* is ignored in that case).299        """300        if self.qsize():301            if self.putters:302                self._schedule_unlock()303            return self._get()304        elif not block and get_hub().greenlet is getcurrent():305            # special case to make get_nowait() runnable in the mainloop greenlet306            # there are no items in the queue; try to fix the situation by unlocking putters307            while self.putters:308                putter = self.putters.pop()309                if putter:310                    putter.switch(putter)311                    if self.qsize():312                        return self._get()313            raise Empty314        elif block:315            waiter = Waiter()316            timeout = Timeout(timeout, Empty)317            try:318                self.getters.add(waiter)319                if self.putters:320                    self._schedule_unlock()321                try:322                    return waiter.wait()323                except:324                    self._schedule_unlock()325                    raise326            finally:327                self.getters.discard(waiter)328                timeout.cancel()329        else:330            raise Empty331 332    def get_nowait(self):333        """Remove and return an item from the queue without blocking.334 335        Only get an item if one is immediately available. Otherwise336        raise the :class:`Empty` exception.337        """338        return self.get(False)339 340    def _unlock(self):341        try:342            while True:343                if self.qsize() and self.getters:344                    getter = self.getters.pop()345                    if getter:346                        try:347                            item = self._get()348                        except:349                            getter.throw(*sys.exc_info())350                        else:351                            getter.switch(item)352                elif self.putters and self.getters:353                    putter = self.putters.pop()354                    if putter:355                        getter = self.getters.pop()356                        if getter:357                            item = putter.item358                            # this makes greenlet calling put() not to call _put() again359                            putter.item = _NONE360                            self._put(item)361                            item = self._get()362                            getter.switch(item)363                            putter.switch(putter)364                        else:365                            self.putters.add(putter)366                elif self.putters and (self.getters or367                                       self.maxsize is None or368                                       self.qsize() < self.maxsize):369                    putter = self.putters.pop()370                    putter.switch(putter)371                elif self.putters and not self.getters:372                    full = [p for p in self.putters if not p.block]373                    if not full:374                        break375                    for putter in full:376                        self.putters.discard(putter)377                        get_hub().schedule_call_global(378                            0, putter.greenlet.throw, Full)379                else:380                    break381        finally:382            self._event_unlock = None  # QQQ maybe it's possible to obtain this info from libevent?383            # i.e. whether this event is pending _OR_ currently executing384        # testcase: 2 greenlets: while True: q.put(q.get()) - nothing else has a change to execute385        # to avoid this, schedule unlock with timer(0, ...) once in a while386 387    def _schedule_unlock(self):388        if self._event_unlock is None:389            self._event_unlock = get_hub().schedule_call_global(0, self._unlock)390 391 392class ItemWaiter(Waiter):393    __slots__ = ['item', 'block']394 395    def __init__(self, item, block):396        Waiter.__init__(self)397        self.item = item398        self.block = block399 400 401class Queue(LightQueue):402    '''Create a queue object with a given maximum size.403 404    If *maxsize* is less than zero or ``None``, the queue size is infinite.405 406    ``Queue(0)`` is a channel, that is, its :meth:`put` method always blocks407    until the item is delivered. (This is unlike the standard408    :class:`Stdlib_Queue`, where 0 means infinite size).409 410    In all other respects, this Queue class resembles the standard library,411    :class:`Stdlib_Queue`.412    '''413 414    def __init__(self, maxsize=None):415        LightQueue.__init__(self, maxsize)416        self.unfinished_tasks = 0417        self._cond = Event()418 419    def _format(self):420        result = LightQueue._format(self)421        if self.unfinished_tasks:422            result += ' tasks=%s _cond=%s' % (self.unfinished_tasks, self._cond)423        return result424 425    def _put(self, item):426        LightQueue._put(self, item)427        self._put_bookkeeping()428 429    def _put_bookkeeping(self):430        self.unfinished_tasks += 1431        if self._cond.ready():432            self._cond.reset()433 434    def task_done(self):435        '''Indicate that a formerly enqueued task is complete. Used by queue consumer threads.436        For each :meth:`get <Queue.get>` used to fetch a task, a subsequent call to437        :meth:`task_done` tells the queue that the processing on the task is complete.438 439        If a :meth:`join` is currently blocking, it will resume when all items have been processed440        (meaning that a :meth:`task_done` call was received for every item that had been441        :meth:`put <Queue.put>` into the queue).442 443        Raises a :exc:`ValueError` if called more times than there were items placed in the queue.444        '''445 446        if self.unfinished_tasks <= 0:447            raise ValueError('task_done() called too many times')448        self.unfinished_tasks -= 1449        if self.unfinished_tasks == 0:450            self._cond.send(None)451 452    def join(self):453        '''Block until all items in the queue have been gotten and processed.454 455        The count of unfinished tasks goes up whenever an item is added to the queue.456        The count goes down whenever a consumer thread calls :meth:`task_done` to indicate457        that the item was retrieved and all work on it is complete. When the count of458        unfinished tasks drops to zero, :meth:`join` unblocks.459        '''460        if self.unfinished_tasks > 0:461            self._cond.wait()462 463 464class PriorityQueue(Queue):465    '''A subclass of :class:`Queue` that retrieves entries in priority order (lowest first).466 467    Entries are typically tuples of the form: ``(priority number, data)``.468    '''469 470    def _init(self, maxsize):471        self.queue = []472 473    def _put(self, item, heappush=heapq.heappush):474        heappush(self.queue, item)475        self._put_bookkeeping()476 477    def _get(self, heappop=heapq.heappop):478        return heappop(self.queue)479 480 481class LifoQueue(Queue):482    '''A subclass of :class:`Queue` that retrieves most recently added entries first.'''483 484    def _init(self, maxsize):485        self.queue = []486 487    def _put(self, item):488        self.queue.append(item)489        self._put_bookkeeping()490 491    def _get(self):492        return self.queue.pop()493 
codekingpro/portable-devtools · Team Ai