feat: balancer dashboard metrics + fix WebSocket Hijacker
This commit is contained in:
@@ -1,10 +1,13 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"net"
|
||||
"net/http"
|
||||
"time"
|
||||
)
|
||||
@@ -82,3 +85,11 @@ func (w *loggingResponseWriter) WriteHeader(code int) {
|
||||
func (w *loggingResponseWriter) Write(b []byte) (int, error) {
|
||||
return w.ResponseWriter.Write(b)
|
||||
}
|
||||
|
||||
func (w *loggingResponseWriter) Hijack() (net.Conn, *bufio.ReadWriter, error) {
|
||||
hijacker, ok := w.ResponseWriter.(http.Hijacker)
|
||||
if !ok {
|
||||
return nil, nil, fmt.Errorf("loggingResponseWriter: underlying ResponseWriter does not implement http.Hijacker")
|
||||
}
|
||||
return hijacker.Hijack()
|
||||
}
|
||||
|
||||
@@ -14,6 +14,7 @@ func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
// Нет зарегистрированных нод
|
||||
if nodeID == "" {
|
||||
a.engine.IncrementRouteFallbacks()
|
||||
writeJSON(w, http.StatusServiceUnavailable, models.RouteResponse{
|
||||
Error: "no_nodes_registered",
|
||||
})
|
||||
@@ -22,6 +23,7 @@ func (a *API) handleRoute(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
// Все ноды unhealthy — fallback
|
||||
if fallback || score < 0 {
|
||||
a.engine.IncrementRouteFallbacks()
|
||||
fallbackGW, _ := a.findFallbackGateway()
|
||||
nodes := a.getRouteNodeInfo()
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sync"
|
||||
|
||||
"github.com/pulse-lets-go/internal/config"
|
||||
"github.com/pulse-lets-go/internal/engine"
|
||||
@@ -23,6 +24,11 @@ type API struct {
|
||||
logFormat string
|
||||
routeLimiter *rateLimiter
|
||||
apiLimiter *rateLimiter
|
||||
|
||||
gatewayMu sync.RWMutex
|
||||
gatewayStates map[string]string // trunkID → "up"/"down"
|
||||
fsStats esl.FsStats
|
||||
fsStatsMu sync.RWMutex
|
||||
}
|
||||
|
||||
// NewAPI создаёт новый HTTP API с заданными зависимостями.
|
||||
@@ -35,6 +41,7 @@ func NewAPI(eng *engine.Engine, cfgMgr *config.Manager, jwtSecret, monitoringAPI
|
||||
wsHub: newWSHub(),
|
||||
natsConnected: natsFn,
|
||||
logFormat: cfg.Log.Format,
|
||||
gatewayStates: make(map[string]string),
|
||||
}
|
||||
if cfg.RateLimit.Enabled {
|
||||
a.routeLimiter = newRateLimiter(cfg.RateLimit.RoutePerSec, cfg.RateLimit.RoutePerSec)
|
||||
@@ -130,6 +137,10 @@ func corsMiddleware(next http.Handler) http.Handler {
|
||||
|
||||
// BroadcastGatewayEvent транслирует статус gateway всем WS-клиентам.
|
||||
func (a *API) BroadcastGatewayEvent(gatewayName, trunkID, status string) {
|
||||
a.gatewayMu.Lock()
|
||||
a.gatewayStates[trunkID] = status
|
||||
a.gatewayMu.Unlock()
|
||||
|
||||
msg := &models.WsMetricsMessage{
|
||||
Type: "gateway_update",
|
||||
Payload: map[string]string{
|
||||
@@ -141,6 +152,28 @@ func (a *API) BroadcastGatewayEvent(gatewayName, trunkID, status string) {
|
||||
a.wsHub.broadcast(msg)
|
||||
}
|
||||
|
||||
func (a *API) gatewayStats() (total, up, down int) {
|
||||
a.gatewayMu.RLock()
|
||||
defer a.gatewayMu.RUnlock()
|
||||
for _, s := range a.gatewayStates {
|
||||
total++
|
||||
switch s {
|
||||
case "up":
|
||||
up++
|
||||
case "down":
|
||||
down++
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// UpdateFsStats обновляет кэш метрик FreeSWITCH (вызывается из ticker).
|
||||
func (a *API) UpdateFsStats(stats esl.FsStats) {
|
||||
a.fsStatsMu.Lock()
|
||||
a.fsStats = stats
|
||||
a.fsStatsMu.Unlock()
|
||||
}
|
||||
|
||||
// maxBodySize — максимальный размер тела запроса в байтах (1 MiB).
|
||||
const maxBodySize = 1 << 20
|
||||
|
||||
|
||||
+42
-7
@@ -104,12 +104,7 @@ func (a *API) handleWS(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
// Отправляем текущее состояние при подключении
|
||||
go func() {
|
||||
nodes := a.engine.GetAllNodes()
|
||||
msg := &models.WsMetricsMessage{
|
||||
Type: "nodes_update",
|
||||
Payload: nodes,
|
||||
}
|
||||
a.wsHub.broadcast(msg)
|
||||
a.BroadcastAll()
|
||||
}()
|
||||
|
||||
// Читаем из вебсокета (ping/pong + закрытие)
|
||||
@@ -125,7 +120,7 @@ func (a *API) handleWS(w http.ResponseWriter, r *http.Request) {
|
||||
}()
|
||||
}
|
||||
|
||||
// BroadcastMetrics рассылает метрики всем WS-клиентам (вызывается engine при обновлении).
|
||||
// BroadcastMetrics рассылает метрики нод всем WS-клиентам.
|
||||
func (a *API) BroadcastMetrics() {
|
||||
nodes := a.engine.GetAllNodes()
|
||||
msg := &models.WsMetricsMessage{
|
||||
@@ -135,6 +130,46 @@ func (a *API) BroadcastMetrics() {
|
||||
a.wsHub.broadcast(msg)
|
||||
}
|
||||
|
||||
// BroadcastAll рассылает ноды и состояние балансировщика.
|
||||
func (a *API) BroadcastAll() {
|
||||
a.BroadcastMetrics()
|
||||
a.broadcastBalancerUpdate()
|
||||
}
|
||||
|
||||
func (a *API) broadcastBalancerUpdate() {
|
||||
stats := a.engine.GetHealthStats()
|
||||
gwTotal, gwUp, gwDown := a.gatewayStats()
|
||||
|
||||
info := models.BalancerInfo{
|
||||
NatsConnected: a.natsConnected != nil && a.natsConnected(),
|
||||
GatewaysTotal: gwTotal,
|
||||
GatewaysUp: gwUp,
|
||||
GatewaysDown: gwDown,
|
||||
RouteTotal: stats.RouteRequests,
|
||||
RouteFallbacks: stats.RouteFallbacks,
|
||||
UptimeSec: stats.UptimeSeconds,
|
||||
}
|
||||
|
||||
if a.eslClient != nil {
|
||||
eslStats := a.eslClient.GetStats()
|
||||
info.EslConnected = eslStats.Status == "connected"
|
||||
info.EslUptimeSec = eslStats.UptimeSec
|
||||
info.EslReconnects = eslStats.Reconnects
|
||||
|
||||
a.fsStatsMu.RLock()
|
||||
info.FsActiveCalls = a.fsStats.ActiveCalls
|
||||
info.FsMaxSessions = a.fsStats.MaxSessions
|
||||
info.FsUptime = a.fsStats.FsUptime
|
||||
a.fsStatsMu.RUnlock()
|
||||
}
|
||||
|
||||
msg := &models.WsMetricsMessage{
|
||||
Type: "balancer_update",
|
||||
Payload: info,
|
||||
}
|
||||
a.wsHub.broadcast(msg)
|
||||
}
|
||||
|
||||
// validateToken проверяет JWT и возвращает UserInfo.
|
||||
func (a *API) validateToken(tokenStr string) (*UserInfo, error) {
|
||||
token, err := parseJWT(tokenStr, a.jwtSecret)
|
||||
|
||||
Reference in New Issue
Block a user