codekingpro/portable-devtools
114k
1import collections2 3import eventlet4from eventlet import hubs5 6 7class Semaphore(object):8 9 """An unbounded semaphore.10 Optionally initialize with a resource *count*, then :meth:`acquire` and11 :meth:`release` resources as needed. Attempting to :meth:`acquire` when12 *count* is zero suspends the calling greenthread until *count* becomes13 nonzero again.14 15 This is API-compatible with :class:`threading.Semaphore`.16 17 It is a context manager, and thus can be used in a with block::18 19 sem = Semaphore(2)20 with sem:21 do_some_stuff()22 23 If not specified, *value* defaults to 1.24 25 It is possible to limit acquire time::26 27 sem = Semaphore()28 ok = sem.acquire(timeout=0.1)29 # True if acquired, False if timed out.30 31 """32 33 def __init__(self, value=1):34 try:35 value = int(value)36 except ValueError as e:37 msg = 'Semaphore() expect value :: int, actual: {0} {1}'.format(type(value), str(e))38 raise TypeError(msg)39 if value < 0:40 msg = 'Semaphore() expect value >= 0, actual: {0}'.format(repr(value))41 raise ValueError(msg)42 self.counter = value43 self._waiters = collections.deque()44 45 def __repr__(self):46 params = (self.__class__.__name__, hex(id(self)),47 self.counter, len(self._waiters))48 return '<%s at %s c=%s _w[%s]>' % params49 50 def __str__(self):51 params = (self.__class__.__name__, self.counter, len(self._waiters))52 return '<%s c=%s _w[%s]>' % params53 54 def locked(self):55 """Returns true if a call to acquire would block.56 """57 return self.counter <= 058 59 def bounded(self):60 """Returns False; for consistency with61 :class:`~eventlet.semaphore.CappedSemaphore`.62 """63 return False64 65 def acquire(self, blocking=True, timeout=None):66 """Acquire a semaphore.67 68 When invoked without arguments: if the internal counter is larger than69 zero on entry, decrement it by one and return immediately. If it is zero70 on entry, block, waiting until some other thread has called release() to71 make it larger than zero. This is done with proper interlocking so that72 if multiple acquire() calls are blocked, release() will wake exactly one73 of them up. The implementation may pick one at random, so the order in74 which blocked threads are awakened should not be relied on. There is no75 return value in this case.76 77 When invoked with blocking set to true, do the same thing as when called78 without arguments, and return true.79 80 When invoked with blocking set to false, do not block. If a call without81 an argument would block, return false immediately; otherwise, do the82 same thing as when called without arguments, and return true.83 84 Timeout value must be strictly positive.85 """86 if timeout == -1:87 timeout = None88 if timeout is not None and timeout < 0:89 raise ValueError("timeout value must be strictly positive")90 if not blocking:91 if timeout is not None:92 raise ValueError("can't specify timeout for non-blocking acquire")93 timeout = 094 if not blocking and self.locked():95 return False96 97 current_thread = eventlet.getcurrent()98 99 if self.counter <= 0 or self._waiters:100 if current_thread not in self._waiters:101 self._waiters.append(current_thread)102 try:103 if timeout is not None:104 ok = False105 with eventlet.Timeout(timeout, False):106 while self.counter <= 0:107 hubs.get_hub().switch()108 ok = True109 if not ok:110 return False111 else:112 # If someone else is already in this wait loop, give them113 # a chance to get out.114 while True:115 hubs.get_hub().switch()116 if self.counter > 0:117 break118 finally:119 try:120 self._waiters.remove(current_thread)121 except ValueError:122 # Fine if its already been dropped.123 pass124 125 self.counter -= 1126 return True127 128 def __enter__(self):129 self.acquire()130 131 def release(self, blocking=True):132 """Release a semaphore, incrementing the internal counter by one. When133 it was zero on entry and another thread is waiting for it to become134 larger than zero again, wake up that thread.135 136 The *blocking* argument is for consistency with CappedSemaphore and is137 ignored138 """139 self.counter += 1140 if self._waiters:141 hubs.get_hub().schedule_call_global(0, self._do_acquire)142 return True143 144 def _do_acquire(self):145 if self._waiters and self.counter > 0:146 waiter = self._waiters.popleft()147 waiter.switch()148 149 def __exit__(self, typ, val, tb):150 self.release()151 152 @property153 def balance(self):154 """An integer value that represents how many new calls to155 :meth:`acquire` or :meth:`release` would be needed to get the counter to156 0. If it is positive, then its value is the number of acquires that can157 happen before the next acquire would block. If it is negative, it is158 the negative of the number of releases that would be required in order159 to make the counter 0 again (one more release would push the counter to160 1 and unblock acquirers). It takes into account how many greenthreads161 are currently blocking in :meth:`acquire`.162 """163 # positive means there are free items164 # zero means there are no free items but nobody has requested one165 # negative means there are requests for items, but no items166 return self.counter - len(self._waiters)167 168 169class BoundedSemaphore(Semaphore):170 171 """A bounded semaphore checks to make sure its current value doesn't exceed172 its initial value. If it does, ValueError is raised. In most situations173 semaphores are used to guard resources with limited capacity. If the174 semaphore is released too many times it's a sign of a bug. If not given,175 *value* defaults to 1.176 """177 178 def __init__(self, value=1):179 super(BoundedSemaphore, self).__init__(value)180 self.original_counter = value181 182 def release(self, blocking=True):183 """Release a semaphore, incrementing the internal counter by one. If184 the counter would exceed the initial value, raises ValueError. When185 it was zero on entry and another thread is waiting for it to become186 larger than zero again, wake up that thread.187 188 The *blocking* argument is for consistency with :class:`CappedSemaphore`189 and is ignored190 """191 if self.counter >= self.original_counter:192 raise ValueError("Semaphore released too many times")193 return super(BoundedSemaphore, self).release(blocking)194 195 196class CappedSemaphore(object):197 198 """A blockingly bounded semaphore.199 200 Optionally initialize with a resource *count*, then :meth:`acquire` and201 :meth:`release` resources as needed. Attempting to :meth:`acquire` when202 *count* is zero suspends the calling greenthread until count becomes nonzero203 again. Attempting to :meth:`release` after *count* has reached *limit*204 suspends the calling greenthread until *count* becomes less than *limit*205 again.206 207 This has the same API as :class:`threading.Semaphore`, though its208 semantics and behavior differ subtly due to the upper limit on calls209 to :meth:`release`. It is **not** compatible with210 :class:`threading.BoundedSemaphore` because it blocks when reaching *limit*211 instead of raising a ValueError.212 213 It is a context manager, and thus can be used in a with block::214 215 sem = CappedSemaphore(2)216 with sem:217 do_some_stuff()218 """219 220 def __init__(self, count, limit):221 if count < 0:222 raise ValueError("CappedSemaphore must be initialized with a "223 "positive number, got %s" % count)224 if count > limit:225 # accidentally, this also catches the case when limit is None226 raise ValueError("'count' cannot be more than 'limit'")227 self.lower_bound = Semaphore(count)228 self.upper_bound = Semaphore(limit - count)229 230 def __repr__(self):231 params = (self.__class__.__name__, hex(id(self)),232 self.balance, self.lower_bound, self.upper_bound)233 return '<%s at %s b=%s l=%s u=%s>' % params234 235 def __str__(self):236 params = (self.__class__.__name__, self.balance,237 self.lower_bound, self.upper_bound)238 return '<%s b=%s l=%s u=%s>' % params239 240 def locked(self):241 """Returns true if a call to acquire would block.242 """243 return self.lower_bound.locked()244 245 def bounded(self):246 """Returns true if a call to release would block.247 """248 return self.upper_bound.locked()249 250 def acquire(self, blocking=True):251 """Acquire a semaphore.252 253 When invoked without arguments: if the internal counter is larger than254 zero on entry, decrement it by one and return immediately. If it is zero255 on entry, block, waiting until some other thread has called release() to256 make it larger than zero. This is done with proper interlocking so that257 if multiple acquire() calls are blocked, release() will wake exactly one258 of them up. The implementation may pick one at random, so the order in259 which blocked threads are awakened should not be relied on. There is no260 return value in this case.261 262 When invoked with blocking set to true, do the same thing as when called263 without arguments, and return true.264 265 When invoked with blocking set to false, do not block. If a call without266 an argument would block, return false immediately; otherwise, do the267 same thing as when called without arguments, and return true.268 """269 if not blocking and self.locked():270 return False271 self.upper_bound.release()272 try:273 return self.lower_bound.acquire()274 except:275 self.upper_bound.counter -= 1276 # using counter directly means that it can be less than zero.277 # however I certainly don't need to wait here and I don't seem to have278 # a need to care about such inconsistency279 raise280 281 def __enter__(self):282 self.acquire()283 284 def release(self, blocking=True):285 """Release a semaphore. In this class, this behaves very much like286 an :meth:`acquire` but in the opposite direction.287 288 Imagine the docs of :meth:`acquire` here, but with every direction289 reversed. When calling this method, it will block if the internal290 counter is greater than or equal to *limit*.291 """292 if not blocking and self.bounded():293 return False294 self.lower_bound.release()295 try:296 return self.upper_bound.acquire()297 except:298 self.lower_bound.counter -= 1299 raise300 301 def __exit__(self, typ, val, tb):302 self.release()303 304 @property305 def balance(self):306 """An integer value that represents how many new calls to307 :meth:`acquire` or :meth:`release` would be needed to get the counter to308 0. If it is positive, then its value is the number of acquires that can309 happen before the next acquire would block. If it is negative, it is310 the negative of the number of releases that would be required in order311 to make the counter 0 again (one more release would push the counter to312 1 and unblock acquirers). It takes into account how many greenthreads313 are currently blocking in :meth:`acquire` and :meth:`release`.314 """315 return self.lower_bound.balance - self.upper_bound.balance316 