feat: weighted random routing + mode switching + EWMA smoothing + stale eviction

- Weighted random routing (math/rand/v2 per-goroutine ChaCha8)
- PUT /api/route/mode for runtime mode switch (best / weighted_random)
- EWMA score smoothing (configurable smoothing_factor, default 0.3)
- Periodic stale node eviction (5 min threshold, 60s interval)
- Non-blocking WS broadcast (per-client buffered channel + writePump)
- Background rate limiter cleanup goroutine
- In-memory trunk gateway cache (zero disk I/O in hot path)
- Configurable load_avg multiplier (default 50.0, backward compat)
- UI mode indicator + admin Switch button in dashboard
- 42 tests: weighted random, mode switching, EWMA, eviction, API integration
This commit is contained in:
Maksim Totmin
2026-06-25 20:42:46 +07:00
parent 3a082923cf
commit b0697da8e4
15 changed files with 1026 additions and 65 deletions
+23 -15
View File
@@ -14,7 +14,6 @@ type rateLimiter struct {
rate float64 // токенов в секунду
burst int // максимальный размер бакета
cleanupInterval time.Duration
lastCleanup time.Time
}
type tokenBucket struct {
@@ -23,17 +22,22 @@ type tokenBucket struct {
}
// newRateLimiter создаёт rate limiter с заданными параметрами.
// Запускает фоновую горутину очистки устаревших бакетов.
func newRateLimiter(ratePerSec, burst int) *rateLimiter {
if burst <= 0 {
burst = ratePerSec
}
return &rateLimiter{
rl := &rateLimiter{
buckets: make(map[string]*tokenBucket),
rate: float64(ratePerSec),
burst: burst,
cleanupInterval: 60 * time.Second,
lastCleanup: time.Now(),
}
// Фоновая горутина очистки — не блокирует hot path
go rl.cleanupLoop()
return rl
}
// allow проверяет, разрешён ли запрос для данного ключа.
@@ -44,22 +48,11 @@ func (rl *rateLimiter) allow(key string) bool {
now := time.Now()
// Периодическая очистка устаревших бакетов
if now.Sub(rl.lastCleanup) > rl.cleanupInterval {
for k, b := range rl.buckets {
if now.Sub(b.lastSeen) > rl.cleanupInterval {
delete(rl.buckets, k)
}
}
rl.lastCleanup = now
}
b, exists := rl.buckets[key]
if !exists {
// Новый бакет с полным запасом токенов
b = &tokenBucket{tokens: float64(rl.burst), lastSeen: now}
rl.buckets[key] = b
b.tokens-- // расходуем один токен
b.tokens--
return true
}
@@ -78,6 +71,21 @@ func (rl *rateLimiter) allow(key string) bool {
return false
}
// cleanupLoop — фоновая очистка устаревших бакетов.
func (rl *rateLimiter) cleanupLoop() {
ticker := time.NewTicker(rl.cleanupInterval)
defer ticker.Stop()
for range ticker.C {
rl.mu.Lock()
for k, b := range rl.buckets {
if time.Since(b.lastSeen) > rl.cleanupInterval {
delete(rl.buckets, k)
}
}
rl.mu.Unlock()
}
}
// rateLimitMiddleware создаёт middleware с rate limiting по IP.
// Ключ — RemoteAddr (или X-Forwarded-For, если за проксей).
func rateLimitMiddleware(limiter *rateLimiter) func(http.Handler) http.Handler {
+43 -25
View File
@@ -10,7 +10,7 @@ import (
func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
a.engine.IncrementRouteRequests()
nodeID, score, fallback := a.engine.GetBestNode()
nodeID, score, fallback := a.engine.PickNode()
// Нет зарегистрированных нод
if nodeID == "" {
@@ -24,7 +24,7 @@ func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
// Все ноды unhealthy — fallback
if fallback || score < 0 {
a.engine.IncrementRouteFallbacks()
fallbackGW, _ := a.findFallbackGateway()
fallbackGW := a.cachedFallbackGateway()
nodes := a.getRouteNodeInfo()
writeJSON(w, http.StatusOK, models.RouteResponse{
@@ -36,8 +36,8 @@ func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
return
}
// Нормальный маршрут
gw, ok := a.findBalanceGateway(nodeID)
// Нормальный маршрут — O(1) из in-memory кэша
gw, ok := a.cachedBalanceGateway(nodeID)
if !ok {
writeError(w, http.StatusInternalServerError, "gateway не найден для ноды")
return
@@ -50,32 +50,50 @@ func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
})
}
// findBalanceGateway ищет gateway для balance-транка, привязанного к nodeID.
func (a *API) findBalanceGateway(nodeID string) (string, bool) {
trunks, err := a.configManager.ReadTrunks()
if err != nil {
return "", false
// handleSetRouteMode — PUT /api/route/mode (admin only).
func (a *API) handleSetRouteMode(w http.ResponseWriter, r *http.Request) {
var req struct {
Mode string `json:"mode"`
}
for _, t := range trunks {
if t.Type == models.TrunkBalance && t.NodeID == nodeID && t.Enabled {
return t.Gateway, true
if err := decodeJSON(r, &req); err != nil {
writeError(w, http.StatusBadRequest, "некорректный JSON")
return
}
if req.Mode != "best" && req.Mode != "weighted_random" {
writeError(w, http.StatusBadRequest, "mode должен быть 'best' или 'weighted_random'")
return
}
// Сохраняем в config.json атомарно
if a.configManager != nil {
cfg, err := a.configManager.ReadConfig()
if err == nil {
cfg.BalanceMode = req.Mode
a.configManager.SaveConfig(cfg)
}
}
return "", false
a.engine.SetBalancingMode(req.Mode)
// Немедленный WS-push для обновления UI
go a.broadcastBalancerUpdate()
writeJSON(w, http.StatusOK, map[string]string{"mode": req.Mode})
}
// findFallbackGateway ищет gateway для fallback-транка.
func (a *API) findFallbackGateway() (string, bool) {
trunks, err := a.configManager.ReadTrunks()
if err != nil {
return "", false
}
for _, t := range trunks {
if t.Type == models.TrunkFallback && t.Enabled {
return t.Gateway, true
}
}
return "", false
// cachedBalanceGateway возвращает gateway для balance-транка из in-memory кэша (O(1)).
func (a *API) cachedBalanceGateway(nodeID string) (string, bool) {
a.trunkCacheMu.RLock()
defer a.trunkCacheMu.RUnlock()
gw, ok := a.balanceGateways[nodeID]
return gw, ok
}
// cachedFallbackGateway возвращает fallback gateway из in-memory кэша (O(1)).
func (a *API) cachedFallbackGateway() string {
a.trunkCacheMu.RLock()
defer a.trunkCacheMu.RUnlock()
return a.fallbackGateway
}
// getRouteNodeInfo возвращает список нод с их скорами для ответа маршрутизации.
+154
View File
@@ -0,0 +1,154 @@
package api
import (
"bytes"
"encoding/json"
"net/http/httptest"
"testing"
"time"
"github.com/golang-jwt/jwt/v5"
"github.com/pulse-lets-go/internal/config"
"github.com/pulse-lets-go/internal/engine"
)
func testCfg() *config.Config {
return &config.Config{
StaleThresholdSec: 20,
Log: config.LogConfig{Level: "info", Format: "text"},
RateLimit: config.RateLimitConfig{Enabled: false},
BalanceMode: "weighted_random",
Scoring: config.ScoringConfig{
IdleCPUMin: 5.0,
CallFailureRateLethal: 15.0,
Weights: config.ScoringWeights{CallScore: 0.40, LoadScore: 0.30, IdleScore: 0.20, FailScore: 0.10},
},
}
}
func createTestJWT(secret, username, role string) string {
claims := jwt.MapClaims{
"user_id": "user-test",
"username": username,
"role": role,
"exp": time.Now().Add(1 * time.Hour).Unix(),
}
token := jwt.NewWithClaims(jwt.SigningMethodHS256, claims)
s, _ := token.SignedString([]byte(secret))
return s
}
func TestRouteModeEndpoint_NoAuth(t *testing.T) {
cfg := testCfg()
eng := engine.NewEngine(cfg)
a := NewAPI(eng, nil, "test-secret", "test-api-key", func() bool { return true }, cfg)
handler := a.Handler()
body, _ := json.Marshal(map[string]string{"mode": "best"})
req := httptest.NewRequest("PUT", "/api/route/mode", bytes.NewReader(body))
req.Header.Set("Content-Type", "application/json")
rr := httptest.NewRecorder()
handler.ServeHTTP(rr, req)
if rr.Code == 404 {
t.Fatalf("PUT /api/route/mode вернул 404 — маршрут не найден! body=%s", rr.Body.String())
}
if rr.Code == 401 {
t.Log("OK: без токена получаем 401 (роутинг работает)")
} else {
t.Logf("Код ответа: %d, body: %s", rr.Code, rr.Body.String())
}
}
func TestRouteModeEndpoint_ViewerForbidden(t *testing.T) {
cfg := testCfg()
eng := engine.NewEngine(cfg)
a := NewAPI(eng, nil, "test-secret", "test-api-key", func() bool { return true }, cfg)
handler := a.Handler()
token := createTestJWT("test-secret", "pbx-viewer", "viewer")
body, _ := json.Marshal(map[string]string{"mode": "weighted_random"})
req := httptest.NewRequest("PUT", "/api/route/mode", bytes.NewReader(body))
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+token)
rr := httptest.NewRecorder()
handler.ServeHTTP(rr, req)
if rr.Code == 404 {
t.Fatalf("PUT /api/route/mode вернул 404 (viewer) — маршрут не найден!")
}
if rr.Code != 403 {
t.Errorf("viewer должен получить 403, получен %d", rr.Code)
}
t.Logf("viewer: код=%d, body=%s", rr.Code, rr.Body.String())
}
func TestRouteModeEndpoint_AdminSuccess(t *testing.T) {
cfg := testCfg()
eng := engine.NewEngine(cfg)
a := NewAPI(eng, nil, "test-secret", "test-api-key", func() bool { return true }, cfg)
handler := a.Handler()
token := createTestJWT("test-secret", "admin", "admin")
body, _ := json.Marshal(map[string]string{"mode": "best"})
req := httptest.NewRequest("PUT", "/api/route/mode", bytes.NewReader(body))
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+token)
rr := httptest.NewRecorder()
handler.ServeHTTP(rr, req)
if rr.Code == 404 {
t.Fatalf("PUT /api/route/mode вернул 404 (admin) — маршрут не найден!")
}
if rr.Code != 200 {
t.Fatalf("admin должен получить 200, получен %d body=%s", rr.Code, rr.Body.String())
}
var resp map[string]string
json.Unmarshal(rr.Body.Bytes(), &resp)
if resp["mode"] != "best" {
t.Errorf("ожидался mode=best, получен %s", resp["mode"])
}
if eng.GetBalancingMode() != "best" {
t.Errorf("engine mode должен быть 'best', получен %s", eng.GetBalancingMode())
}
t.Logf("admin: код=%d, body=%s, engine_mode=%s", rr.Code, rr.Body.String(), eng.GetBalancingMode())
}
func TestRouteModeEndpoint_InvalidMode(t *testing.T) {
cfg := testCfg()
eng := engine.NewEngine(cfg)
a := NewAPI(eng, nil, "test-secret", "test-api-key", func() bool { return true }, cfg)
handler := a.Handler()
token := createTestJWT("test-secret", "admin", "admin")
body, _ := json.Marshal(map[string]string{"mode": "invalid"})
req := httptest.NewRequest("PUT", "/api/route/mode", bytes.NewReader(body))
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+token)
rr := httptest.NewRecorder()
handler.ServeHTTP(rr, req)
if rr.Code != 400 {
t.Errorf("невалидный mode должен дать 400, получен %d", rr.Code)
}
}
func TestRouteEndpoint_UsesPickNode(t *testing.T) {
cfg := testCfg()
eng := engine.NewEngine(cfg)
a := NewAPI(eng, nil, "test-secret", "test-api-key", func() bool { return true }, cfg)
a.balanceGateways["test-node"] = "sip:test:5060"
a.fallbackGateway = "sip:fallback:5060"
handler := a.Handler()
req := httptest.NewRequest("GET", "/api/route", nil)
rr := httptest.NewRecorder()
handler.ServeHTTP(rr, req)
if rr.Code != 503 {
t.Errorf("GET /api/route без нод должен дать 503, получен %d", rr.Code)
}
}
+37
View File
@@ -25,6 +25,11 @@ type API struct {
routeLimiter *rateLimiter
apiLimiter *rateLimiter
// In-memory кэш gateway для /api/route (устраняет disk I/O в hot path)
trunkCacheMu sync.RWMutex
balanceGateways map[string]string // nodeID → gateway
fallbackGateway string
gatewayMu sync.RWMutex
gatewayStates map[string]string // trunkID → "up"/"down"
fsStats esl.FsStats
@@ -42,6 +47,7 @@ func NewAPI(eng *engine.Engine, cfgMgr *config.Manager, jwtSecret, monitoringAPI
natsConnected: natsFn,
logFormat: cfg.Log.Format,
gatewayStates: make(map[string]string),
balanceGateways: make(map[string]string),
}
if cfg.RateLimit.Enabled {
a.routeLimiter = newRateLimiter(cfg.RateLimit.RoutePerSec, cfg.RateLimit.RoutePerSec)
@@ -81,6 +87,7 @@ func (a *API) Handler() http.Handler {
mux.Handle("GET /api/nodes", a.authMiddleware(http.HandlerFunc(a.handleGetNodes)))
mux.Handle("GET /api/nodes/{id}/metrics", a.authMiddleware(http.HandlerFunc(a.handleGetNodeMetrics)))
mux.Handle("PUT /api/nodes/{id}/toggle", a.authMiddleware(a.adminOnly(http.HandlerFunc(a.handleToggleNode))))
mux.Handle("PUT /api/route/mode", a.authMiddleware(a.adminOnly(http.HandlerFunc(a.handleSetRouteMode))))
// --- Админка: Trunks (admin only) ---
@@ -167,6 +174,36 @@ func (a *API) gatewayStats() (total, up, down int) {
return
}
// RebuildTrunkCache перестраивает in-memory кэш gateway для /api/route.
// Вызывается при старте и после любого CRUD-оперирования с транками.
func (a *API) RebuildTrunkCache() {
trunks, err := a.configManager.ReadTrunks()
if err != nil {
return
}
a.trunkCacheMu.Lock()
defer a.trunkCacheMu.Unlock()
a.balanceGateways = make(map[string]string, len(trunks))
a.fallbackGateway = ""
for _, t := range trunks {
if t.Enabled {
switch t.Type {
case models.TrunkBalance:
if t.NodeID != "" {
a.balanceGateways[t.NodeID] = t.Gateway
}
case models.TrunkFallback:
if a.fallbackGateway == "" {
a.fallbackGateway = t.Gateway
}
}
}
}
}
// UpdateFsStats обновляет кэш метрик FreeSWITCH (вызывается из ticker).
func (a *API) UpdateFsStats(stats esl.FsStats) {
a.fsStatsMu.Lock()
+6
View File
@@ -92,6 +92,8 @@ func (a *API) handleCreateTrunk(w http.ResponseWriter, r *http.Request) {
// ESL push: создаём gateway на FS для ingress-транка
a.eslPushGatewayAdd(trunk)
a.RebuildTrunkCache()
writeJSON(w, http.StatusCreated, trunk)
}
@@ -166,6 +168,8 @@ func (a *API) handleUpdateTrunk(w http.ResponseWriter, r *http.Request) {
// ESL push: обновляем gateway на FS
a.eslPushGatewayUpdate(*t)
a.RebuildTrunkCache()
writeJSON(w, http.StatusOK, t)
}
@@ -200,6 +204,8 @@ func (a *API) handleDeleteTrunk(w http.ResponseWriter, r *http.Request) {
return
}
a.RebuildTrunkCache()
writeJSON(w, http.StatusNoContent, nil)
}
+33 -11
View File
@@ -20,10 +20,12 @@ var upgrader = websocket.Upgrader{
},
}
// wsClient — соединение WebSocket одного клиента.
const wsSendBufferSize = 32 // размер буфера per-client канала
// wsClient — соединение WebSocket одного клиента с выделенной write-горутиной.
type wsClient struct {
conn *websocket.Conn
mu sync.Mutex
send chan *models.WsMetricsMessage
}
// wsHub управляет подключёнными WebSocket-клиентами и рассылает метрики.
@@ -46,32 +48,45 @@ func (h *wsHub) add(c *wsClient) {
h.mu.Unlock()
}
// remove удаляет клиента из хаба.
// remove удаляет клиента из хаба и закрывает его канал.
// Идемпотентно — может вызываться многократно (writePump + readPump).
func (h *wsHub) remove(c *wsClient) {
h.mu.Lock()
if h.clients[c] {
c.conn.Close()
delete(h.clients, c)
close(c.send)
c.conn.Close()
}
h.mu.Unlock()
}
// broadcast рассылает метрики всем подключённым клиентам.
// broadcast рассылает метрики всем клиентам неблокирующе.
// Медленные клиенты пропускают сообщения (best-effort delivery).
func (h *wsHub) broadcast(msg *models.WsMetricsMessage) {
h.mu.RLock()
defer h.mu.RUnlock()
for client := range h.clients {
client.mu.Lock()
err := client.conn.WriteJSON(msg)
client.mu.Unlock()
if err != nil {
log.Printf("[ws] ошибка отправки: %v", err)
select {
case client.send <- msg:
default:
// Клиент слишком медленный — пропускаем сообщение
go h.remove(client)
}
}
}
// writePump — выделенная горутина записи для одного клиента.
// Читает из client.send и пишет в WebSocket.
func (c *wsClient) writePump(hub *wsHub) {
defer hub.remove(c)
for msg := range c.send {
if err := c.conn.WriteJSON(msg); err != nil {
return
}
}
}
// handleWS — WebSocket /ws/metrics.
func (a *API) handleWS(w http.ResponseWriter, r *http.Request) {
token := r.URL.Query().Get("token")
@@ -97,11 +112,17 @@ func (a *API) handleWS(w http.ResponseWriter, r *http.Request) {
return
}
client := &wsClient{conn: conn}
client := &wsClient{
conn: conn,
send: make(chan *models.WsMetricsMessage, wsSendBufferSize),
}
a.wsHub.add(client)
log.Printf("[ws] клиент подключён: %s (%s)", userInfo.Username, r.RemoteAddr)
// Запускаем write-горутину
go client.writePump(a.wsHub)
// Отправляем текущее состояние при подключении
go func() {
a.BroadcastAll()
@@ -148,6 +169,7 @@ func (a *API) broadcastBalancerUpdate() {
RouteTotal: stats.RouteRequests,
RouteFallbacks: stats.RouteFallbacks,
UptimeSec: stats.UptimeSeconds,
BalanceMode: a.engine.GetBalancingMode(),
}
if a.eslClient != nil {