feat: auto-trunk creation, SIP gateway in metrics, gateway WS broadcast
- SIPGateway field in NodeMetric model - HasNode() engine method - OnNewNode callback in NATS subscriber - Auto-create balance trunk from new node metric - BroadcastGatewayEvent on WS hub - Gateway status indicators in trunks UI (green/red/gray) - SPA fallback for frontend routing
This commit is contained in:
@@ -8,6 +8,7 @@ import (
|
||||
"github.com/pulse-lets-go/internal/config"
|
||||
"github.com/pulse-lets-go/internal/engine"
|
||||
"github.com/pulse-lets-go/internal/esl"
|
||||
"github.com/pulse-lets-go/internal/models"
|
||||
)
|
||||
|
||||
// API — главная структура HTTP API, агрегирует все зависимости.
|
||||
@@ -127,6 +128,19 @@ func corsMiddleware(next http.Handler) http.Handler {
|
||||
})
|
||||
}
|
||||
|
||||
// BroadcastGatewayEvent транслирует статус gateway всем WS-клиентам.
|
||||
func (a *API) BroadcastGatewayEvent(gatewayName, trunkID, status string) {
|
||||
msg := &models.WsMetricsMessage{
|
||||
Type: "gateway_update",
|
||||
Payload: map[string]string{
|
||||
"gateway": gatewayName,
|
||||
"trunk_id": trunkID,
|
||||
"status": status,
|
||||
},
|
||||
}
|
||||
a.wsHub.broadcast(msg)
|
||||
}
|
||||
|
||||
// decodeJSON декодирует тело запроса в структуру v.
|
||||
func decodeJSON(r *http.Request, v interface{}) error {
|
||||
defer r.Body.Close()
|
||||
|
||||
@@ -144,6 +144,14 @@ func (e *Engine) GetNodeHistory(nodeID string) []NodeSnapshot {
|
||||
return ring.snapshot()
|
||||
}
|
||||
|
||||
// HasNode возвращает true, если нода уже зарегистрирована в engine.
|
||||
func (e *Engine) HasNode(nodeID string) bool {
|
||||
e.mu.RLock()
|
||||
defer e.mu.RUnlock()
|
||||
_, exists := e.nodes[nodeID]
|
||||
return exists
|
||||
}
|
||||
|
||||
// ToggleNode включает/выключает ноду из распределения.
|
||||
func (e *Engine) ToggleNode(nodeID string, disabled bool, reason string) error {
|
||||
e.mu.Lock()
|
||||
|
||||
@@ -15,6 +15,7 @@ type NodeMetric struct {
|
||||
IdleCPU float64 `json:"idle_cpu"` // 0..100 %
|
||||
LoadAvg float64 `json:"load_avg"` // system load average (1m)
|
||||
CallFailureRate float64 `json:"call_failure_rate"` // 0.0..100.0 %
|
||||
SIPGateway string `json:"sip_gateway,omitempty"` // SIP-адрес ноды для авто-создания транка
|
||||
}
|
||||
|
||||
// --- In-memory состояние ноды в engine ---
|
||||
|
||||
@@ -18,11 +18,12 @@ const subject = "pulse.metrics.>"
|
||||
|
||||
// Subscriber подписывается на NATS и обновляет engine + пишет лог.
|
||||
type Subscriber struct {
|
||||
nc *natsgo.Conn
|
||||
sub *natsgo.Subscription
|
||||
engine *engine.Engine
|
||||
logger *filelog.Logger
|
||||
wg sync.WaitGroup
|
||||
nc *natsgo.Conn
|
||||
sub *natsgo.Subscription
|
||||
engine *engine.Engine
|
||||
logger *filelog.Logger
|
||||
onNewNode func(nodeID, sipGateway string) // вызывается при первой метрике от новой ноды
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
// NewSubscriber подключается к NATS и запускает подписку.
|
||||
@@ -57,6 +58,11 @@ func NewSubscriber(natsURL, user, password string, eng *engine.Engine, l *filelo
|
||||
return s, nil
|
||||
}
|
||||
|
||||
// SetOnNewNode устанавливает callback, вызываемый при первой метрике от новой ноды.
|
||||
func (s *Subscriber) SetOnNewNode(fn func(nodeID, sipGateway string)) {
|
||||
s.onNewNode = fn
|
||||
}
|
||||
|
||||
// handleMetric — обработчик входящих сообщений NATS.
|
||||
func (s *Subscriber) handleMetric(msg *natsgo.Msg) {
|
||||
var metric models.NodeMetric
|
||||
@@ -74,8 +80,16 @@ func (s *Subscriber) handleMetric(msg *natsgo.Msg) {
|
||||
metric.MaxCalls = 250 // разумное значение по умолчанию
|
||||
}
|
||||
|
||||
// Проверка: это новая нода с sip_gateway?
|
||||
isNew := !s.engine.HasNode(metric.NodeID)
|
||||
|
||||
ns := s.engine.UpdateMetric(&metric)
|
||||
s.logger.Log(ns)
|
||||
|
||||
// Авто-создание balance-транка для новой ноды
|
||||
if isNew && metric.SIPGateway != "" && s.onNewNode != nil {
|
||||
s.onNewNode(metric.NodeID, metric.SIPGateway)
|
||||
}
|
||||
}
|
||||
|
||||
// IsConnected возвращает true, если NATS-соединение активно.
|
||||
|
||||
Reference in New Issue
Block a user