first commit
This commit is contained in:
@@ -0,0 +1,98 @@
|
||||
// Package nats содержит подписчика на NATS-метрики от нижестоящих АТС.
|
||||
package nats
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"sync"
|
||||
|
||||
natsgo "github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/pulse-lets-go/internal/models"
|
||||
"github.com/pulse-lets-go/internal/engine"
|
||||
filelog "github.com/pulse-lets-go/internal/log"
|
||||
)
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
// NewSubscriber подключается к NATS и запускает подписку.
|
||||
func NewSubscriber(natsURL, user, password string, eng *engine.Engine, l *filelog.Logger) (*Subscriber, error) {
|
||||
opts := []natsgo.Option{
|
||||
natsgo.ReconnectWait(natsgo.DefaultReconnectWait),
|
||||
natsgo.MaxReconnects(-1),
|
||||
}
|
||||
if user != "" && password != "" {
|
||||
opts = append(opts, natsgo.UserInfo(user, password))
|
||||
}
|
||||
|
||||
nc, err := natsgo.Connect(natsURL, opts...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("nats connect: %w", err)
|
||||
}
|
||||
|
||||
s := &Subscriber{
|
||||
nc: nc,
|
||||
engine: eng,
|
||||
logger: l,
|
||||
}
|
||||
|
||||
sub, err := nc.Subscribe(subject, s.handleMetric)
|
||||
if err != nil {
|
||||
nc.Close()
|
||||
return nil, fmt.Errorf("nats subscribe: %w", err)
|
||||
}
|
||||
s.sub = sub
|
||||
|
||||
log.Printf("[nats] подписан на %s, подключён к %s", subject, natsURL)
|
||||
return s, nil
|
||||
}
|
||||
|
||||
// handleMetric — обработчик входящих сообщений NATS.
|
||||
func (s *Subscriber) handleMetric(msg *natsgo.Msg) {
|
||||
var metric models.NodeMetric
|
||||
if err := json.Unmarshal(msg.Data, &metric); err != nil {
|
||||
log.Printf("[nats] ошибка парсинга метрики: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Валидация минимальных полей
|
||||
if metric.NodeID == "" {
|
||||
log.Printf("[nats] метрика без node_id, игнорируется")
|
||||
return
|
||||
}
|
||||
if metric.MaxCalls <= 0 {
|
||||
metric.MaxCalls = 250 // разумное значение по умолчанию
|
||||
}
|
||||
|
||||
ns := s.engine.UpdateMetric(&metric)
|
||||
s.logger.Log(ns)
|
||||
}
|
||||
|
||||
// IsConnected возвращает true, если NATS-соединение активно.
|
||||
func (s *Subscriber) IsConnected() bool {
|
||||
if s.nc == nil {
|
||||
return false
|
||||
}
|
||||
return s.nc.IsConnected()
|
||||
}
|
||||
|
||||
// Close останавливает подписку и закрывает соединение.
|
||||
func (s *Subscriber) Close() {
|
||||
if s.sub != nil {
|
||||
s.sub.Unsubscribe()
|
||||
}
|
||||
if s.nc != nil {
|
||||
s.nc.Close()
|
||||
}
|
||||
s.wg.Wait()
|
||||
}
|
||||
Reference in New Issue
Block a user