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