Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
112 changes: 112 additions & 0 deletions ratelimit.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
// ratelimit.go - Rate limit por peer (sender) para prevenir decrypt burst
//
// Contexto: quando WhatsApp faz sender key rotation em grupos ativos, msgs podem
// chegar em burst com iterations muito à frente. libsignal-go fará ratchet
// forward de até 2000 steps por msg (groups/GroupCipher.go:getSenderKey). Se
// 100+ msgs chegam próximas do mesmo peer, CPU/mem satura, decrypt trunca
// payloads longos.
//
// Fix: token bucket por peer. Se peer excede N msgs em janela T, dropar
// excedente. Msgs legítimas de conversação (< 5 msgs/segundo) passam livres.
// Bombardeio de rotação (>50 msgs/segundo) é achatado.
//
// Config via env:
// RATE_LIMIT_PEER_MSG_PER_SEC=20 (default; ajuste conforme uso)
// RATE_LIMIT_PEER_ENABLED=true (default; setar false pra desligar)

package main

import (
"os"
"strconv"
"sync"
"time"

"github.com/rs/zerolog/log"
)

type peerBucket struct {
tokens float64
lastRefill time.Time
mu sync.Mutex
}

type peerRateLimiter struct {
buckets sync.Map // key = peer JID string, value = *peerBucket
rate float64 // tokens per second
capacity float64 // max tokens
enabled bool
}

var globalPeerLimiter *peerRateLimiter
var globalPeerLimiterOnce sync.Once

func getPeerLimiter() *peerRateLimiter {
globalPeerLimiterOnce.Do(func() {
rate := parseEnvFloat("RATE_LIMIT_PEER_MSG_PER_SEC", 20.0)
enabled := parseEnvBool("RATE_LIMIT_PEER_ENABLED", true)
globalPeerLimiter = &peerRateLimiter{
rate: rate,
capacity: rate * 2, // burst tolerance = 2s
enabled: enabled,
}
log.Info().
Float64("rate_per_sec", rate).
Float64("capacity", rate*2).
Bool("enabled", enabled).
Msg("[PEER_RATE_LIMIT] initialized")
})
return globalPeerLimiter
}

// Allow returns true if peer msg should be processed, false if dropped.
func (p *peerRateLimiter) Allow(peerKey string) bool {
if !p.enabled {
return true
}
now := time.Now()
bucketI, _ := p.buckets.LoadOrStore(peerKey, &peerBucket{
tokens: p.capacity,
lastRefill: now,
})
b := bucketI.(*peerBucket)
b.mu.Lock()
defer b.mu.Unlock()

elapsed := now.Sub(b.lastRefill).Seconds()
b.tokens += elapsed * p.rate
if b.tokens > p.capacity {
b.tokens = p.capacity
}
b.lastRefill = now

if b.tokens < 1 {
return false
}
b.tokens--
return true
}

func parseEnvFloat(key string, def float64) float64 {
v := os.Getenv(key)
if v == "" {
return def
}
f, err := strconv.ParseFloat(v, 64)
if err != nil {
return def
}
return f
}

func parseEnvBool(key string, def bool) bool {
v := os.Getenv(key)
if v == "" {
return def
}
b, err := strconv.ParseBool(v)
if err != nil {
return def
}
return b
}
16 changes: 16 additions & 0 deletions wmiau.go
Original file line number Diff line number Diff line change
Expand Up @@ -1091,6 +1091,22 @@ func (mycli *MyClient) myEventHandler(rawEvt interface{}) {
return
case *events.Message:

// [PEER_RATE_LIMIT] descarta msgs em burst do mesmo peer (fix decrypt burst)
{
peerKey := evt.Info.Sender.String()
if evt.Info.IsGroup {
peerKey = evt.Info.Chat.String() + "|" + peerKey
}
if !getPeerLimiter().Allow(peerKey) {
log.Warn().
Str("peer", peerKey).
Str("msg_id", evt.Info.ID).
Bool("is_group", evt.Info.IsGroup).
Msg("[PEER_RATE_LIMIT] msg dropped (peer excedeu rate)")
return
}
}

var s3Config struct {
Enabled string `db:"s3_enabled"`
MediaDelivery string `db:"media_delivery"`
Expand Down