Documentation
¶
Index ¶
Examples ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func NewRateLimiter ¶
func NewRateLimiter[T any](limiter *limiter.Limiter, keyGetter func(T) string) func(destination ro.Observable[T]) ro.Observable[T]
NewRateLimiter rate limits observable values using ulule/limiter with custom key extraction. Play: https://go.dev/play/p/V4meCiGc3bx
Example ¶
// Create a rate limiter with 5 requests per second
store := memory.NewStore()
rate := limiter.Rate{
Period: time.Second,
Limit: 5,
}
limiter := limiter.New(store, rate)
// Rate limit by user ID
observable := ro.Pipe1(
ro.Just("user1", "user2", "user1", "user3", "user1", "user2", "user1"),
NewRateLimiter[string](limiter, func(userID string) string {
return userID
}),
)
subscription := observable.Subscribe(ro.PrintObserver[string]())
defer subscription.Unsubscribe()
Output: Next: user1 Next: user2 Next: user1 Next: user3 Next: user1 Next: user2 Next: user1 Completed
Example (WithCompositeKey) ¶
// Create a rate limiter with 5 requests per minute
store := memory.NewStore()
rate := limiter.Rate{
Period: time.Minute,
Limit: 5,
}
limiter := limiter.New(store, rate)
type APIRequest struct {
IPAddress string
Endpoint string
Method string
}
// Rate limit by IP + endpoint combination
observable := ro.Pipe1(
ro.Just(
APIRequest{IPAddress: "192.168.1.1", Endpoint: "/api/users", Method: "GET"},
APIRequest{IPAddress: "192.168.1.2", Endpoint: "/api/users", Method: "GET"},
APIRequest{IPAddress: "192.168.1.1", Endpoint: "/api/posts", Method: "GET"},
APIRequest{IPAddress: "192.168.1.1", Endpoint: "/api/users", Method: "POST"},
APIRequest{IPAddress: "192.168.1.2", Endpoint: "/api/posts", Method: "GET"},
),
NewRateLimiter[APIRequest](limiter, func(req APIRequest) string {
return req.IPAddress + ":" + req.Endpoint
}),
)
subscription := observable.Subscribe(ro.PrintObserver[APIRequest]())
defer subscription.Unsubscribe()
Output: Next: {192.168.1.1 /api/users GET} Next: {192.168.1.2 /api/users GET} Next: {192.168.1.1 /api/posts GET} Next: {192.168.1.1 /api/users POST} Next: {192.168.1.2 /api/posts GET} Completed
Example (WithContext) ¶
// Create a rate limiter with 2 requests per second
store := memory.NewStore()
rate := limiter.Rate{
Period: time.Second,
Limit: 2,
}
limiter := limiter.New(store, rate)
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
// Rate limit with context
observable := ro.Pipe1(
ro.Just("user1", "user2", "user1", "user3", "user1", "user2", "user1"),
NewRateLimiter[string](limiter, func(userID string) string {
return userID
}),
)
subscription := observable.SubscribeWithContext(
ctx,
ro.NewObserverWithContext(
func(ctx context.Context, value string) {
// Handle rate-limited value with context
},
func(ctx context.Context, err error) {
// Handle error with context
},
func(ctx context.Context) {
// Handle completion with context
},
),
)
defer subscription.Unsubscribe()
Example (WithEndpoint) ¶
// Create a rate limiter with 2 requests per second
store := memory.NewStore()
rate := limiter.Rate{
Period: time.Second,
Limit: 2,
}
limiter := limiter.New(store, rate)
type APIRequest struct {
IPAddress string
Endpoint string
Method string
}
// Rate limit by endpoint
observable := ro.Pipe1(
ro.Just(
APIRequest{IPAddress: "192.168.1.1", Endpoint: "/api/users", Method: "GET"},
APIRequest{IPAddress: "192.168.1.2", Endpoint: "/api/posts", Method: "GET"},
APIRequest{IPAddress: "192.168.1.1", Endpoint: "/api/users", Method: "POST"},
APIRequest{IPAddress: "192.168.1.3", Endpoint: "/api/users", Method: "GET"},
APIRequest{IPAddress: "192.168.1.1", Endpoint: "/api/comments", Method: "GET"},
),
NewRateLimiter[APIRequest](limiter, func(req APIRequest) string {
return req.Endpoint
}),
)
subscription := observable.Subscribe(ro.PrintObserver[APIRequest]())
defer subscription.Unsubscribe()
Output: Next: {192.168.1.1 /api/users GET} Next: {192.168.1.2 /api/posts GET} Next: {192.168.1.1 /api/users POST} Next: {192.168.1.1 /api/comments GET} Completed
Example (WithErrorHandling) ¶
// Create a rate limiter with 3 requests per second
store := memory.NewStore()
rate := limiter.Rate{
Period: time.Second,
Limit: 3,
}
limiter := limiter.New(store, rate)
// Rate limit with error handling
observable := ro.Pipe1(
ro.Just("user1", "user2", "user1", "user3", "user1", "user2", "user1"),
NewRateLimiter[string](limiter, func(userID string) string {
return userID
}),
)
subscription := observable.Subscribe(
ro.NewObserver(
func(value string) {
// Handle successful rate-limited value
},
func(err error) {
// Handle rate limiting error
// This could be due to:
// - Store errors
// - Context cancellation
// - Other limiter errors
},
func() {
// Handle completion
},
),
)
defer subscription.Unsubscribe()
Example (WithIPAddress) ¶
// Create a rate limiter with 10 requests per minute
store := memory.NewStore()
rate := limiter.Rate{
Period: time.Minute,
Limit: 10,
}
limiter := limiter.New(store, rate)
type APIRequest struct {
IPAddress string
Endpoint string
Method string
}
// Rate limit by IP address
observable := ro.Pipe1(
ro.Just(
APIRequest{IPAddress: "192.168.1.1", Endpoint: "/api/users", Method: "GET"},
APIRequest{IPAddress: "192.168.1.2", Endpoint: "/api/users", Method: "GET"},
APIRequest{IPAddress: "192.168.1.1", Endpoint: "/api/posts", Method: "POST"},
APIRequest{IPAddress: "192.168.1.3", Endpoint: "/api/users", Method: "GET"},
APIRequest{IPAddress: "192.168.1.1", Endpoint: "/api/comments", Method: "GET"},
),
NewRateLimiter[APIRequest](limiter, func(req APIRequest) string {
return req.IPAddress
}),
)
subscription := observable.Subscribe(ro.PrintObserver[APIRequest]())
defer subscription.Unsubscribe()
Output: Next: {192.168.1.1 /api/users GET} Next: {192.168.1.2 /api/users GET} Next: {192.168.1.1 /api/posts POST} Next: {192.168.1.3 /api/users GET} Next: {192.168.1.1 /api/comments GET} Completed
Example (WithStructs) ¶
// Create a rate limiter with 3 requests per minute
store := memory.NewStore()
rate := limiter.Rate{
Period: time.Minute,
Limit: 3,
}
limiter := limiter.New(store, rate)
type Request struct {
UserID string
Action string
Data string
}
// Rate limit by user ID
observable := ro.Pipe1(
ro.Just(
Request{UserID: "user1", Action: "login", Data: "data1"},
Request{UserID: "user2", Action: "login", Data: "data2"},
Request{UserID: "user1", Action: "logout", Data: "data3"},
Request{UserID: "user3", Action: "login", Data: "data4"},
Request{UserID: "user1", Action: "update", Data: "data5"},
Request{UserID: "user2", Action: "logout", Data: "data6"},
),
NewRateLimiter[Request](limiter, func(req Request) string {
return req.UserID
}),
)
subscription := observable.Subscribe(ro.PrintObserver[Request]())
defer subscription.Unsubscribe()
Output: Next: {user1 login data1} Next: {user2 login data2} Next: {user1 logout data3} Next: {user3 login data4} Next: {user1 update data5} Next: {user2 logout data6} Completed
Types ¶
This section is empty.
Click to show internal directories.
Click to hide internal directories.