package main import ( "fmt" "time" ) // TokenBucket represents a token bucket rate limiter type TokenBucket struct { capacity int // Maximum number of tokens the bucket can hold rate time.Duration // Time to add one token tokens chan struct{} // Channel acting as the token reservoir lastRefill time.Time } // NewTokenBucket creates and starts a new token bucket func NewTokenBucket(capacity int, refillRate time.Duration) *TokenBucket { tb := &TokenBucket{ capacity: capacity, rate: refillRate, tokens: make(chan struct{}, capacity), lastRefill: time.Now(), } // Fill the bucket initially for i := 0; i < capacity; i++ { tb.tokens <- struct{}{} } // Start a goroutine to continuously refill the bucket go tb.refiller() return tb } // refiller continuously adds tokens to the bucket func (tb *TokenBucket) refiller() { ticker := time.NewTicker(tb.rate) defer ticker.Stop() for range ticker.C { select { case tb.tokens <- struct{}{}: // Token added successfully default: // Bucket is full, do nothing } } } // Take attempts to take a token from the bucket. // It blocks until a token is available or returns false if non-blocking is used. func (tb *TokenBucket) Take() bool { select { case <-tb.tokens: return true // Token acquired default: return false // No token available } } func main() { // Create a bucket with capacity 5, refilling 1 token every 200ms limiter := NewTokenBucket(5, 200*time.Millisecond) fmt.Println("Starting rate limiter test...") for i := 0; i < 20; i++ { if limiter.Take() { fmt.Printf("Request %d: ALLOWED ", i+1) } else { fmt.Printf("Request %d: BLOCKED (Rate limit exceeded) ", i+1) } time.Sleep(50 * time.Millisecond) // Simulate incoming requests } fmt.Println("Rate limiter test finished.") time.Sleep(2 * time.Second) // Give refiller a chance to run }