PATCH /api/v1/maps/{id} з оптимістичним блокуванням за revision, плюс
створення, видалення й автопобудова з виявленої топології.
Усі скалярні поля патча — вказівники: перетягування шле лише x/y, і якби
відсутні поля означали порожні, кожен рух миші стирав би стиль, розмір і
прив'язку до пристрою.
Ребро може посилатися на вузол, створений тим же патчем, за client_id.
Той, хто спізнився з ревізією, отримує 409, а не тихо затирає чужу правку.
Знімок пишеться тією ж транзакцією, що й зміна, — інакше в історії лишався
б крок, якого в мапі немає.
Автопобудова ідемпотентна: повторний запуск не дублює вузлів і не скидає
ручну розкладку.
Тестами знайдено: revision <= $2 - $3 з двома нетипізованими параметрами
дає "operator is not unique: unknown - unknown" — потрібні явні касти.
Перевірено: 23 інтеграційні тести API (-race), плюс живий прогін проти
даних, зібраних агентом.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
377 lines
10 KiB
Go
377 lines
10 KiB
Go
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)
|
||
}
|
||
}
|
||
|
||
// BroadcastMapUpdate повідомляє інші відкриті полотна, що мапу змінили.
|
||
//
|
||
// Шлемо лише номер ревізії, а не сам патч: клієнт сам вирішить, чи
|
||
// перечитувати полотно. Розсилати зміни як дельти означало б тримати
|
||
// на сервері модель того, що бачить кожен клієнт, — а це вже спільне
|
||
// редагування з CRDT, окрема задача.
|
||
func (h *Hub) BroadcastMapUpdate(tenantID, mapID string, revision int64) {
|
||
msg, err := json.Marshal(map[string]any{
|
||
"type": "map.updated",
|
||
"map_id": mapID,
|
||
"revision": revision,
|
||
"at": time.Now(),
|
||
})
|
||
if err != nil {
|
||
return
|
||
}
|
||
h.broadcast(tenantID, mapID, msg)
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// З'єднання
|
||
// ---------------------------------------------------------------------
|
||
|
||
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
|
||
}
|