ratelimit.mx raw

   1  // Package ratelimit provides a cooperative token-bucket rate limiter.
   2  // No mutexes - designed for single-threaded cooperative scheduling.
   3  package ratelimit
   4  
   5  import "time"
   6  
   7  type bucket struct {
   8  	tokens float64
   9  	last   int64 // unix nano
  10  }
  11  
  12  // Limiter is a per-key token bucket rate limiter.
  13  //
  14  // buckets holds bucket state BY VALUE. A *bucket would be allocated by Allow,
  15  // and Allow mutates only through a map update - no store through the receiver -
  16  // so the compiler does not treat it as a mutating method of a sovereign
  17  // receiver and its allocations land in whatever frame arena is current. That
  18  // arena is released at the end of the turn, so the bucket did not persist and
  19  // every write was allowed: the limiter only ever worked by leaking (the event
  20  // loop's arena used to be permanent). Keeping the bucket in the map puts it in
  21  // the map's own storage, which lives as long as the limiter.
  22  type Limiter struct {
  23  	buckets map[string]bucket
  24  	rate    float64 // tokens per second
  25  	burst   int32     // max tokens
  26  }
  27  
  28  // New creates a rate limiter. Rate is tokens/second, burst is the
  29  // maximum tokens that can accumulate.
  30  func New(rate float64, burst int32) (l *Limiter) {
  31  	return &Limiter{
  32  		buckets: map[string]bucket{},
  33  		rate:    rate,
  34  		burst:   burst,
  35  	}
  36  }
  37  
  38  // Allow checks whether key has a token available. Consumes one token
  39  // if allowed.
  40  func (l *Limiter) Allow(key []byte) (ok bool) {
  41  	k := string(key)
  42  	now := time.Now().UnixNano()
  43  	b, found := l.buckets[k]
  44  	if !found {
  45  		b = bucket{tokens: float64(l.burst), last: now}
  46  	}
  47  	elapsed := float64(now-b.last) / 1e9
  48  	b.tokens += elapsed * l.rate
  49  	if b.tokens > float64(l.burst) {
  50  		b.tokens = float64(l.burst)
  51  	}
  52  	b.last = now
  53  	if b.tokens < 1.0 {
  54  		l.buckets[k] = b
  55  		return false
  56  	}
  57  	// Explicit write-back of the whole bucket: the value came out of the map by
  58  	// copy, so the decrement has to be stored back for the next call to see it.
  59  	b.tokens = b.tokens - 1.0
  60  	l.buckets[k] = b
  61  	return true
  62  }
  63  
  64  // Cleanup removes entries older than maxAge to prevent unbounded growth.
  65  func (l *Limiter) Cleanup(maxAge time.Duration) {
  66  	cutoff := time.Now().Add(-maxAge).UnixNano()
  67  	for k, b := range l.buckets {
  68  		if b.last < cutoff {
  69  			delete(l.buckets, k)
  70  		}
  71  	}
  72  }
  73