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)
|
||||
|
||||
+18
-10
@@ -29,6 +29,7 @@ type Engine struct {
|
||||
// Статистика для health endpoint
|
||||
startTime time.Time
|
||||
routeRequests atomic.Int64
|
||||
routeFallbacks atomic.Int64
|
||||
lastMetricTime atomic.Int64 // unix ts последней полученной метрики
|
||||
}
|
||||
|
||||
@@ -188,13 +189,19 @@ func (e *Engine) IncrementRouteRequests() {
|
||||
e.routeRequests.Add(1)
|
||||
}
|
||||
|
||||
// IncrementRouteFallbacks увеличивает счётчик fallback-запросов.
|
||||
func (e *Engine) IncrementRouteFallbacks() {
|
||||
e.routeFallbacks.Add(1)
|
||||
}
|
||||
|
||||
// HealthStats возвращает статистику для health endpoint.
|
||||
type HealthStats struct {
|
||||
TotalNodes int `json:"total_nodes"`
|
||||
HealthyNodes int `json:"healthy_nodes"`
|
||||
RouteRequests int64 `json:"route_requests_total"`
|
||||
UptimeSeconds int64 `json:"uptime_seconds"`
|
||||
LastMetricTS int64 `json:"last_metric_ts"`
|
||||
TotalNodes int `json:"total_nodes"`
|
||||
HealthyNodes int `json:"healthy_nodes"`
|
||||
RouteRequests int64 `json:"route_requests_total"`
|
||||
RouteFallbacks int64 `json:"route_fallbacks_total"`
|
||||
UptimeSeconds int64 `json:"uptime_seconds"`
|
||||
LastMetricTS int64 `json:"last_metric_ts"`
|
||||
}
|
||||
|
||||
// GetHealthStats возвращает текущую HealthStats (потокобезопасно).
|
||||
@@ -209,11 +216,12 @@ func (e *Engine) GetHealthStats() HealthStats {
|
||||
}
|
||||
}
|
||||
return HealthStats{
|
||||
TotalNodes: len(e.nodes),
|
||||
HealthyNodes: healthy,
|
||||
RouteRequests: e.routeRequests.Load(),
|
||||
UptimeSeconds: int64(time.Since(e.startTime).Seconds()),
|
||||
LastMetricTS: e.lastMetricTime.Load(),
|
||||
TotalNodes: len(e.nodes),
|
||||
HealthyNodes: healthy,
|
||||
RouteRequests: e.routeRequests.Load(),
|
||||
RouteFallbacks: e.routeFallbacks.Load(),
|
||||
UptimeSeconds: int64(time.Since(e.startTime).Seconds()),
|
||||
LastMetricTS: e.lastMetricTime.Load(),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -401,3 +401,70 @@ func (c *Client) GetStats() Stats {
|
||||
func (c *Client) IncGatewayOps() {
|
||||
c.gatewayOps.Add(1)
|
||||
}
|
||||
|
||||
// --- FS stats ---
|
||||
|
||||
// FsStats — показатели FreeSWITCH балансировщика.
|
||||
type FsStats struct {
|
||||
ActiveCalls int
|
||||
MaxSessions int
|
||||
FsUptime string
|
||||
}
|
||||
|
||||
// FetchFsStats запрашивает метрики FreeSWITCH через api команды.
|
||||
func (c *Client) FetchFsStats() (FsStats, error) {
|
||||
if !c.connected.Load() {
|
||||
return FsStats{}, fmt.Errorf("esl не подключён")
|
||||
}
|
||||
var stats FsStats
|
||||
|
||||
_, body, err := c.Send("api show calls count")
|
||||
if err != nil {
|
||||
return FsStats{}, err
|
||||
}
|
||||
stats.ActiveCalls = parseFirstInt(body)
|
||||
|
||||
_, body, err = c.Send("api status")
|
||||
if err != nil {
|
||||
return stats, err
|
||||
}
|
||||
stats.MaxSessions = parseMaxSessions(body)
|
||||
stats.FsUptime = parseUptime(body)
|
||||
|
||||
return stats, nil
|
||||
}
|
||||
|
||||
func parseFirstInt(s string) int {
|
||||
s = strings.TrimSpace(s)
|
||||
for i, c := range s {
|
||||
if c < '0' || c > '9' {
|
||||
if i == 0 {
|
||||
return 0
|
||||
}
|
||||
v, err := strconv.Atoi(s[:i])
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
return v
|
||||
}
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func parseMaxSessions(body string) int {
|
||||
for _, line := range strings.Split(body, "\n") {
|
||||
if strings.Contains(line, "session(s) max") {
|
||||
return parseFirstInt(line)
|
||||
}
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func parseUptime(body string) string {
|
||||
for _, line := range strings.Split(body, "\n") {
|
||||
if strings.HasPrefix(line, "UP ") {
|
||||
return strings.TrimPrefix(line, "UP ")
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
@@ -209,10 +209,29 @@ type ZabbixResponse struct {
|
||||
Data []ZabbixNode `json:"data"`
|
||||
}
|
||||
|
||||
// --- Balancer ---
|
||||
|
||||
// BalancerInfo — состояние балансировщика для дашборда.
|
||||
type BalancerInfo struct {
|
||||
EslConnected bool `json:"esl_connected"`
|
||||
EslUptimeSec int64 `json:"esl_uptime_sec,omitempty"`
|
||||
EslReconnects int64 `json:"esl_reconnects,omitempty"`
|
||||
NatsConnected bool `json:"nats_connected"`
|
||||
GatewaysTotal int `json:"gateways_total"`
|
||||
GatewaysUp int `json:"gateways_up"`
|
||||
GatewaysDown int `json:"gateways_down"`
|
||||
FsActiveCalls int `json:"fs_active_calls"`
|
||||
FsMaxSessions int `json:"fs_max_sessions"`
|
||||
FsUptime string `json:"fs_uptime"`
|
||||
RouteTotal int64 `json:"route_total"`
|
||||
RouteFallbacks int64 `json:"route_fallbacks"`
|
||||
UptimeSec int64 `json:"uptime_sec"`
|
||||
}
|
||||
|
||||
// --- WebSocket ---
|
||||
|
||||
// WsMetricsMessage — сообщение, отправляемое по WebSocket.
|
||||
type WsMetricsMessage struct {
|
||||
Type string `json:"type"` // "nodes_update", "node_toggle", etc
|
||||
Type string `json:"type"` // "nodes_update", "balancer_update", "gateway_update", etc
|
||||
Payload interface{} `json:"payload"`
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user