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