diff --git a/HISTORY.md b/HISTORY.md index 351258c..bab9ea1 100644 --- a/HISTORY.md +++ b/HISTORY.md @@ -373,9 +373,70 @@ server/cmd/netpulse-secret/ заведення шифрованих с Це повний шлях даних для анімації трафіку на мапі. Усі три набори тестів проходять з `-race`. +--- + +## 2026-08-14 — Етап 3: REST/WebSocket API для UI + +### Створено + +``` +server/ +├── API.md ендпоїнти, протокол WebSocket, приклади +├── cmd/netpulse-api/ окремий процес: HTTP + WS +└── internal/ + ├── httpapi/server.go роутер, Bearer-автентифікація, обробники + ├── httpapi/ws.go hub, насос подій, насос завантаження каналів + └── store/ + ├── maps.go стан полотна з живими статусами + ├── events.go читання core.event_outbox + ├── inventory.go пристрої, зонди + └── apitokens.go автентифікація токенів UI +``` + +### Прийняті рішення + +1. **API — окремий процес від AgentService.** Зонди й браузери мають різні профілі + навантаження, периметри й цикли релізів. Спільний лише шар `store`. +2. **Стан мапи віддається одним викликом разом зі статусами.** Без цього мапа + малювалася б сірою й доганяла кольори сотнею дозапитів. +3. **`link_status` виводиться з кінців лінка, а не читається з колонки.** + `topo.links.status` ніхто не підтримує; писати туди означало б оновлювати всі + лінки пристрою на кожну зміну статусу. Стан лінка — похідна величина. +4. **Події пишуться тією ж транзакцією, що й зміна.** Інакше WebSocket міг би + розповісти про перехід, якого в базі ще (або вже) немає. +5. **Опитування outbox замість LISTEN/NOTIFY.** NOTIFY не переживає падіння + підписника й обмежений 8 КБ; тут потрібна гарантія, що зміна статусу не + загубиться між перезапусками API. +6. **Дві частоти розсилки:** статус — подія (миттєво), завантаження — величина + (раз на 5 с). Частіше за оновлення лічильників (60 с) — це та сама цифра по колу. +7. **Токен WebSocket їде підпротоколом**, бо браузер не дозволяє довільні + заголовки; в URL він не потрапляє, а отже й у логи проксі. +8. **Підписник, що не встигає читати, відключається**, а не гальмує решту. + +### Перевірено + +10 інтеграційних тестів проти живої БД, справжнього HTTP і WebSocket — усі з `-race`. +Найцінніші: `TestWebSocketDeliversStatusChange` проганяє справжній батч телеметрії +через `applyDeviceStatus` → outbox → hub → браузер; `TestLinkStatusFollowsEndpoints` +перевіряє, що лінк червоніє від падіння кінця або порту без змін у `topo.links`. + +**Живий прогін** — агент, `netpulse-server` і `netpulse-api` разом проти справжнього +`snmpd`: мапа з 2 вузлами `up` (RTT 0.113 і 0.827 мс), ребро з портом `eth0`, +`лінк=up`, живий `util_pct`, зонд `online` з RSS 11.6 МБ і `dropped_samples=0`. + +`target_port` порожній — і це правильно: лінк знайдено через ARP, а ARP не повідомляє +порт віддаленої сторони. + +### Чого ще немає + +- Усі ендпоїнти read-only: редактор мапи потребує `PATCH` з оптимістичним + блокуванням за `revision` (колонка є, обробника немає). +- Немає `GET /api/v1/metrics` для графіків і віддачі `alr.alerts`. +- Автопобудова мапи з `topo.links` — поки SQL-скрипт, не кнопка. + ### Далі +- Запис у мапу (`PATCH`) — без нього немає редактора. +- Фронтенд: React Flow поверх цього API. - `EnrollmentService` + видача сертифікатів зондам. - Планувальник NCM: бекап за cron і за Syslog-подією. -- REST/WebSocket API для UI: у базі вже є все для мапи — вузли, лінки, - статуси й завантаження каналів. diff --git a/server/API.md b/server/API.md new file mode 100644 index 0000000..9956b40 --- /dev/null +++ b/server/API.md @@ -0,0 +1,175 @@ +# NetPulse API — REST і WebSocket для фронтенду + +```bash +go build -o netpulse-api ./cmd/netpulse-api +``` + +```bash +./netpulse-api -listen :8080 -dsn "postgres://..." -cert api.pem -key api.key +``` + +Окремий процес від `netpulse-server` (AgentService) навмисно: зонди й браузери мають +різні профілі навантаження, різні мережеві периметри й різні цикли релізів. Спільним +лишається лише шар `store`, тому дані обидва бачать однакові. + +## Автентифікація + +Bearer-токен із `core.api_tokens`; у БД лежить лише `sha256`, сам токен показується +користувачу один раз при створенні. + +``` +Authorization: Bearer np_ui_xxxxxxxx +``` + +Порожній `scopes` означає повний доступ — так поводяться токени, створені власником +тенанта для себе. Інакше перевіряються права `maps:read`, `devices:read`, `agents:read`. + +`GET /healthz` свідомо відкритий: його опитує балансувальник. + +## Ендпоїнти + +| Метод | Шлях | Призначення | +|-------|------|-------------| +| `GET` | `/healthz` | стан процесу й кількість WebSocket-підписників | +| `GET` | `/api/v1/maps` | перелік мап із лічильниками вузлів і ребер | +| `GET` | `/api/v1/maps/{id}` | **повний стан полотна разом із живими статусами** | +| `GET` | `/api/v1/devices` | інвентар | +| `GET` | `/api/v1/agents` | зонди, версії, самометрики | +| `GET` | `/api/v1/ws` | WebSocket: події та завантаження каналів | + +### `GET /api/v1/maps/{id}` — головний запит продукту + +Одним викликом повертає готове до рендеру полотно: мапу, підкладки, вузли, ребра — +**і живий стан**. Це не оптимізація: без підмішаного статусу мапа малювалася б сірою +й лише потім доганяла кольори сотнею дозапитів. + +```jsonc +{ + "id": "…", "name": "Автомапа", "layout_algo": "manual", + "viewport": {"x":0,"y":0,"zoom":1}, + "grid": {"enabled":true,"size":16,"snap":true}, + "clustering": {"enabled":true,"zoom_threshold":0.4}, + + "backgrounds": [ + {"kind":"image","storage_key":"s3://…/floorplan.svg","opacity":0.6, + "x":0,"y":0,"width":2400,"height":1600,"locked":true} + ], + + "nodes": [ + {"id":"…","kind":"device","label":"core-sw","x":100,"y":200, + "device_id":"…","style":{"icon":"switch"}, + "status":"up","rtt_ms":0.113,"loss_pct":0} // ← фарбує вузол + ], + + "edges": [ + {"id":"…","source_node_id":"…","target_node_id":"…", + "source_port":"Gi0/1","target_port":"ether1","label":"Gi0/1 → ether1", + "style":"smoothstep","waypoints":[], + "animation":{"enabled":true,"speed_source":"utilization"}, + "thresholds":{"warn_pct":70,"crit_pct":90}, + "link_status":"up","util_pct":7.2,"capacity_bps":10000000000} + ] +} +``` + +**`link_status` виводиться з кінців лінка, а не читається з колонки.** У схемі є +`topo.links.status`, але її ніхто не підтримує: писати туди означало б оновлювати всі +лінки пристрою на кожну зміну його статусу й тримати це узгодженим. Стан лінка — +похідна величина, тож рахується на читанні. Порядок гілок важливий: обрив перекриває +все інше, а «невідомо» стоїть перед «up», щоб мапа не малювала зеленим те, чого ще +жодного разу не опитували. + +**`rtt_ms` і `loss_pct` беруться лише за останні 15 хвилин** (`ts.device_last_icmp`). +Старіші дані не характеризують поточний стан, і показувати їх означало б брехати про +живість пристрою. + +## WebSocket + +``` +GET /api/v1/ws +Sec-WebSocket-Protocol: netpulse.token.<токен> +``` + +Браузерний WebSocket API не дозволяє довільні заголовки, тому токен їде підпротоколом. +Сам токен при цьому не потрапляє в URL, а отже і в логи проксі. + +**Клієнт → сервер** + +```json +{"type": "subscribe", "map_id": "…"} +``` + +Мапа перевіряється на належність тенанту: вгаданий id не відкриє чужу топологію. + +**Сервер → клієнт** + +| `type` | Коли | Вміст | +|--------|------|-------| +| `hello` | одразу після рукостискання | `tenant_id` | +| `subscribed` | підтвердження підписки | `map_id` | +| `error` | напр. чужа мапа | `code`, `map_id` | +| `device.status` | зміна статусу пристрою | `device_id`, `status`, `previous_status`, `reason` | +| `link.load` | кожні 5 с для підписаної мапи | `map_id`, `links: [{link_id, status, util_pct}]` | + +Дві частоти навмисно різні. Зміна статусу — **подія**: рідка, але має дійти майже +миттєво, інакше мапа бреше про стан мережі. Завантаження каналу — **величина**: вона +змінюється весь час, і слати її частіше, ніж оновлюються лічильники (60 с), означає +слати ту саму цифру по колу. + +### Звідки беруться події + +`core.event_outbox` наповнюється **тією ж транзакцією**, що й зміна, яку описує. +Інакше WebSocket міг би розповісти про перехід, якого в базі ще (або вже) немає. + +Транспорт — опитування таблиці раз на секунду, а не `LISTEN/NOTIFY`. NOTIFY не +переживає падіння підписника й обмежений 8 КБ на повідомлення, а тут потрібна +гарантія, що жодна зміна статусу не загубиться між перезапусками API. Ціна — один +дешевий запит за індексом. + +Нова сесія починає з кінця журналу: клієнт щойно завантажив повний стан мапи, і все +старіше в ньому вже враховано. + +Підписник, який не встигає читати, **відключається**, а не сповільнює решту: тримати +його чергу означало б віддавати пам'ять і затримувати всіх інших. + +## Стан перевірки + +Тести працюють проти справжньої БД, справжнього HTTP і справжнього WebSocket +(`NETPULSE_TEST_DSN`; без змінної пропускаються). `go test -race` — усі проходять. + +| Тест | Що доводить | +|------|-------------| +| `TestAuthRequired` | без токена й з чужим токеном — 401; `healthz` відкритий | +| `TestMapStateIsRenderReady` | одним викликом приходять підкладка, координати, статус вузла, RTT, порти ребра, `util_pct`, `capacity_bps`, налаштування анімації | +| `TestLinkStatusFollowsEndpoints` | лінк червоніє, коли впав кінець або порт, — без жодних змін у `topo.links` | +| `TestTenantIsolation` | чужа мапа за точним id → 404; у переліку пристроїв немає чужого тенанта | +| `TestBadMapID` | некоректний uuid → 400, а не 500 | +| `TestWebSocketRequiresToken` | без токена з'єднання не відкривається | +| `TestWebSocketDeliversStatusChange` | справжній батч телеметрії → `applyDeviceStatus` → outbox → hub → браузер отримав `device.status` із `previous_status` | +| `TestWebSocketRejectsForeignMap` | підписка на чужу мапу відхилена | +| `TestWebSocketPushesLinkLoads` | періодичний `link.load` із реальним `util_pct` | + +### Живий прогін + +Агент, `netpulse-server` і `netpulse-api` запущені разом проти справжнього `snmpd`: + +``` +мапа: Автомапа | вузлів: 2 | ребер: 1 +вузол snmp-host статус=up rtt=0.113 мс loss=0 +вузол gateway статус=up rtt=0.827 мс loss=0 +ребро eth0 → ? порти=eth0->None лінк=up util=7.0e-06% capacity=10000000000 +зонд probe-snmp: online, linux/amd64, пристроїв=2, + health={rss_bytes: 11624464, dropped_samples: 0, clock_skew_ms: 0} +``` + +`target_port` порожній — і це правильно: лінк знайдено через ARP, а ARP не повідомляє +порт віддаленої сторони. Заповнить його LLDP, коли поруч буде обладнання, що його шле. + +## Чого ще немає + +- **Запис**: усі ендпоїнти read-only. Редактор мапи (перетягування вузлів, малювання + зв'язків) потребує `PATCH /api/v1/maps/{id}` з оптимістичним блокуванням за + `revision` — колонка вже є, обробника ще немає. +- **Історія метрик**: `GET /api/v1/metrics` для графіків із `ts.samples_5m` не написано. +- **Алерти**: `alr.alerts` не віддаються, хоча схема готова. +- **Автопобудова мапи** з `topo.links` робиться поки SQL-скриптом, не кнопкою. diff --git a/server/cmd/netpulse-api/main.go b/server/cmd/netpulse-api/main.go new file mode 100644 index 0000000..96f3d04 --- /dev/null +++ b/server/cmd/netpulse-api/main.go @@ -0,0 +1,145 @@ +// Команда netpulse-api — REST і WebSocket для фронтенду. +// +// Окремий процес від netpulse-server (AgentService): зонди й браузери +// мають різні профілі навантаження й різні мережеві периметри. Спільний +// у них лише шар store, тому дані обидва бачать однакові. +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "log/slog" + "net/http" + "os" + "os/signal" + "syscall" + "time" + + "github.com/netpulse/netpulse/server/internal/httpapi" + "github.com/netpulse/netpulse/server/internal/store" +) + +var version = "dev" + +func main() { + if err := run(); err != nil { + fmt.Fprintln(os.Stderr, "netpulse-api:", err) + os.Exit(1) + } +} + +func run() error { + var ( + listen = flag.String("listen", envOr("NETPULSE_API_LISTEN", ":8080"), "адреса HTTP") + dsn = flag.String("dsn", os.Getenv("NETPULSE_DSN"), "DSN PostgreSQL") + certFile = flag.String("cert", os.Getenv("NETPULSE_API_CERT"), "сертифікат TLS") + keyFile = flag.String("key", os.Getenv("NETPULSE_API_KEY"), "приватний ключ TLS") + logLevel = flag.String("log-level", envOr("NETPULSE_LOG_LEVEL", "info"), "debug|info|warn|error") + pruneAge = flag.Duration("event-retention", 24*time.Hour, + "скільки тримати доставлені події в core.event_outbox") + ) + flag.Parse() + + if *dsn == "" { + return errors.New("не вказано -dsn (або NETPULSE_DSN)") + } + + log := newLogger(*logLevel) + + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + + st, err := store.New(ctx, *dsn) + if err != nil { + return fmt.Errorf("підключення до БД: %w", err) + } + defer st.Close() + + api := httpapi.New(st, log) + + go api.Hub().Run(ctx) + go pruneLoop(ctx, st, log, *pruneAge) + + srv := &http.Server{ + Addr: *listen, + Handler: api.Handler(), + // Читання й заголовки обмежені, а от загального WriteTimeout + // немає навмисно: він рубав би довгі WebSocket-з'єднання + // рівно посеред роботи NOC-екрана. + ReadHeaderTimeout: 10 * time.Second, + ReadTimeout: 30 * time.Second, + IdleTimeout: 120 * time.Second, + } + + log.Info("запуск", "version", version, "listen", *listen, "tls", *certFile != "") + + errCh := make(chan error, 1) + go func() { + if *certFile != "" && *keyFile != "" { + errCh <- srv.ListenAndServeTLS(*certFile, *keyFile) + return + } + log.Warn("запуск без TLS — припустимо лише за зворотним проксі") + errCh <- srv.ListenAndServe() + }() + + select { + case <-ctx.Done(): + log.Info("зупинка", "підписників", api.Hub().SubscriberCount()) + shutCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + return srv.Shutdown(shutCtx) + case err := <-errCh: + if errors.Is(err, http.ErrServerClosed) { + return nil + } + return err + } +} + +// pruneLoop прибирає доставлені події. +// +// Без нього core.event_outbox росла б вічно: подій там небагато, але +// «небагато» помножене на роки — це той самий терабайт. +func pruneLoop(ctx context.Context, st *store.Store, log *slog.Logger, age time.Duration) { + t := time.NewTicker(time.Hour) + defer t.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-t.C: + n, err := st.PruneEvents(ctx, age) + if err != nil { + log.Warn("прибирання журналу подій", "err", err) + continue + } + if n > 0 { + log.Info("прибрано подій", "рядків", n) + } + } + } +} + +func envOr(key, def string) string { + if v := os.Getenv(key); v != "" { + return v + } + return def +} + +func newLogger(level string) *slog.Logger { + lv := slog.LevelInfo + switch level { + case "debug": + lv = slog.LevelDebug + case "warn": + lv = slog.LevelWarn + case "error": + lv = slog.LevelError + } + return slog.New(slog.NewJSONHandler(os.Stderr, &slog.HandlerOptions{Level: lv})) +} diff --git a/server/go.mod b/server/go.mod index 6cd95f0..43a5c55 100644 --- a/server/go.mod +++ b/server/go.mod @@ -3,6 +3,7 @@ module github.com/netpulse/netpulse/server go 1.25.0 require ( + github.com/coder/websocket v1.8.15 github.com/jackc/pgx/v5 v5.7.6 github.com/netpulse/netpulse/gen/go v0.0.0 google.golang.org/grpc v1.83.0 diff --git a/server/go.sum b/server/go.sum index a9cc096..50c7f87 100644 --- a/server/go.sum +++ b/server/go.sum @@ -1,5 +1,7 @@ github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/coder/websocket v1.8.15 h1:6B2JPeOGlpff2Uz6vOEH1Vzpi0iUz20A+lPVhPHtNUA= +github.com/coder/websocket v1.8.15/go.mod h1:NX3SzP+inril6yawo5CQXx8+fk145lPDC6pumgx0mVg= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= diff --git a/server/internal/httpapi/api_test.go b/server/internal/httpapi/api_test.go new file mode 100644 index 0000000..28ebf0e --- /dev/null +++ b/server/internal/httpapi/api_test.go @@ -0,0 +1,586 @@ +// Наскрізний тест API: справжня БД, справжній HTTP, справжній WebSocket. +// +// NETPULSE_TEST_DSN="postgres://netpulse:netpulse@localhost/netpulse_it" go test ./... +package httpapi_test + +import ( + "context" + "crypto/sha256" + "encoding/json" + "fmt" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "os" + "strings" + "testing" + "time" + + "github.com/coder/websocket" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/netpulse/netpulse/server/internal/httpapi" + "github.com/netpulse/netpulse/server/internal/store" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" + "google.golang.org/protobuf/types/known/timestamppb" +) + +type fixture struct { + pool *pgxpool.Pool + store *store.Store + srv *httptest.Server + ctx context.Context + + tenantID string + token string + agentID string + deviceID string + peerID string + ifaceID string + peerIfID string + linkID string + mapID string + nodeID string +} + +func setup(t *testing.T) *fixture { + t.Helper() + + dsn := os.Getenv("NETPULSE_TEST_DSN") + if dsn == "" { + t.Skip("NETPULSE_TEST_DSN не задано — інтеграційний тест пропущено") + } + + ctx := context.Background() + st, err := store.New(ctx, dsn) + if err != nil { + t.Fatalf("БД: %v", err) + } + t.Cleanup(st.Close) + + f := &fixture{pool: st.Pool(), store: st, ctx: ctx} + f.seed(t) + + api := httpapi.New(st, slog.New(slog.NewTextHandler(io.Discard, nil))) + + hubCtx, cancel := context.WithCancel(ctx) + go api.Hub().Run(hubCtx) + t.Cleanup(cancel) + + f.srv = httptest.NewServer(api.Handler()) + t.Cleanup(f.srv.Close) + + return f +} + +func (f *fixture) seed(t *testing.T) { + t.Helper() + ctx := f.ctx + + slug := fmt.Sprintf("api-%d", time.Now().UnixNano()) + f.token = "np_api_" + slug + sum := sha256.Sum256([]byte(f.token)) + agentToken := sha256.Sum256([]byte("agent_" + slug)) + + scan := func(dest *string, q string, args ...any) { + t.Helper() + if err := f.pool.QueryRow(ctx, q, args...).Scan(dest); err != nil { + t.Fatalf("seed %q: %v", strings.TrimSpace(q)[:40], err) + } + } + exec := func(q string, args ...any) { + t.Helper() + if _, err := f.pool.Exec(ctx, q, args...); err != nil { + t.Fatalf("seed %q: %v", strings.TrimSpace(q)[:40], err) + } + } + + scan(&f.tenantID, `INSERT INTO core.tenants (slug, name, status) + VALUES ($1, $2, 'active') RETURNING id::text`, slug, slug) + + t.Cleanup(func() { + bg := context.Background() + _, _ = f.pool.Exec(bg, `DELETE FROM ts.icmp_samples WHERE tenant_id = $1`, f.tenantID) + _, _ = f.pool.Exec(bg, `DELETE FROM ts.if_counters WHERE tenant_id = $1`, f.tenantID) + _, _ = f.pool.Exec(bg, `DELETE FROM ts.device_status_history WHERE tenant_id = $1`, f.tenantID) + _, _ = f.pool.Exec(bg, `DELETE FROM core.tenants WHERE id = $1`, f.tenantID) + }) + + scan(&f.agentID, `INSERT INTO core.agents (tenant_id, name, token_hash, status) + VALUES ($1, 'probe-api', $2, 'online') RETURNING id::text`, f.tenantID, agentToken[:]) + + exec(`INSERT INTO core.api_tokens (tenant_id, name, prefix, token_hash, scopes) + VALUES ($1, 'ui', 'np_api_', $2, '{}')`, f.tenantID, sum[:]) + + scan(&f.deviceID, `INSERT INTO inv.devices (tenant_id, agent_id, name, address, kind, status) + VALUES ($1, $2, 'core-sw', '10.20.0.1', 'switch', 'up') RETURNING id::text`, + f.tenantID, f.agentID) + scan(&f.peerID, `INSERT INTO inv.devices (tenant_id, agent_id, name, address, kind, status) + VALUES ($1, $2, 'edge-rtr', '10.20.0.2', 'router', 'up') RETURNING id::text`, + f.tenantID, f.agentID) + + scan(&f.ifaceID, `INSERT INTO inv.interfaces (tenant_id, device_id, if_index, name, speed_bps, oper_status) + VALUES ($1, $2, 1, 'Gi0/1', 1000000000, 'up') RETURNING id::text`, f.tenantID, f.deviceID) + scan(&f.peerIfID, `INSERT INTO inv.interfaces (tenant_id, device_id, if_index, name, speed_bps, oper_status) + VALUES ($1, $2, 1, 'ether1', 1000000000, 'up') RETURNING id::text`, f.tenantID, f.peerID) + + scan(&f.linkID, `INSERT INTO topo.links + (tenant_id, a_device_id, a_interface_id, b_device_id, b_interface_id, + kind, capacity_bps, discovered_by, status) + VALUES ($1,$2,$3,$4,$5,'physical',1000000000,'lldp','up') RETURNING id::text`, + f.tenantID, f.deviceID, f.ifaceID, f.peerID, f.peerIfID) + + scan(&f.mapID, `INSERT INTO topo.maps (tenant_id, name, slug, kind) + VALUES ($1, 'NOC', 'noc', 'logical') RETURNING id::text`, f.tenantID) + + exec(`INSERT INTO topo.map_backgrounds + (tenant_id, map_id, kind, storage_key, mime_type, width, height, opacity) + VALUES ($1,$2,'image','s3://plan.svg','image/svg+xml',2400,1600,0.6)`, + f.tenantID, f.mapID) + + scan(&f.nodeID, `INSERT INTO topo.map_nodes (tenant_id, map_id, kind, device_id, label, x, y) + VALUES ($1,$2,'device',$3,'core-sw',100,200) RETURNING id::text`, + f.tenantID, f.mapID, f.deviceID) + + var peerNode string + scan(&peerNode, `INSERT INTO topo.map_nodes (tenant_id, map_id, kind, device_id, label, x, y) + VALUES ($1,$2,'device',$3,'edge-rtr',400,200) RETURNING id::text`, + f.tenantID, f.mapID, f.peerID) + + exec(`INSERT INTO topo.map_edges + (tenant_id, map_id, source_node_id, target_node_id, + source_interface_id, target_interface_id, link_id, label) + VALUES ($1,$2,$3,$4,$5,$6,$7,'Gi0/1 → ether1')`, + f.tenantID, f.mapID, f.nodeID, peerNode, f.ifaceID, f.peerIfID, f.linkID) + + // Живі лічильники: без них у ребра не буде util_pct, тобто нічим + // керувати анімацією. + exec(`INSERT INTO ts.if_counters + (ts, interface_id, device_id, tenant_id, in_bps, out_bps, util_in_pct, util_out_pct, oper_up) + VALUES (now(), $1, $2, $3, 420e6, 780e6, 42, 78, true)`, + f.ifaceID, f.deviceID, f.tenantID) + + exec(`INSERT INTO ts.icmp_samples + (ts, device_id, tenant_id, agent_id, rtt_avg_ms, loss_pct, packets_sent, packets_recv, reachable) + VALUES (now(), $1, $2, $3, 1.5, 0, 3, 3, true)`, + f.deviceID, f.tenantID, f.agentID) +} + +func (f *fixture) get(t *testing.T, path, token string) (int, []byte) { + t.Helper() + req, err := http.NewRequest(http.MethodGet, f.srv.URL+path, nil) + if err != nil { + t.Fatalf("запит: %v", err) + } + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + resp, err := f.srv.Client().Do(req) + if err != nil { + t.Fatalf("виклик %s: %v", path, err) + } + defer resp.Body.Close() + body, _ := io.ReadAll(resp.Body) + return resp.StatusCode, body +} + +// --------------------------------------------------------------------- + +func TestAuthRequired(t *testing.T) { + f := setup(t) + + if code, _ := f.get(t, "/api/v1/maps", ""); code != http.StatusUnauthorized { + t.Fatalf("без токена код %d", code) + } + if code, _ := f.get(t, "/api/v1/maps", "np_api_нема-такого"); code != http.StatusUnauthorized { + t.Fatalf("з чужим токеном код %d", code) + } + // healthz свідомо відкритий: його опитує балансувальник. + if code, _ := f.get(t, "/healthz", ""); code != http.StatusOK { + t.Fatalf("healthz код %d", code) + } +} + +func TestListMaps(t *testing.T) { + f := setup(t) + + code, body := f.get(t, "/api/v1/maps", f.token) + if code != http.StatusOK { + t.Fatalf("код %d: %s", code, body) + } + + var out struct { + Maps []store.MapSummary `json:"maps"` + } + if err := json.Unmarshal(body, &out); err != nil { + t.Fatalf("розбір: %v", err) + } + if len(out.Maps) != 1 { + t.Fatalf("мап: %d", len(out.Maps)) + } + m := out.Maps[0] + if m.Name != "NOC" || m.NodeCount != 2 || m.EdgeCount != 1 { + t.Fatalf("зведення мапи неправильне: %+v", m) + } +} + +// Головний запит продукту: одним викликом фронтенд має отримати готове +// до рендеру полотно разом із живими статусами — інакше мапа малювалася б +// сірою й лише потім доганяла кольори сотнею дозапитів. +func TestMapStateIsRenderReady(t *testing.T) { + f := setup(t) + + code, body := f.get(t, "/api/v1/maps/"+f.mapID, f.token) + if code != http.StatusOK { + t.Fatalf("код %d: %s", code, body) + } + + var st store.MapState + if err := json.Unmarshal(body, &st); err != nil { + t.Fatalf("розбір: %v", err) + } + + if st.Name != "NOC" || st.LayoutAlgo == "" { + t.Fatalf("мапа: %+v", st) + } + if len(st.Viewport) == 0 || len(st.Grid) == 0 { + t.Fatal("немає viewport/grid — полотно нема куди відновити") + } + + if len(st.Backgrounds) != 1 || st.Backgrounds[0].StorageKey == "" { + t.Fatalf("підкладки: %+v", st.Backgrounds) + } + if st.Backgrounds[0].Opacity != 0.6 { + t.Fatalf("прозорість підкладки: %v", st.Backgrounds[0].Opacity) + } + + if len(st.Nodes) != 2 { + t.Fatalf("вузлів: %d", len(st.Nodes)) + } + var found *store.MapNode + for i := range st.Nodes { + if st.Nodes[i].DeviceID == f.deviceID { + found = &st.Nodes[i] + } + } + if found == nil { + t.Fatal("вузол пристрою не повернувся") + } + if found.X != 100 || found.Y != 200 { + t.Fatalf("координати вузла: %v,%v", found.X, found.Y) + } + // Саме це фарбує вузол. + if found.Status != "up" { + t.Fatalf("статус вузла: %q", found.Status) + } + if found.RttMs == nil || *found.RttMs < 1.4 || *found.RttMs > 1.6 { + t.Fatalf("RTT не підмішався: %v", found.RttMs) + } + + if len(st.Edges) != 1 { + t.Fatalf("ребер: %d", len(st.Edges)) + } + e := st.Edges[0] + if e.SourcePort != "Gi0/1" || e.TargetPort != "ether1" { + t.Fatalf("порти ребра: %q → %q", e.SourcePort, e.TargetPort) + } + // Стан лінка виводиться з кінців, а не читається з колонки: + // обидва пристрої up і порти up → лінк up. + if e.LinkStatus != "up" { + t.Fatalf("статус лінка: %q", e.LinkStatus) + } + // Джерело швидкості анімації. + if e.UtilPct == nil || *e.UtilPct != 78 { + t.Fatalf("util_pct: %v", e.UtilPct) + } + if e.CapacityBps == nil || *e.CapacityBps != 1000000000 { + t.Fatalf("capacity_bps: %v", e.CapacityBps) + } + if len(e.Animation) == 0 || len(e.Thresholds) == 0 { + t.Fatal("немає налаштувань анімації/порогів") + } +} + +// Стан лінка — похідна від його кінців. Колонка topo.links.status +// ніким не підтримується, тому читати її означало б завжди показувати +// "unknown" і брехати про обрив. +func TestLinkStatusFollowsEndpoints(t *testing.T) { + f := setup(t) + + edgeStatus := func() string { + t.Helper() + code, body := f.get(t, "/api/v1/maps/"+f.mapID, f.token) + if code != http.StatusOK { + t.Fatalf("код %d: %s", code, body) + } + var st store.MapState + if err := json.Unmarshal(body, &st); err != nil { + t.Fatalf("розбір: %v", err) + } + if len(st.Edges) != 1 { + t.Fatalf("ребер: %d", len(st.Edges)) + } + return st.Edges[0].LinkStatus + } + + if got := edgeStatus(); got != "up" { + t.Fatalf("обидва кінці живі, а лінк %q", got) + } + + // Один кінець упав — лінк має почервоніти, хоча в topo.links + // нічого не змінювалось. + if _, err := f.pool.Exec(f.ctx, + `UPDATE inv.devices SET status = 'down' WHERE id = $1`, f.peerID); err != nil { + t.Fatalf("зміна статусу: %v", err) + } + if got := edgeStatus(); got != "down" { + t.Fatalf("кінець лежить, а лінк %q", got) + } + + // Порт адміністративно вимкнено — теж обрив. + if _, err := f.pool.Exec(f.ctx, + `UPDATE inv.devices SET status = 'up' WHERE id = $1`, f.peerID); err != nil { + t.Fatalf("відновлення: %v", err) + } + if _, err := f.pool.Exec(f.ctx, + `UPDATE inv.interfaces SET oper_status = 'down' WHERE id = $1`, f.ifaceID); err != nil { + t.Fatalf("зміна порту: %v", err) + } + if got := edgeStatus(); got != "down" { + t.Fatalf("порт лежить, а лінк %q", got) + } +} + +// Токен одного тенанта не має відкривати мапу іншого навіть за точним id. +func TestTenantIsolation(t *testing.T) { + f := setup(t) + other := setup(t) + + code, body := f.get(t, "/api/v1/maps/"+other.mapID, f.token) + if code != http.StatusNotFound { + t.Fatalf("чужа мапа віддалась із кодом %d: %s", code, body) + } + + code, body = f.get(t, "/api/v1/devices", f.token) + if code != http.StatusOK { + t.Fatalf("код %d", code) + } + if strings.Contains(string(body), other.deviceID) { + t.Fatal("у переліку пристроїв видно чужий тенант") + } +} + +func TestBadMapID(t *testing.T) { + f := setup(t) + + if code, _ := f.get(t, "/api/v1/maps/не-uuid", f.token); code != http.StatusBadRequest { + t.Fatalf("некоректний id дав код %d, очікували 400", code) + } +} + +func TestListDevicesAndAgents(t *testing.T) { + f := setup(t) + + code, body := f.get(t, "/api/v1/devices", f.token) + if code != http.StatusOK { + t.Fatalf("код %d", code) + } + var dev struct { + Devices []store.DeviceSummary `json:"devices"` + } + if err := json.Unmarshal(body, &dev); err != nil { + t.Fatalf("розбір: %v", err) + } + if len(dev.Devices) != 2 { + t.Fatalf("пристроїв: %d", len(dev.Devices)) + } + for _, d := range dev.Devices { + if d.Name == "core-sw" && d.IfaceCount != 1 { + t.Fatalf("лічильник інтерфейсів: %d", d.IfaceCount) + } + } + + code, body = f.get(t, "/api/v1/agents", f.token) + if code != http.StatusOK { + t.Fatalf("код %d", code) + } + var ag struct { + Agents []store.AgentSummary `json:"agents"` + } + if err := json.Unmarshal(body, &ag); err != nil { + t.Fatalf("розбір: %v", err) + } + if len(ag.Agents) != 1 || ag.Agents[0].DeviceCount != 2 { + t.Fatalf("зонди: %+v", ag.Agents) + } +} + +// --------------------------------------------------------------------- +// WebSocket +// --------------------------------------------------------------------- + +func dialWS(t *testing.T, ctx context.Context, f *fixture, token string) *websocket.Conn { + t.Helper() + + url := "ws" + strings.TrimPrefix(f.srv.URL, "http") + "/api/v1/ws" + conn, _, err := websocket.Dial(ctx, url, &websocket.DialOptions{ + HTTPClient: f.srv.Client(), + // Браузерний WebSocket не дозволяє довільні заголовки, тому + // токен їде підпротоколом — перевіряємо саме цей шлях. + Subprotocols: []string{"netpulse.token." + token}, + }) + if err != nil { + t.Fatalf("WebSocket: %v", err) + } + t.Cleanup(func() { conn.CloseNow() }) + return conn +} + +func readMsg(t *testing.T, ctx context.Context, conn *websocket.Conn) map[string]any { + t.Helper() + _, data, err := conn.Read(ctx) + if err != nil { + t.Fatalf("читання WebSocket: %v", err) + } + var m map[string]any + if err := json.Unmarshal(data, &m); err != nil { + t.Fatalf("розбір повідомлення: %v", err) + } + return m +} + +// waitFor читає, доки не побачить повідомлення потрібного типу. +func waitFor(t *testing.T, ctx context.Context, conn *websocket.Conn, want string) map[string]any { + t.Helper() + for { + m := readMsg(t, ctx, conn) + if m["type"] == want { + return m + } + if m["type"] == "error" { + t.Fatalf("сервер повернув помилку: %v", m) + } + } +} + +func TestWebSocketRequiresToken(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.ctx, 10*time.Second) + defer cancel() + + url := "ws" + strings.TrimPrefix(f.srv.URL, "http") + "/api/v1/ws" + if conn, _, err := websocket.Dial(ctx, url, &websocket.DialOptions{ + HTTPClient: f.srv.Client(), + }); err == nil { + conn.CloseNow() + t.Fatal("WebSocket відкрився без токена") + } +} + +// Зміна статусу пристрою має долетіти до браузера сама — заради цього +// й існує вся зв'язка outbox → hub → WebSocket. +func TestWebSocketDeliversStatusChange(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.ctx, 30*time.Second) + defer cancel() + + conn := dialWS(t, ctx, f, f.token) + waitFor(t, ctx, conn, "hello") + + if err := conn.Write(ctx, websocket.MessageText, + []byte(`{"type":"subscribe","map_id":"`+f.mapID+`"}`)); err != nil { + t.Fatalf("підписка: %v", err) + } + sub := waitFor(t, ctx, conn, "subscribed") + if sub["map_id"] != f.mapID { + t.Fatalf("підписка на іншу мапу: %v", sub) + } + + // Проганяємо справжній шлях: батч телеметрії з недосяжним + // пристроєм → зміна статусу → подія в outbox → трансляція. + agent := &store.Agent{ID: f.agentID, TenantID: f.tenantID} + table := store.NewSeriesTable() + if _, err := f.store.WriteBatch(ctx, agent, &npv1.TelemetryBatch{ + BatchId: 1, AgentId: f.agentID, + Icmp: []*npv1.IcmpResult{{ + DeviceId: f.deviceID, Ts: timestamppb.Now(), + LossPct: 100, PacketsSent: 3, PacketsRecv: 0, Reachable: false, + }}, + }, table); err != nil { + t.Fatalf("запис батчу: %v", err) + } + + msg := waitFor(t, ctx, conn, "device.status") + payload, ok := msg["payload"].(map[string]any) + if !ok { + t.Fatalf("подія без корисного навантаження: %v", msg) + } + if payload["device_id"] != f.deviceID { + t.Fatalf("подія про інший пристрій: %v", payload) + } + if payload["status"] != "down" { + t.Fatalf("статус у події: %v", payload["status"]) + } + if payload["previous_status"] != "up" { + t.Fatalf("попередній статус: %v", payload["previous_status"]) + } +} + +// Підписка на чужу мапу за вгаданим id не має відкривати топологію. +func TestWebSocketRejectsForeignMap(t *testing.T) { + f := setup(t) + other := setup(t) + + ctx, cancel := context.WithTimeout(f.ctx, 15*time.Second) + defer cancel() + + conn := dialWS(t, ctx, f, f.token) + waitFor(t, ctx, conn, "hello") + + if err := conn.Write(ctx, websocket.MessageText, + []byte(`{"type":"subscribe","map_id":"`+other.mapID+`"}`)); err != nil { + t.Fatalf("підписка: %v", err) + } + + m := readMsg(t, ctx, conn) + if m["type"] != "error" || m["code"] != "map_not_found" { + t.Fatalf("чужа мапа прийнята до підписки: %v", m) + } +} + +// Завантаження каналів штовхається періодично — це те, що рухає +// анімацію без перемальовування полотна. +func TestWebSocketPushesLinkLoads(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.ctx, 30*time.Second) + defer cancel() + + conn := dialWS(t, ctx, f, f.token) + waitFor(t, ctx, conn, "hello") + + if err := conn.Write(ctx, websocket.MessageText, + []byte(`{"type":"subscribe","map_id":"`+f.mapID+`"}`)); err != nil { + t.Fatalf("підписка: %v", err) + } + waitFor(t, ctx, conn, "subscribed") + + msg := waitFor(t, ctx, conn, "link.load") + if msg["map_id"] != f.mapID { + t.Fatalf("оновлення для іншої мапи: %v", msg) + } + + links, ok := msg["links"].([]any) + if !ok || len(links) == 0 { + t.Fatalf("немає лінків у оновленні: %v", msg) + } + first, _ := links[0].(map[string]any) + if first["link_id"] != f.linkID { + t.Fatalf("інший лінк: %v", first) + } + if util, ok := first["util_pct"].(float64); !ok || util != 78 { + t.Fatalf("завантаження не доїхало: %v", first["util_pct"]) + } +} diff --git a/server/internal/httpapi/server.go b/server/internal/httpapi/server.go new file mode 100644 index 0000000..6015874 --- /dev/null +++ b/server/internal/httpapi/server.go @@ -0,0 +1,273 @@ +// Package httpapi — REST і WebSocket для фронтенду. +// +// Окремий процес від AgentService навмисно: зонди й браузери мають різні +// профілі навантаження, різні мережеві периметри й різні цикли релізів. +// Спільним лишається лише шар store, тому обидва бачать однакові дані. +package httpapi + +import ( + "encoding/json" + "errors" + "log/slog" + "net/http" + "strings" + "time" + + "github.com/netpulse/netpulse/server/internal/store" +) + +type ctxKey string + +const tokenCtxKey ctxKey = "netpulse.token" + +type Server struct { + store *store.Store + log *slog.Logger + hub *Hub +} + +func New(st *store.Store, log *slog.Logger) *Server { + if log == nil { + log = slog.Default() + } + s := &Server{store: st, log: log} + s.hub = NewHub(st, log) + return s +} + +// Hub — доступ до трансляції для зовнішнього коду (тести, метрики). +func (s *Server) Hub() *Hub { return s.hub } + +// Handler збирає маршрути. +// +// Роутер стандартної бібліотеки: Go 1.22 вміє шаблони з методом і +// параметрами шляху, і цього тут вистачає. Зовнішній роутер додав би +// залежність заради синтаксису. +func (s *Server) Handler() http.Handler { + mux := http.NewServeMux() + + mux.HandleFunc("GET /healthz", s.handleHealth) + + mux.Handle("GET /api/v1/maps", s.authenticated(s.handleListMaps)) + mux.Handle("GET /api/v1/maps/{id}", s.authenticated(s.handleGetMap)) + mux.Handle("GET /api/v1/devices", s.authenticated(s.handleListDevices)) + mux.Handle("GET /api/v1/agents", s.authenticated(s.handleListAgents)) + + // WebSocket теж під автентифікацією: браузер шле токен у + // заголовку через підпротокол — див. ws.go. + mux.Handle("GET /api/v1/ws", s.authenticated(s.handleWS)) + + return s.withRecovery(s.withLogging(mux)) +} + +// --------------------------------------------------------------------- +// Проміжні шари +// --------------------------------------------------------------------- + +func (s *Server) authenticated(next func(http.ResponseWriter, *http.Request, *store.APIToken)) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + raw := bearerToken(r) + if raw == "" { + writeError(w, http.StatusUnauthorized, "no_token", "потрібен Bearer-токен") + return + } + + tok, err := s.store.AuthenticateAPIToken(r.Context(), raw) + if errors.Is(err, store.ErrTokenInvalid) { + writeError(w, http.StatusUnauthorized, "invalid_token", "токен недійсний") + return + } + if err != nil { + s.log.Error("автентифікація", "err", err) + writeError(w, http.StatusInternalServerError, "internal", "внутрішня помилка") + return + } + + next(w, r, tok) + }) +} + +// bearerToken дістає токен із заголовка або з підпротоколу WebSocket. +// +// Другий шлях потрібен тому, що браузерний WebSocket API не дозволяє +// задати довільні заголовки — токен передається як підпротокол +// "netpulse.token.". Це загальноприйнятий обхід; сам токен при +// цьому не потрапляє в URL, а отже і в логи проксі. +func bearerToken(r *http.Request) string { + if v := r.Header.Get("Authorization"); v != "" { + if after, ok := strings.CutPrefix(v, "Bearer "); ok { + return strings.TrimSpace(after) + } + } + for _, proto := range websocketProtocols(r) { + if after, ok := strings.CutPrefix(proto, "netpulse.token."); ok { + return after + } + } + return "" +} + +func websocketProtocols(r *http.Request) []string { + raw := r.Header.Get("Sec-WebSocket-Protocol") + if raw == "" { + return nil + } + parts := strings.Split(raw, ",") + for i := range parts { + parts[i] = strings.TrimSpace(parts[i]) + } + return parts +} + +func (s *Server) withLogging(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + start := time.Now() + rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK} + next.ServeHTTP(rec, r) + + level := slog.LevelDebug + if rec.status >= 500 { + level = slog.LevelError + } + s.log.Log(r.Context(), level, "http", + "method", r.Method, "path", r.URL.Path, + "status", rec.status, "ms", time.Since(start).Milliseconds()) + }) +} + +// withRecovery не дає паніці в одному запиті вбити весь процес разом +// із живими WebSocket-підписниками. +func (s *Server) withRecovery(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + defer func() { + if v := recover(); v != nil { + s.log.Error("паніка в обробнику", "path", r.URL.Path, "panic", v) + writeError(w, http.StatusInternalServerError, "internal", "внутрішня помилка") + } + }() + next.ServeHTTP(w, r) + }) +} + +type statusRecorder struct { + http.ResponseWriter + status int + written bool +} + +func (r *statusRecorder) WriteHeader(code int) { + if r.written { + return + } + r.written = true + r.status = code + r.ResponseWriter.WriteHeader(code) +} + +// Hijack потрібен, бо WebSocket перехоплює з'єднання, а обгортка +// логування стоїть у ланцюжку вище. +func (r *statusRecorder) Unwrap() http.ResponseWriter { return r.ResponseWriter } + +// --------------------------------------------------------------------- +// Обробники +// --------------------------------------------------------------------- + +func (s *Server) handleHealth(w http.ResponseWriter, r *http.Request) { + writeJSON(w, http.StatusOK, map[string]any{ + "status": "ok", + "subscribers": s.hub.SubscriberCount(), + }) +} + +func (s *Server) handleListMaps(w http.ResponseWriter, r *http.Request, tok *store.APIToken) { + if !tok.Can("maps:read") { + writeError(w, http.StatusForbidden, "forbidden", "немає права maps:read") + return + } + + maps, err := s.store.ListMaps(r.Context(), tok.TenantID) + if err != nil { + s.log.Error("перелік мап", "err", err) + writeError(w, http.StatusInternalServerError, "internal", "внутрішня помилка") + return + } + if maps == nil { + maps = []store.MapSummary{} + } + writeJSON(w, http.StatusOK, map[string]any{"maps": maps}) +} + +func (s *Server) handleGetMap(w http.ResponseWriter, r *http.Request, tok *store.APIToken) { + if !tok.Can("maps:read") { + writeError(w, http.StatusForbidden, "forbidden", "немає права maps:read") + return + } + + state, err := s.store.GetMapState(r.Context(), tok.TenantID, r.PathValue("id")) + if errors.Is(err, store.ErrNotFound) { + writeError(w, http.StatusNotFound, "not_found", "мапу не знайдено") + return + } + if err != nil { + // Невалідний uuid у шляху доходить сюди помилкою розбору — + // для клієнта це 400, а не 500. + if strings.Contains(err.Error(), "invalid input syntax for type uuid") { + writeError(w, http.StatusBadRequest, "bad_id", "некоректний ідентифікатор мапи") + return + } + s.log.Error("стан мапи", "err", err) + writeError(w, http.StatusInternalServerError, "internal", "внутрішня помилка") + return + } + writeJSON(w, http.StatusOK, state) +} + +func (s *Server) handleListDevices(w http.ResponseWriter, r *http.Request, tok *store.APIToken) { + if !tok.Can("devices:read") { + writeError(w, http.StatusForbidden, "forbidden", "немає права devices:read") + return + } + + devices, err := s.store.ListDevices(r.Context(), tok.TenantID) + if err != nil { + s.log.Error("перелік пристроїв", "err", err) + writeError(w, http.StatusInternalServerError, "internal", "внутрішня помилка") + return + } + if devices == nil { + devices = []store.DeviceSummary{} + } + writeJSON(w, http.StatusOK, map[string]any{"devices": devices}) +} + +func (s *Server) handleListAgents(w http.ResponseWriter, r *http.Request, tok *store.APIToken) { + if !tok.Can("agents:read") { + writeError(w, http.StatusForbidden, "forbidden", "немає права agents:read") + return + } + + agents, err := s.store.ListAgents(r.Context(), tok.TenantID) + if err != nil { + s.log.Error("перелік зондів", "err", err) + writeError(w, http.StatusInternalServerError, "internal", "внутрішня помилка") + return + } + if agents == nil { + agents = []store.AgentSummary{} + } + writeJSON(w, http.StatusOK, map[string]any{"agents": agents}) +} + +// --------------------------------------------------------------------- + +func writeJSON(w http.ResponseWriter, status int, v any) { + w.Header().Set("Content-Type", "application/json; charset=utf-8") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(v) +} + +func writeError(w http.ResponseWriter, status int, code, message string) { + writeJSON(w, status, map[string]any{ + "error": map[string]string{"code": code, "message": message}, + }) +} diff --git a/server/internal/httpapi/ws.go b/server/internal/httpapi/ws.go new file mode 100644 index 0000000..816f6bf --- /dev/null +++ b/server/internal/httpapi/ws.go @@ -0,0 +1,358 @@ +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 +} diff --git a/server/internal/store/apitokens.go b/server/internal/store/apitokens.go new file mode 100644 index 0000000..fc72056 --- /dev/null +++ b/server/internal/store/apitokens.go @@ -0,0 +1,87 @@ +package store + +import ( + "context" + "crypto/sha256" + "crypto/subtle" + "errors" + "time" + + "github.com/jackc/pgx/v5" +) + +var ErrTokenInvalid = errors.New("токен недійсний") + +// APIToken — ідентичність клієнта UI/інтеграції після автентифікації. +type APIToken struct { + ID string + TenantID string + Name string + Scopes []string +} + +// Can перевіряє право. Порожній набір scopes означає повний доступ: +// так поводяться токени, створені власником тенанта для себе. +func (t *APIToken) Can(scope string) bool { + if len(t.Scopes) == 0 { + return true + } + for _, s := range t.Scopes { + if s == scope || s == "*" { + return true + } + } + return false +} + +// AuthenticateAPIToken знаходить токен за значенням. +// +// У БД лежить лише sha256 — сам токен показується користувачу один раз +// при створенні й ніде більше не зберігається. Прострочені й відкликані +// відсіюються тут же, щоб жоден виклик далі не мусив про це пам'ятати. +func (s *Store) AuthenticateAPIToken(ctx context.Context, token string) (*APIToken, error) { + if token == "" { + return nil, ErrTokenInvalid + } + sum := sha256.Sum256([]byte(token)) + + var ( + t APIToken + hash []byte + expiresAt *time.Time + revokedAt *time.Time + ) + + err := s.pool.QueryRow(ctx, ` + SELECT id::text, tenant_id::text, name, scopes, token_hash, expires_at, revoked_at + FROM core.api_tokens + WHERE token_hash = $1 + `, sum[:]).Scan(&t.ID, &t.TenantID, &t.Name, &t.Scopes, &hash, &expiresAt, &revokedAt) + + if errors.Is(err, pgx.ErrNoRows) { + return nil, ErrTokenInvalid + } + if err != nil { + return nil, err + } + if subtle.ConstantTimeCompare(hash, sum[:]) != 1 { + return nil, ErrTokenInvalid + } + if revokedAt != nil { + return nil, ErrTokenInvalid + } + if expiresAt != nil && expiresAt.Before(time.Now()) { + return nil, ErrTokenInvalid + } + + // last_used_at оновлюємо без очікування: це діагностика, і платити + // за неї затримкою кожного запиту API немає сенсу. + go func() { + bg, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) + defer cancel() + _, _ = s.pool.Exec(bg, + `UPDATE core.api_tokens SET last_used_at = now() WHERE id = $1`, t.ID) + }() + + return &t, nil +} diff --git a/server/internal/store/events.go b/server/internal/store/events.go new file mode 100644 index 0000000..b9a81c5 --- /dev/null +++ b/server/internal/store/events.go @@ -0,0 +1,115 @@ +package store + +import ( + "context" + "encoding/json" + "time" + + "github.com/jackc/pgx/v5" +) + +// Подія для UI. Пишеться в core.event_outbox тією ж транзакцією, що й +// зміна, яку описує, — інакше WebSocket міг би розповісти про перехід, +// якого в базі ще (або вже) немає. +type Event struct { + ID int64 `json:"id"` + TenantID string `json:"-"` + Topic string `json:"topic"` + Payload json.RawMessage `json:"payload"` + At time.Time `json:"at"` +} + +// FetchEvents читає нові події після afterID. +// +// Транспорт навмисно простий: опитування таблиці замість LISTEN/NOTIFY. +// NOTIFY не переживає падіння підписника й обмежений 8 КБ на +// повідомлення, а тут потрібна гарантія, що жодна зміна статусу не +// загубиться між перезапусками API. Ціна — один дешевий запит за +// індексом раз на секунду. +func (s *Store) FetchEvents(ctx context.Context, afterID int64, limit int) ([]Event, error) { + if limit <= 0 || limit > 1000 { + limit = 500 + } + + rows, err := s.pool.Query(ctx, ` + SELECT id, tenant_id::text, topic, payload::text, created_at + FROM core.event_outbox + WHERE id > $1 + ORDER BY id + LIMIT $2 + `, afterID, limit) + if err != nil { + return nil, err + } + defer rows.Close() + + var out []Event + for rows.Next() { + var e Event + var payload string + if err := rows.Scan(&e.ID, &e.TenantID, &e.Topic, &payload, &e.At); err != nil { + return nil, err + } + e.Payload = json.RawMessage(payload) + out = append(out, e) + } + return out, rows.Err() +} + +// LatestEventID — з якого місця починати читати новому підписнику. +// +// Нова сесія WebSocket не має отримувати вчорашні події: клієнт щойно +// завантажив повний стан мапи, і все старіше в ньому вже враховано. +func (s *Store) LatestEventID(ctx context.Context) (int64, error) { + var id *int64 + if err := s.pool.QueryRow(ctx, + `SELECT max(id) FROM core.event_outbox`).Scan(&id); err != nil { + return 0, err + } + if id == nil { + return 0, nil + } + return *id, nil +} + +// MarkEventsPublished позначає події доставленими. +// +// Позначка не керує доставкою (підписники йдуть за id), вона потрібна +// прибиральнику: невідправлені події видаляти не можна, а відправлені — +// можна, і без цього поля таблиця росла б вічно. +func (s *Store) MarkEventsPublished(ctx context.Context, throughID int64) error { + _, err := s.pool.Exec(ctx, ` + UPDATE core.event_outbox + SET published_at = now() + WHERE id <= $1 AND published_at IS NULL + `, throughID) + return err +} + +// PruneEvents видаляє доставлені події, старші за вказаний вік. +func (s *Store) PruneEvents(ctx context.Context, olderThan time.Duration) (int64, error) { + tag, err := s.pool.Exec(ctx, ` + DELETE FROM core.event_outbox + WHERE published_at IS NOT NULL AND created_at < now() - $1::interval + `, olderThan.String()) + if err != nil { + return 0, err + } + return tag.RowsAffected(), nil +} + +// PublishEvent кладе подію в outbox поза чужою транзакцією. +// Для змін, що вже мають власну транзакцію, подія пишеться прямо в ній. +func (s *Store) PublishEvent(ctx context.Context, tenantID, topic string, payload any) error { + data, err := json.Marshal(payload) + if err != nil { + return err + } + return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { + _, err := tx.Exec(ctx, ` + INSERT INTO core.event_outbox (tenant_id, topic, payload) + VALUES ($1, $2, $3::jsonb) + `, tenantID, topic, string(data)) + return err + }) +} diff --git a/server/internal/store/inventory.go b/server/internal/store/inventory.go new file mode 100644 index 0000000..0278fb0 --- /dev/null +++ b/server/internal/store/inventory.go @@ -0,0 +1,130 @@ +package store + +import ( + "context" + "time" + + "github.com/jackc/pgx/v5" +) + +type DeviceSummary struct { + ID string `json:"id"` + Name string `json:"name"` + Address string `json:"address,omitempty"` + Kind string `json:"kind"` + Vendor string `json:"vendor,omitempty"` + Model string `json:"model,omitempty"` + SiteName string `json:"site_name,omitempty"` + Status string `json:"status"` + Enabled bool `json:"enabled"` + LastSeenAt *time.Time `json:"last_seen_at,omitempty"` + AgentID string `json:"agent_id,omitempty"` + IfaceCount int `json:"interface_count"` +} + +// ListDevices — інвентар для таблиці в UI. +func (s *Store) ListDevices(ctx context.Context, tenantID string) ([]DeviceSummary, error) { + var out []DeviceSummary + + err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { + rows, err := tx.Query(ctx, ` + SELECT d.id::text, d.name, COALESCE(host(d.address), ''), d.kind::text, + COALESCE(d.vendor,''), COALESCE(d.model,''), COALESCE(st.name,''), + d.status::text, d.enabled, d.last_seen_at, COALESCE(d.agent_id::text,''), + (SELECT count(*) FROM inv.interfaces i WHERE i.device_id = d.id) + FROM inv.devices d + LEFT JOIN inv.sites st ON st.id = d.site_id + WHERE d.tenant_id = $1 AND d.deleted_at IS NULL + ORDER BY d.name + `, tenantID) + if err != nil { + return err + } + defer rows.Close() + + for rows.Next() { + var d DeviceSummary + if err := rows.Scan(&d.ID, &d.Name, &d.Address, &d.Kind, &d.Vendor, &d.Model, + &d.SiteName, &d.Status, &d.Enabled, &d.LastSeenAt, &d.AgentID, + &d.IfaceCount); err != nil { + return err + } + out = append(out, d) + } + return rows.Err() + }) + + return out, err +} + +type AgentSummary struct { + ID string `json:"id"` + Name string `json:"name"` + Status string `json:"status"` + Version string `json:"version,omitempty"` + OS string `json:"os,omitempty"` + Arch string `json:"arch,omitempty"` + Hostname string `json:"hostname,omitempty"` + LastHeartbeatAt *time.Time `json:"last_heartbeat_at,omitempty"` + EnabledModules []string `json:"enabled_modules"` + Health map[string]any `json:"health"` + DeviceCount int `json:"device_count"` +} + +// ListAgents — стан зондів для UI. +// +// health віддається як є: там лежить зведення з останнього heartbeat, +// зокрема dropped_samples. Оператор має бачити діру в даних без +// походу в графіки самометрик. +func (s *Store) ListAgents(ctx context.Context, tenantID string) ([]AgentSummary, error) { + var out []AgentSummary + + err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { + rows, err := tx.Query(ctx, ` + SELECT a.id::text, a.name, a.status::text, + COALESCE(a.version,''), COALESCE(a.os,''), COALESCE(a.arch,''), + COALESCE(a.hostname,''), a.last_heartbeat_at, a.enabled_modules, a.health, + (SELECT count(*) FROM inv.devices d + WHERE d.agent_id = a.id AND d.deleted_at IS NULL) + FROM core.agents a + WHERE a.tenant_id = $1 + ORDER BY a.name + `, tenantID) + if err != nil { + return err + } + defer rows.Close() + + for rows.Next() { + var a AgentSummary + if err := rows.Scan(&a.ID, &a.Name, &a.Status, &a.Version, &a.OS, &a.Arch, + &a.Hostname, &a.LastHeartbeatAt, &a.EnabledModules, &a.Health, + &a.DeviceCount); err != nil { + return err + } + if a.EnabledModules == nil { + a.EnabledModules = []string{} + } + if a.Health == nil { + a.Health = map[string]any{} + } + out = append(out, a) + } + return rows.Err() + }) + + return out, err +} + +// MapTenant повертає тенанта мапи — потрібно WebSocket-у, щоб не дати +// підписатись на чужу мапу за вгаданим id. +func (s *Store) MapExists(ctx context.Context, tenantID, mapID string) (bool, error) { + var exists bool + err := s.pool.QueryRow(ctx, ` + SELECT EXISTS ( + SELECT 1 FROM topo.maps + WHERE id = $1 AND tenant_id = $2 AND deleted_at IS NULL + ) + `, mapID, tenantID).Scan(&exists) + return exists, err +} diff --git a/server/internal/store/maps.go b/server/internal/store/maps.go new file mode 100644 index 0000000..12b1a38 --- /dev/null +++ b/server/internal/store/maps.go @@ -0,0 +1,364 @@ +package store + +import ( + "context" + "encoding/json" + "errors" + "time" + + "github.com/jackc/pgx/v5" +) + +// Тут живлять мапу: усе, що потрібно полотну, віддається одним викликом. +// +// Структури мають json-теги, бо їдуть у фронтенд як є. Це свідомо: +// проміжний DTO-шар між БД і React додав би роботи й нічого не дав — +// форма полотна й так визначена схемою (topo.map_nodes / map_edges). + +type MapSummary struct { + ID string `json:"id"` + Name string `json:"name"` + Slug string `json:"slug"` + Kind string `json:"kind"` + SiteID string `json:"site_id,omitempty"` + IsDefault bool `json:"is_default"` + NodeCount int `json:"node_count"` + EdgeCount int `json:"edge_count"` + Revision int64 `json:"revision"` + UpdatedAt time.Time `json:"updated_at"` +} + +type MapState struct { + ID string `json:"id"` + Name string `json:"name"` + Slug string `json:"slug"` + Kind string `json:"kind"` + LayoutAlgo string `json:"layout_algo"` + Viewport json.RawMessage `json:"viewport"` + Grid json.RawMessage `json:"grid"` + Clustering json.RawMessage `json:"clustering"` + Revision int64 `json:"revision"` + + Backgrounds []MapBackground `json:"backgrounds"` + Nodes []MapNode `json:"nodes"` + Edges []MapEdge `json:"edges"` +} + +type MapBackground struct { + ID string `json:"id"` + Kind string `json:"kind"` + StorageKey string `json:"storage_key,omitempty"` + MimeType string `json:"mime_type,omitempty"` + X float64 `json:"x"` + Y float64 `json:"y"` + Width *float64 `json:"width,omitempty"` + Height *float64 `json:"height,omitempty"` + Rotation float64 `json:"rotation"` + Opacity float64 `json:"opacity"` + Locked bool `json:"locked"` + ZIndex int `json:"z_index"` + Geo json.RawMessage `json:"geo,omitempty"` + Rack json.RawMessage `json:"rack,omitempty"` +} + +type MapNode struct { + ID string `json:"id"` + Kind string `json:"kind"` + Label string `json:"label,omitempty"` + DeviceID string `json:"device_id,omitempty"` + ParentID string `json:"parent_id,omitempty"` + X float64 `json:"x"` + Y float64 `json:"y"` + Width *float64 `json:"width,omitempty"` + Height *float64 `json:"height,omitempty"` + ZIndex int `json:"z_index"` + Style json.RawMessage `json:"style"` + Data json.RawMessage `json:"data"` + Collapsed bool `json:"collapsed"` + Locked bool `json:"locked"` + + // Живий стан пристрою. Саме це фарбує вузол. + Status string `json:"status,omitempty"` + LastSeenAt *time.Time `json:"last_seen_at,omitempty"` + RttMs *float32 `json:"rtt_ms,omitempty"` + LossPct *float32 `json:"loss_pct,omitempty"` +} + +type MapEdge struct { + ID string `json:"id"` + SourceNodeID string `json:"source_node_id"` + TargetNodeID string `json:"target_node_id"` + SourcePort string `json:"source_port,omitempty"` + TargetPort string `json:"target_port,omitempty"` + LinkID string `json:"link_id,omitempty"` + Label string `json:"label,omitempty"` + Style string `json:"style"` + Dash string `json:"dash"` + Color string `json:"color,omitempty"` + WidthPx float64 `json:"width_px"` + Waypoints json.RawMessage `json:"waypoints"` + Animation json.RawMessage `json:"animation"` + Thresholds json.RawMessage `json:"thresholds"` + ShowMetrics bool `json:"show_metrics"` + + // Живе завантаження — джерело швидкості анімації. + LinkStatus string `json:"link_status,omitempty"` + UtilPct *float64 `json:"util_pct,omitempty"` + CapacityBps *int64 `json:"capacity_bps,omitempty"` +} + +var ErrNotFound = errors.New("не знайдено") + +// ListMaps віддає перелік мап тенанта з лічильниками. +func (s *Store) ListMaps(ctx context.Context, tenantID string) ([]MapSummary, error) { + var out []MapSummary + + err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { + rows, err := tx.Query(ctx, ` + SELECT m.id::text, m.name, m.slug, m.kind::text, + COALESCE(m.site_id::text, ''), m.is_default, m.revision, m.updated_at, + (SELECT count(*) FROM topo.map_nodes n WHERE n.map_id = m.id), + (SELECT count(*) FROM topo.map_edges e WHERE e.map_id = m.id) + FROM topo.maps m + WHERE m.tenant_id = $1 AND m.deleted_at IS NULL + ORDER BY m.is_default DESC, m.name + `, tenantID) + if err != nil { + return err + } + defer rows.Close() + + for rows.Next() { + var m MapSummary + if err := rows.Scan(&m.ID, &m.Name, &m.Slug, &m.Kind, &m.SiteID, + &m.IsDefault, &m.Revision, &m.UpdatedAt, &m.NodeCount, &m.EdgeCount); err != nil { + return err + } + out = append(out, m) + } + return rows.Err() + }) + + return out, err +} + +// GetMapState збирає повний стан полотна разом із живими статусами. +// +// Чотири запити в одній транзакції замість N+1 на кожен вузол: фронтенд +// отримує готове до рендеру полотно, а не сотні дозапитів. Живий стан +// підмішується тут же — інакше мапа малювалася б сірою й лише потім +// доганяла кольори. +func (s *Store) GetMapState(ctx context.Context, tenantID, mapID string) (*MapState, error) { + var st *MapState + + err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { + m := &MapState{} + err := tx.QueryRow(ctx, ` + SELECT id::text, name, slug, kind::text, layout_algo::text, + viewport::text, grid::text, clustering::text, revision + FROM topo.maps + WHERE tenant_id = $1 AND id = $2 AND deleted_at IS NULL + `, tenantID, mapID).Scan(&m.ID, &m.Name, &m.Slug, &m.Kind, &m.LayoutAlgo, + &m.Viewport, &m.Grid, &m.Clustering, &m.Revision) + if errors.Is(err, pgx.ErrNoRows) { + return ErrNotFound + } + if err != nil { + return err + } + + if m.Backgrounds, err = mapBackgrounds(ctx, tx, mapID); err != nil { + return err + } + if m.Nodes, err = mapNodes(ctx, tx, tenantID, mapID); err != nil { + return err + } + if m.Edges, err = mapEdges(ctx, tx, tenantID, mapID); err != nil { + return err + } + + st = m + return nil + }) + + return st, err +} + +func mapBackgrounds(ctx context.Context, tx pgx.Tx, mapID string) ([]MapBackground, error) { + rows, err := tx.Query(ctx, ` + SELECT id::text, kind::text, COALESCE(storage_key,''), COALESCE(mime_type,''), + x, y, width, height, rotation, opacity, locked, z_index, + COALESCE(geo::text,''), COALESCE(rack::text,'') + FROM topo.map_backgrounds + WHERE map_id = $1 + ORDER BY z_index + `, mapID) + if err != nil { + return nil, err + } + defer rows.Close() + + out := []MapBackground{} + for rows.Next() { + var b MapBackground + var geo, rack string + if err := rows.Scan(&b.ID, &b.Kind, &b.StorageKey, &b.MimeType, + &b.X, &b.Y, &b.Width, &b.Height, &b.Rotation, &b.Opacity, + &b.Locked, &b.ZIndex, &geo, &rack); err != nil { + return nil, err + } + b.Geo = rawOrNil(geo) + b.Rack = rawOrNil(rack) + out = append(out, b) + } + return out, rows.Err() +} + +func mapNodes(ctx context.Context, tx pgx.Tx, tenantID, mapID string) ([]MapNode, error) { + // ts.device_last_icmp вже обмежений останніми 15 хвилинами: старіші + // дані не характеризують поточний стан, і показувати їх на мапі + // означало б брехати про живість пристрою. + rows, err := tx.Query(ctx, ` + SELECT n.id::text, n.kind::text, COALESCE(n.label,''), + COALESCE(n.device_id::text,''), COALESCE(n.parent_node_id::text,''), + n.x, n.y, n.width, n.height, n.z_index, + n.style::text, n.data::text, n.collapsed, n.locked, + COALESCE(d.status::text,''), d.last_seen_at, + i.rtt_avg_ms, i.loss_pct + FROM topo.map_nodes n + LEFT JOIN inv.devices d ON d.id = n.device_id AND d.tenant_id = $1 + LEFT JOIN ts.device_last_icmp i ON i.device_id = n.device_id + WHERE n.map_id = $2 AND n.tenant_id = $1 AND NOT n.hidden + ORDER BY n.z_index, n.id + `, tenantID, mapID) + if err != nil { + return nil, err + } + defer rows.Close() + + out := []MapNode{} + for rows.Next() { + var n MapNode + if err := rows.Scan(&n.ID, &n.Kind, &n.Label, &n.DeviceID, &n.ParentID, + &n.X, &n.Y, &n.Width, &n.Height, &n.ZIndex, + &n.Style, &n.Data, &n.Collapsed, &n.Locked, + &n.Status, &n.LastSeenAt, &n.RttMs, &n.LossPct); err != nil { + return nil, err + } + out = append(out, n) + } + return out, rows.Err() +} + +// linkStatusExpr виводить стан лінка з його кінців. +// +// Колонка topo.links.status існує, але її ніхто не підтримує: писати +// туди означало б оновлювати всі лінки пристрою на кожну зміну його +// статусу й тримати це узгодженим. Дешевше й чесніше порахувати на +// читанні — стан лінка є похідною величиною, а не фактом. +// +// Порядок гілок важливий: обрив (down) перекриває все інше, бо саме він +// вимагає уваги оператора; «невідомо» стоїть перед «up», щоб мапа не +// малювала зеленим те, чого ще жодного разу не опитували. +const linkStatusExpr = ` + CASE + WHEN da.status = 'down' OR db.status = 'down' THEN 'down' + WHEN COALESCE(ia.oper_status::text,'up') NOT IN ('up','unknown') + OR COALESCE(ib.oper_status::text,'up') NOT IN ('up','unknown') THEN 'down' + WHEN da.status = 'maintenance' OR db.status = 'maintenance' THEN 'maintenance' + WHEN da.status = 'unknown' OR db.status = 'unknown' THEN 'unknown' + WHEN da.status = 'warning' OR db.status = 'warning' THEN 'warning' + ELSE 'up' + END` + +func mapEdges(ctx context.Context, tx pgx.Tx, tenantID, mapID string) ([]MapEdge, error) { + rows, err := tx.Query(ctx, ` + SELECT e.id::text, e.source_node_id::text, e.target_node_id::text, + COALESCE(si.name,''), COALESCE(ti.name,''), + COALESCE(e.link_id::text,''), COALESCE(e.label,''), + e.style::text, e.dash::text, COALESCE(e.color,''), e.width_px, + e.waypoints::text, e.animation::text, e.thresholds::text, e.show_metrics, + CASE WHEN l.id IS NULL THEN '' ELSE `+linkStatusExpr+` END, + lv.util_pct, l.capacity_bps + FROM topo.map_edges e + LEFT JOIN inv.interfaces si ON si.id = e.source_interface_id + LEFT JOIN inv.interfaces ti ON ti.id = e.target_interface_id + LEFT JOIN topo.links l ON l.id = e.link_id + LEFT JOIN inv.devices da ON da.id = l.a_device_id + LEFT JOIN inv.devices db ON db.id = l.b_device_id + LEFT JOIN inv.interfaces ia ON ia.id = l.a_interface_id + LEFT JOIN inv.interfaces ib ON ib.id = l.b_interface_id + LEFT JOIN topo.link_live lv ON lv.link_id = e.link_id + WHERE e.map_id = $2 AND e.tenant_id = $1 AND NOT e.hidden + ORDER BY e.z_index, e.id + `, tenantID, mapID) + if err != nil { + return nil, err + } + defer rows.Close() + + out := []MapEdge{} + for rows.Next() { + var e MapEdge + if err := rows.Scan(&e.ID, &e.SourceNodeID, &e.TargetNodeID, + &e.SourcePort, &e.TargetPort, &e.LinkID, &e.Label, + &e.Style, &e.Dash, &e.Color, &e.WidthPx, + &e.Waypoints, &e.Animation, &e.Thresholds, &e.ShowMetrics, + &e.LinkStatus, &e.UtilPct, &e.CapacityBps); err != nil { + return nil, err + } + out = append(out, e) + } + return out, rows.Err() +} + +// LinkLoad — поточне завантаження лінка для WebSocket-оновлень. +type LinkLoad struct { + LinkID string `json:"link_id"` + Status string `json:"status"` + UtilPct *float64 `json:"util_pct,omitempty"` +} + +// MapLinkLoads віддає завантаження всіх лінків мапи. +// +// Окремо від повного стану: WebSocket штовхає лише ці кілька чисел +// кожні кілька секунд, а не перемальовує полотно. +func (s *Store) MapLinkLoads(ctx context.Context, tenantID, mapID string) ([]LinkLoad, error) { + var out []LinkLoad + + err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { + rows, err := tx.Query(ctx, ` + SELECT DISTINCT l.id::text, `+linkStatusExpr+`, lv.util_pct + FROM topo.map_edges e + JOIN topo.links l ON l.id = e.link_id + LEFT JOIN inv.devices da ON da.id = l.a_device_id + LEFT JOIN inv.devices db ON db.id = l.b_device_id + LEFT JOIN inv.interfaces ia ON ia.id = l.a_interface_id + LEFT JOIN inv.interfaces ib ON ib.id = l.b_interface_id + LEFT JOIN topo.link_live lv ON lv.link_id = l.id + WHERE e.map_id = $2 AND e.tenant_id = $1 + `, tenantID, mapID) + if err != nil { + return err + } + defer rows.Close() + + for rows.Next() { + var l LinkLoad + if err := rows.Scan(&l.LinkID, &l.Status, &l.UtilPct); err != nil { + return err + } + out = append(out, l) + } + return rows.Err() + }) + + return out, err +} + +func rawOrNil(s string) json.RawMessage { + if s == "" || s == "null" { + return nil + } + return json.RawMessage(s) +} diff --git a/server/internal/store/telemetry.go b/server/internal/store/telemetry.go index 62297c3..03c0109 100644 --- a/server/internal/store/telemetry.go +++ b/server/internal/store/telemetry.go @@ -225,6 +225,9 @@ func (s *Store) applyDeviceStatus(ctx context.Context, tenantID, deviceID, statu } return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { + // Історія й подія для UI пишуться в тій самій транзакції, що й + // сама зміна: інакше мапа могла б показати новий колір без + // запису в історії або навпаки. _, err := tx.Exec(ctx, ` WITH prev AS ( SELECT id, status @@ -241,12 +244,25 @@ func (s *Store) applyDeviceStatus(ctx context.Context, tenantID, deviceID, statu FROM prev WHERE d.id = prev.id RETURNING prev.status AS old_status, d.status AS new_status + ), changed AS ( + SELECT old_status, new_status FROM upd + WHERE old_status IS DISTINCT FROM new_status + ), hist AS ( + INSERT INTO ts.device_status_history + (ts, device_id, tenant_id, status, prev_status, reason) + SELECT $4, $1, $2, new_status, old_status, $5 FROM changed + ON CONFLICT (ts, device_id) DO NOTHING + RETURNING 1 ) - INSERT INTO ts.device_status_history (ts, device_id, tenant_id, status, prev_status, reason) - SELECT $4, $1, $2, new_status, old_status, $5 - FROM upd - WHERE old_status IS DISTINCT FROM new_status - ON CONFLICT (ts, device_id) DO NOTHING + INSERT INTO core.event_outbox (tenant_id, topic, payload) + SELECT $2, 'device.status', jsonb_build_object( + 'device_id', $1::text, + 'status', new_status::text, + 'previous_status', old_status::text, + 'reason', $5::text, + 'ts', $4::timestamptz + ) + FROM changed `, deviceID, tenantID, status, ts, nullString(reason)) return err })