Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
semaphore.py316 linesDownload Raw Back to eventlet
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 
codekingpro/portable-devtools · Team Ai