Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
queue.py384 linesDownload Raw Back to Lib
1'''A multi-producer, multi-consumer queue.'''2 3import threading4import types5from collections import deque6from heapq import heappush, heappop7from time import monotonic as time8try:9    from _queue import SimpleQueue10except ImportError:11    SimpleQueue = None12 13__all__ = [14    'Empty',15    'Full',16    'ShutDown',17    'Queue',18    'PriorityQueue',19    'LifoQueue',20    'SimpleQueue',21]22 23 24try:25    from _queue import Empty26except ImportError:27    class Empty(Exception):28        'Exception raised by Queue.get(block=0)/get_nowait().'29        pass30 31class Full(Exception):32    'Exception raised by Queue.put(block=0)/put_nowait().'33    pass34 35 36class ShutDown(Exception):37    '''Raised when put/get with shut-down queue.'''38 39 40class Queue:41    '''Create a queue object with a given maximum size.42 43    If maxsize is <= 0, the queue size is infinite.44    '''45 46    def __init__(self, maxsize=0):47        self.maxsize = maxsize48        self._init(maxsize)49 50        # mutex must be held whenever the queue is mutating.  All methods51        # that acquire mutex must release it before returning.  mutex52        # is shared between the three conditions, so acquiring and53        # releasing the conditions also acquires and releases mutex.54        self.mutex = threading.Lock()55 56        # Notify not_empty whenever an item is added to the queue; a57        # thread waiting to get is notified then.58        self.not_empty = threading.Condition(self.mutex)59 60        # Notify not_full whenever an item is removed from the queue;61        # a thread waiting to put is notified then.62        self.not_full = threading.Condition(self.mutex)63 64        # Notify all_tasks_done whenever the number of unfinished tasks65        # drops to zero; thread waiting to join() is notified to resume66        self.all_tasks_done = threading.Condition(self.mutex)67        self.unfinished_tasks = 068 69        # Queue shutdown state70        self.is_shutdown = False71 72    def task_done(self):73        '''Indicate that a formerly enqueued task is complete.74 75        Used by Queue consumer threads.  For each get() used to fetch a task,76        a subsequent call to task_done() tells the queue that the processing77        on the task is complete.78 79        If a join() is currently blocking, it will resume when all items80        have been processed (meaning that a task_done() call was received81        for every item that had been put() into the queue).82 83        Raises a ValueError if called more times than there were items84        placed in the queue.85        '''86        with self.all_tasks_done:87            unfinished = self.unfinished_tasks - 188            if unfinished <= 0:89                if unfinished < 0:90                    raise ValueError('task_done() called too many times')91                self.all_tasks_done.notify_all()92            self.unfinished_tasks = unfinished93 94    def join(self):95        '''Blocks until all items in the Queue have been gotten and processed.96 97        The count of unfinished tasks goes up whenever an item is added to the98        queue. The count goes down whenever a consumer thread calls task_done()99        to indicate the item was retrieved and all work on it is complete.100 101        When the count of unfinished tasks drops to zero, join() unblocks.102        '''103        with self.all_tasks_done:104            while self.unfinished_tasks:105                self.all_tasks_done.wait()106 107    def qsize(self):108        '''Return the approximate size of the queue (not reliable!).'''109        with self.mutex:110            return self._qsize()111 112    def empty(self):113        '''Return True if the queue is empty, False otherwise (not reliable!).114 115        This method is likely to be removed at some point.  Use qsize() == 0116        as a direct substitute, but be aware that either approach risks a race117        condition where a queue can grow before the result of empty() or118        qsize() can be used.119 120        To create code that needs to wait for all queued tasks to be121        completed, the preferred technique is to use the join() method.122        '''123        with self.mutex:124            return not self._qsize()125 126    def full(self):127        '''Return True if the queue is full, False otherwise (not reliable!).128 129        This method is likely to be removed at some point.  Use qsize() >= n130        as a direct substitute, but be aware that either approach risks a race131        condition where a queue can shrink before the result of full() or132        qsize() can be used.133        '''134        with self.mutex:135            return 0 < self.maxsize <= self._qsize()136 137    def put(self, item, block=True, timeout=None):138        '''Put an item into the queue.139 140        If optional args 'block' is true and 'timeout' is None (the default),141        block if necessary until a free slot is available. If 'timeout' is142        a non-negative number, it blocks at most 'timeout' seconds and raises143        the Full exception if no free slot was available within that time.144        Otherwise ('block' is false), put an item on the queue if a free slot145        is immediately available, else raise the Full exception ('timeout'146        is ignored in that case).147 148        Raises ShutDown if the queue has been shut down.149        '''150        with self.not_full:151            if self.is_shutdown:152                raise ShutDown153            if self.maxsize > 0:154                if not block:155                    if self._qsize() >= self.maxsize:156                        raise Full157                elif timeout is None:158                    while self._qsize() >= self.maxsize:159                        self.not_full.wait()160                        if self.is_shutdown:161                            raise ShutDown162                elif timeout < 0:163                    raise ValueError("'timeout' must be a non-negative number")164                else:165                    endtime = time() + timeout166                    while self._qsize() >= self.maxsize:167                        remaining = endtime - time()168                        if remaining <= 0.0:169                            raise Full170                        self.not_full.wait(remaining)171                        if self.is_shutdown:172                            raise ShutDown173            self._put(item)174            self.unfinished_tasks += 1175            self.not_empty.notify()176 177    def get(self, block=True, timeout=None):178        '''Remove and return an item from the queue.179 180        If optional args 'block' is true and 'timeout' is None (the default),181        block if necessary until an item is available. If 'timeout' is182        a non-negative number, it blocks at most 'timeout' seconds and raises183        the Empty exception if no item was available within that time.184        Otherwise ('block' is false), return an item if one is immediately185        available, else raise the Empty exception ('timeout' is ignored186        in that case).187 188        Raises ShutDown if the queue has been shut down and is empty,189        or if the queue has been shut down immediately.190        '''191        with self.not_empty:192            if self.is_shutdown and not self._qsize():193                raise ShutDown194            if not block:195                if not self._qsize():196                    raise Empty197            elif timeout is None:198                while not self._qsize():199                    self.not_empty.wait()200                    if self.is_shutdown and not self._qsize():201                        raise ShutDown202            elif timeout < 0:203                raise ValueError("'timeout' must be a non-negative number")204            else:205                endtime = time() + timeout206                while not self._qsize():207                    remaining = endtime - time()208                    if remaining <= 0.0:209                        raise Empty210                    self.not_empty.wait(remaining)211                    if self.is_shutdown and not self._qsize():212                        raise ShutDown213            item = self._get()214            self.not_full.notify()215            return item216 217    def put_nowait(self, item):218        '''Put an item into the queue without blocking.219 220        Only enqueue the item if a free slot is immediately available.221        Otherwise raise the Full exception.222        '''223        return self.put(item, block=False)224 225    def get_nowait(self):226        '''Remove and return an item from the queue without blocking.227 228        Only get an item if one is immediately available. Otherwise229        raise the Empty exception.230        '''231        return self.get(block=False)232 233    def shutdown(self, immediate=False):234        '''Shut-down the queue, making queue gets and puts raise ShutDown.235 236        By default, gets will only raise once the queue is empty. Set237        'immediate' to True to make gets raise immediately instead.238 239        All blocked callers of put() and get() will be unblocked.240 241        If 'immediate', the queue is drained and unfinished tasks242        is reduced by the number of drained tasks.  If unfinished tasks243        is reduced to zero, callers of Queue.join are unblocked.244        '''245        with self.mutex:246            self.is_shutdown = True247            if immediate:248                while self._qsize():249                    self._get()250                    if self.unfinished_tasks > 0:251                        self.unfinished_tasks -= 1252                # release all blocked threads in `join()`253                self.all_tasks_done.notify_all()254            # All getters need to re-check queue-empty to raise ShutDown255            self.not_empty.notify_all()256            self.not_full.notify_all()257 258    # Override these methods to implement other queue organizations259    # (e.g. stack or priority queue).260    # These will only be called with appropriate locks held261 262    # Initialize the queue representation263    def _init(self, maxsize):264        self.queue = deque()265 266    def _qsize(self):267        return len(self.queue)268 269    # Put a new item in the queue270    def _put(self, item):271        self.queue.append(item)272 273    # Get an item from the queue274    def _get(self):275        return self.queue.popleft()276 277    __class_getitem__ = classmethod(types.GenericAlias)278 279 280class PriorityQueue(Queue):281    '''Variant of Queue that retrieves open entries in priority order (lowest first).282 283    Entries are typically tuples of the form:  (priority number, data).284    '''285 286    def _init(self, maxsize):287        self.queue = []288 289    def _qsize(self):290        return len(self.queue)291 292    def _put(self, item):293        heappush(self.queue, item)294 295    def _get(self):296        return heappop(self.queue)297 298 299class LifoQueue(Queue):300    '''Variant of Queue that retrieves most recently added entries first.'''301 302    def _init(self, maxsize):303        self.queue = []304 305    def _qsize(self):306        return len(self.queue)307 308    def _put(self, item):309        self.queue.append(item)310 311    def _get(self):312        return self.queue.pop()313 314 315class _PySimpleQueue:316    '''Simple, unbounded FIFO queue.317 318    This pure Python implementation is not reentrant.319    '''320    # Note: while this pure Python version provides fairness321    # (by using a threading.Semaphore which is itself fair, being based322    #  on threading.Condition), fairness is not part of the API contract.323    # This allows the C version to use a different implementation.324 325    def __init__(self):326        self._queue = deque()327        self._count = threading.Semaphore(0)328 329    def put(self, item, block=True, timeout=None):330        '''Put the item on the queue.331 332        The optional 'block' and 'timeout' arguments are ignored, as this method333        never blocks.  They are provided for compatibility with the Queue class.334        '''335        self._queue.append(item)336        self._count.release()337 338    def get(self, block=True, timeout=None):339        '''Remove and return an item from the queue.340 341        If optional args 'block' is true and 'timeout' is None (the default),342        block if necessary until an item is available. If 'timeout' is343        a non-negative number, it blocks at most 'timeout' seconds and raises344        the Empty exception if no item was available within that time.345        Otherwise ('block' is false), return an item if one is immediately346        available, else raise the Empty exception ('timeout' is ignored347        in that case).348        '''349        if timeout is not None and timeout < 0:350            raise ValueError("'timeout' must be a non-negative number")351        if not self._count.acquire(block, timeout):352            raise Empty353        return self._queue.popleft()354 355    def put_nowait(self, item):356        '''Put an item into the queue without blocking.357 358        This is exactly equivalent to `put(item, block=False)` and is only provided359        for compatibility with the Queue class.360        '''361        return self.put(item, block=False)362 363    def get_nowait(self):364        '''Remove and return an item from the queue without blocking.365 366        Only get an item if one is immediately available. Otherwise367        raise the Empty exception.368        '''369        return self.get(block=False)370 371    def empty(self):372        '''Return True if the queue is empty, False otherwise (not reliable!).'''373        return len(self._queue) == 0374 375    def qsize(self):376        '''Return the approximate size of the queue (not reliable!).'''377        return len(self._queue)378 379    __class_getitem__ = classmethod(types.GenericAlias)380 381 382if SimpleQueue is None:383    SimpleQueue = _PySimpleQueue384 
codekingpro/portable-devtools · Team Ai