mirror of
https://github.com/mmahdium/TGSS.git
synced 2026-08-12 08:32:48 +03:30
feat: implement rate‑limiter cleanup task and register it in the FX lifecycle.
This commit is contained in:
@@ -17,6 +17,7 @@ func main() {
|
||||
infra.NewLogger,
|
||||
config.Load,
|
||||
server.NewGin,
|
||||
server.NewRateLimiter,
|
||||
|
||||
handlers.NewChannelHandler,
|
||||
handlers.NewAuthHandler,
|
||||
@@ -25,6 +26,7 @@ func main() {
|
||||
),
|
||||
telegram.Module,
|
||||
fx.Invoke(server.RegisterRoutes),
|
||||
fx.Invoke((server.RegisterCleanup)),
|
||||
fx.Logger(infra.NewFXLogger()),
|
||||
)
|
||||
|
||||
|
||||
@@ -16,11 +16,12 @@ import (
|
||||
type RouterParams struct {
|
||||
fx.In
|
||||
|
||||
Lifecycle fx.Lifecycle
|
||||
GinEngine *gin.Engine
|
||||
Logger *zap.Logger
|
||||
TgService *telegram.Service
|
||||
Config *config.Config
|
||||
Lifecycle fx.Lifecycle
|
||||
GinEngine *gin.Engine
|
||||
Logger *zap.Logger
|
||||
TgService *telegram.Service
|
||||
Config *config.Config
|
||||
RateLimiter *RateLimiter
|
||||
|
||||
ChannelHandler *handlers.ChannelHandler
|
||||
AuthHandler *handlers.AuthHandler
|
||||
@@ -47,7 +48,7 @@ func RegisterRoutes(params RouterParams) {
|
||||
})
|
||||
|
||||
// params.GinEngine.GET("/feed/:id/json", params.ChannelHandler.GetMessagesJson)
|
||||
params.GinEngine.GET("/feed/:id", RateLimit(), params.ChannelHandler.GetMessagesRSS)
|
||||
params.GinEngine.GET("/feed/:id", params.RateLimiter.RateLimit(), params.ChannelHandler.GetMessagesRSS)
|
||||
|
||||
// TODO: improve accurecy in rss channel fields
|
||||
// TODO: add ui with templates under /setup with fetch
|
||||
|
||||
@@ -1,10 +1,13 @@
|
||||
package server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"go.uber.org/fx"
|
||||
"go.uber.org/zap"
|
||||
"golang.org/x/time/rate"
|
||||
)
|
||||
|
||||
@@ -13,29 +16,70 @@ type ClientLimiter struct {
|
||||
lastSeen time.Time
|
||||
}
|
||||
|
||||
var clients = map[string]*ClientLimiter{}
|
||||
var mu sync.Mutex
|
||||
type RateLimiter struct {
|
||||
logger *zap.Logger
|
||||
clients map[string]*ClientLimiter
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
func getLimiter(ip string) *rate.Limiter {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
func NewRateLimiter(logger *zap.Logger) *RateLimiter {
|
||||
return &RateLimiter{logger: logger, mu: sync.Mutex{}, clients: map[string]*ClientLimiter{}}
|
||||
}
|
||||
|
||||
if c, ok := clients[ip]; ok {
|
||||
func (r *RateLimiter) getLimiter(ip string) *rate.Limiter {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
|
||||
if c, ok := r.clients[ip]; ok {
|
||||
c.lastSeen = time.Now()
|
||||
return c.limiter
|
||||
}
|
||||
|
||||
limiter := rate.NewLimiter(rate.Every(time.Second * 3), 5) // TODO: get from config
|
||||
clients[ip] = &ClientLimiter{limiter: limiter, lastSeen: time.Now()}
|
||||
limiter := rate.NewLimiter(rate.Every(time.Second*3), 5) // TODO: get from config
|
||||
r.clients[ip] = &ClientLimiter{limiter: limiter, lastSeen: time.Now()}
|
||||
|
||||
return limiter
|
||||
}
|
||||
|
||||
func RateLimit() gin.HandlerFunc {
|
||||
func (r *RateLimiter) CleanupRateLimiter() {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
|
||||
for ip, clientLimiter := range r.clients {
|
||||
if time.Since(clientLimiter.lastSeen) > time.Hour { // Delete ip ratelimiters older that 1 hour
|
||||
delete(r.clients, ip)
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func RegisterCleanup(lc fx.Lifecycle, logger *zap.Logger, r *RateLimiter) {
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(ctx context.Context) error {
|
||||
logger.Info("Starting rate limiter cleanup task", zap.Duration("interval", 30*time.Minute))
|
||||
ticker := time.NewTicker(30 * time.Minute)
|
||||
go func() {
|
||||
for range ticker.C {
|
||||
logger.Info("Running rate limiter cleanup")
|
||||
r.CleanupRateLimiter()
|
||||
logger.Info("Rate limiter cleanup completed")
|
||||
}
|
||||
}()
|
||||
return nil
|
||||
},
|
||||
OnStop: func(ctx context.Context) error {
|
||||
r.logger.Info("Stopping rate limiter cleanup task")
|
||||
return nil
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
func (r *RateLimiter) RateLimit() gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
ip := c.ClientIP()
|
||||
|
||||
limiter := getLimiter(ip)
|
||||
limiter := r.getLimiter(ip)
|
||||
|
||||
if !limiter.Allow() {
|
||||
c.AbortWithStatusJSON(429, gin.H{"error": "Slow down pls"})
|
||||
|
||||
Reference in New Issue
Block a user