Netpulse_SasS/server/internal/httpapi/ws.go
zotac 328850be06 Етап 3: REST/WebSocket API для UI
Окремий процес від AgentService: у зондів і браузерів різні профілі
навантаження й периметри, спільний лише шар store.

Стан мапи віддається одним викликом разом із живими статусами — інакше
полотно малювалося б сірим і доганяло кольори сотнею дозапитів.

link_status виводиться з кінців лінка, а не читається з topo.links: цю
колонку ніхто не підтримує, а стан лінка це похідна величина. Живий
прогін показував "unknown" на цілком робочому каналі, поки не порахували.

Події для WebSocket пишуться тією ж транзакцією, що й зміна, яку
описують. Транспорт — опитування outbox, а не LISTEN/NOTIFY: NOTIFY не
переживає падіння підписника, а тут потрібна гарантія доставки.

Перевірено: 10 інтеграційних тестів проти живої БД, HTTP і WebSocket
(-race), плюс живий прогін агента, сервера й API разом проти snmpd.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-14 15:43:02 +03:00

358 lines
9.6 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package httpapi
import (
"context"
"encoding/json"
"errors"
"log/slog"
"net/http"
"sync"
"sync/atomic"
"time"
"github.com/coder/websocket"
"github.com/netpulse/netpulse/server/internal/store"
)
// Дві частоти, бо це дві різні за природою речі.
//
// Зміна статусу — подія: вона рідка, але має дійти майже миттєво,
// інакше мапа бреше про стан мережі. Завантаження каналу — величина:
// вона змінюється весь час, і слати її частіше, ніж оновлюються
// лічильники (60 с), означає слати ту саму цифру по колу.
const (
eventPollInterval = time.Second
linkPushInterval = 5 * time.Second
writeTimeout = 10 * time.Second
// Черга одного підписника. Переповнення означає, що клієнт не
// читає, — таке з'єднання закривається, а не гальмує решту.
subscriberQueue = 64
)
// Hub тримає підписників і транслює їм зміни.
type Hub struct {
store *store.Store
log *slog.Logger
mu sync.RWMutex
subs map[*subscriber]struct{}
lastEventID atomic.Int64
}
type subscriber struct {
tenantID string
send chan []byte
mu sync.RWMutex
mapID string
}
func (s *subscriber) currentMap() string {
s.mu.RLock()
defer s.mu.RUnlock()
return s.mapID
}
func (s *subscriber) setMap(id string) {
s.mu.Lock()
s.mapID = id
s.mu.Unlock()
}
func NewHub(st *store.Store, log *slog.Logger) *Hub {
if log == nil {
log = slog.Default()
}
return &Hub{store: st, log: log, subs: make(map[*subscriber]struct{})}
}
func (h *Hub) SubscriberCount() int {
h.mu.RLock()
defer h.mu.RUnlock()
return len(h.subs)
}
func (h *Hub) add(s *subscriber) {
h.mu.Lock()
h.subs[s] = struct{}{}
h.mu.Unlock()
}
func (h *Hub) remove(s *subscriber) {
h.mu.Lock()
if _, ok := h.subs[s]; ok {
delete(h.subs, s)
close(s.send)
}
h.mu.Unlock()
}
// Run крутить обидва насоси, доки не скасують контекст.
func (h *Hub) Run(ctx context.Context) {
// Стартуємо з кінця журналу: клієнт щойно завантажив повний стан
// мапи, і все старіше в ньому вже враховано.
if id, err := h.store.LatestEventID(ctx); err == nil {
h.lastEventID.Store(id)
} else {
h.log.Warn("не вдалося прочитати позицію журналу подій", "err", err)
}
var wg sync.WaitGroup
wg.Add(2)
go func() { defer wg.Done(); h.pumpEvents(ctx) }()
go func() { defer wg.Done(); h.pumpLinkLoads(ctx) }()
wg.Wait()
}
// pumpEvents читає core.event_outbox і розсилає події тенанта.
func (h *Hub) pumpEvents(ctx context.Context) {
t := time.NewTicker(eventPollInterval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
}
if h.SubscriberCount() == 0 {
// Підписників немає — але позицію все одно рухаємо вперед,
// інакше після довгої паузи перший клієнт отримає лавину
// накопичених подій.
if id, err := h.store.LatestEventID(ctx); err == nil {
h.lastEventID.Store(id)
}
continue
}
events, err := h.store.FetchEvents(ctx, h.lastEventID.Load(), 500)
if err != nil {
h.log.Warn("читання журналу подій", "err", err)
continue
}
if len(events) == 0 {
continue
}
for _, e := range events {
msg, err := json.Marshal(map[string]any{
"type": e.Topic,
"payload": e.Payload,
"at": e.At,
})
if err != nil {
continue
}
h.broadcast(e.TenantID, "", msg)
}
last := events[len(events)-1].ID
h.lastEventID.Store(last)
if err := h.store.MarkEventsPublished(ctx, last); err != nil {
h.log.Warn("позначення подій доставленими", "err", err)
}
}
}
// pumpLinkLoads розсилає завантаження каналів по підписаних мапах.
func (h *Hub) pumpLinkLoads(ctx context.Context) {
t := time.NewTicker(linkPushInterval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
}
// Один запит на мапу, а не на підписника: у NOC на одну мапу
// зазвичай дивиться кілька екранів одночасно.
type key struct{ tenant, mapID string }
wanted := make(map[key]struct{})
h.mu.RLock()
for s := range h.subs {
if m := s.currentMap(); m != "" {
wanted[key{s.tenantID, m}] = struct{}{}
}
}
h.mu.RUnlock()
for k := range wanted {
loads, err := h.store.MapLinkLoads(ctx, k.tenant, k.mapID)
if err != nil {
h.log.Warn("завантаження лінків", "map", k.mapID, "err", err)
continue
}
if len(loads) == 0 {
continue
}
msg, err := json.Marshal(map[string]any{
"type": "link.load",
"map_id": k.mapID,
"links": loads,
"at": time.Now(),
})
if err != nil {
continue
}
h.broadcast(k.tenant, k.mapID, msg)
}
}
}
// broadcast розсилає повідомлення підписникам тенанта.
// Порожній mapID означає «усім у тенанті».
func (h *Hub) broadcast(tenantID, mapID string, msg []byte) {
var stalled []*subscriber
h.mu.RLock()
for s := range h.subs {
if s.tenantID != tenantID {
continue
}
if mapID != "" && s.currentMap() != mapID {
continue
}
select {
case s.send <- msg:
default:
// Клієнт не читає. Тримати його чергу означало б
// віддавати пам'ять і затримувати всіх інших.
stalled = append(stalled, s)
}
}
h.mu.RUnlock()
for _, s := range stalled {
h.log.Warn("підписник не встигає читати — відключаємо")
h.remove(s)
}
}
// ---------------------------------------------------------------------
// З'єднання
// ---------------------------------------------------------------------
type clientMessage struct {
Type string `json:"type"`
MapID string `json:"map_id"`
}
func (s *Server) handleWS(w http.ResponseWriter, r *http.Request, tok *store.APIToken) {
if !tok.Can("maps:read") {
writeError(w, http.StatusForbidden, "forbidden", "немає права maps:read")
return
}
opts := &websocket.AcceptOptions{
// Токен приїхав підпротоколом — його треба підтвердити,
// інакше браузер розірве з'єднання одразу після рукостискання.
Subprotocols: websocketProtocols(r),
// Перевірку Origin робить зворотний проксі разом із політикою
// CORS: тут немає списку дозволених доменів, а вгадувати його
// небезпечніше, ніж не перевіряти.
InsecureSkipVerify: true,
}
conn, err := websocket.Accept(w, r, opts)
if err != nil {
s.log.Warn("WebSocket не прийнято", "err", err)
return
}
defer conn.CloseNow()
sub := &subscriber{tenantID: tok.TenantID, send: make(chan []byte, subscriberQueue)}
s.hub.add(sub)
defer s.hub.remove(sub)
ctx, cancel := context.WithCancel(r.Context())
defer cancel()
// Читач живе окремо від писаря: інакше повільний клієнт блокував би
// власні ж оновлення, а закриття з'єднання не помічалося б, доки
// не настане час наступної розсилки.
go func() {
defer cancel()
for {
_, data, err := conn.Read(ctx)
if err != nil {
return
}
var msg clientMessage
if err := json.Unmarshal(data, &msg); err != nil {
continue
}
if msg.Type != "subscribe" {
continue
}
// Мапа має належати тому ж тенанту: інакше вгаданий id
// відкривав би чужу топологію.
ok, err := s.store.MapExists(ctx, tok.TenantID, msg.MapID)
if err != nil || !ok {
_ = writeWS(ctx, conn, map[string]any{
"type": "error", "code": "map_not_found", "map_id": msg.MapID,
})
continue
}
sub.setMap(msg.MapID)
_ = writeWS(ctx, conn, map[string]any{
"type": "subscribed", "map_id": msg.MapID,
})
}
}()
_ = writeWS(ctx, conn, map[string]any{"type": "hello", "tenant_id": tok.TenantID})
ping := time.NewTicker(30 * time.Second)
defer ping.Stop()
for {
select {
case <-ctx.Done():
return
case msg, ok := <-sub.send:
if !ok {
return
}
wctx, cancelWrite := context.WithTimeout(ctx, writeTimeout)
err := conn.Write(wctx, websocket.MessageText, msg)
cancelWrite()
if err != nil {
return
}
case <-ping.C:
// Без пінгу проміжний проксі тихо закриє простояле
// з'єднання, і NOC-екран замре з останнім кадром.
pctx, cancelPing := context.WithTimeout(ctx, writeTimeout)
err := conn.Ping(pctx)
cancelPing()
if err != nil {
return
}
}
}
}
func writeWS(ctx context.Context, conn *websocket.Conn, v any) error {
data, err := json.Marshal(v)
if err != nil {
return err
}
wctx, cancel := context.WithTimeout(ctx, writeTimeout)
defer cancel()
if err := conn.Write(wctx, websocket.MessageText, data); err != nil {
if errors.Is(err, context.Canceled) {
return nil
}
return err
}
return nil
}