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