Netpulse_SasS/server/cmd/netpulse-api/main.go
byrsapty e6b585dd4c
All checks were successful
CI / hygiene (push) Successful in 11s
CI / web (push) Successful in 1m11s
CI / server (push) Successful in 1m55s
CI / agent (push) Successful in 1m0s
Сторінка «Журнал сервера»: кільцевий буфер у памяті
Продукт показував syslog З ПРИСТРОЇВ, аудит, черги, стенограми команд —
усе про мережу, і нічого про себе. Щоб дізнатись, на що лається сам
NetPulse, треба було заходити по ssh.

Буфер у памяті, не в базі: журнал, що пише в Postgres, замовкає рівно
тоді, коли ляже Postgres — у найцікавіший момент. Стеля в БАЙТАХ
(8 МіБ на процес), а не в рядках: один рядок із текстом SQL-помилки
буває довшим за сотню звичайних, тож «5000 рядків» означало б
непередбачувані десятки мегабайтів.

Два процеси — один перелік із позначкою джерела, а не дві вкладки:
людина знає симптом («о третій ночі перестали йти сповіщення»), а не
те, який із двох процесів за це відповідає.

Маскування — наявним gitstore.Scrub, не своїм: паролі DSN, матеріал
DEK, секрет підпису сесій. Це другий рубіж — відомі шляхи вже почищені
в місці народження, але кільце робить журнал видимим у браузері й
вивантажуваним у файл, що піде в тікет.

Межі написані НА СТОРІНЦІ, а не лише в документації: журнал не
переживає перезапуску й не покаже причини падіння бази. Поруч —
скільки записів витіснено: без цього числа не відрізнити «нічого не
сталося» від «сталося стільки, що початок уже не влазить».
2026-08-27 23:50:03 +03:00

368 lines
19 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.

// Команда 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/alerting"
"github.com/netpulse/netpulse/server/internal/auth"
"github.com/netpulse/netpulse/server/internal/crypto"
"github.com/netpulse/netpulse/server/internal/gitstore"
"github.com/netpulse/netpulse/server/internal/httpapi"
"github.com/netpulse/netpulse/server/internal/logbuf"
"github.com/netpulse/netpulse/server/internal/store"
"github.com/netpulse/netpulse/server/webui"
)
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")
dsnWorker = flag.String("dsn-worker", os.Getenv("NETPULSE_DSN_WORKER"),
"DSN для фонових тактів поверх усіх кабінетів (порожньо — тим самим з'єднанням)")
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")
jwtKey = flag.String("jwt-secret", os.Getenv("NETPULSE_JWT_SECRET"),
"секрет підпису токенів доступу, щонайменше 32 байти")
pruneAge = flag.Duration("event-retention", 24*time.Hour,
"скільки тримати доставлені події в core.event_outbox")
keysFlag = flag.String("dek", os.Getenv("NETPULSE_DEK"),
"ключі шифрування: key_id=<hex|base64>[,...] — потрібні для каналів сповіщень")
alertEvery = flag.Duration("alert-interval", 30*time.Second,
"як часто обчислювати правила алертів; 0 — не запускати движок")
alertKeep = flag.Duration("alert-retention", 7*24*time.Hour,
"скільки тримати закриті алерти до переносу в історію")
gitRoot = flag.String("git-root", envOr("NETPULSE_GIT_ROOT", "/var/lib/netpulse/git"),
"корінь сховища версій конфігів; порожньо — без Git")
licensePub = flag.String("license-pubkey", os.Getenv("NETPULSE_LICENSE_PUBKEY"),
"відкриті ключі перевірки ліцензій: kid=<hex|base64>[,...]; порожньо — ліцензії не перевіряються")
licenseEvery = flag.Duration("license-interval", time.Hour,
"як часто перевіряти стан ліцензії; 0 — лише при старті")
privateHooks = flag.Bool("allow-private-webhooks",
os.Getenv("NETPULSE_ALLOW_PRIVATE_WEBHOOKS") == "1",
"дозволити вебхуки на внутрішні адреси — для self-hosted інсталяцій")
// Увімкнено за замовчуванням: кнопки під сповіщеннями малюються
// завжди, і інсталяція, де вони є, а приймача немає, — це рівно
// той стан, який цей приймач і виправляє. Прапорець лишається
// для мереж, з яких немає виходу на api.telegram.org: там
// опитування давало б лише потік помилок у журналі.
telegramBot = flag.Bool("telegram-callbacks",
os.Getenv("NETPULSE_TELEGRAM_CALLBACKS") != "0",
"приймати натискання кнопок під сповіщеннями Telegram (довге опитування)")
// Типове значення — ім'я служби з docker-compose, тобто рівно та
// адреса, за якою колектор стоїть у кожній інсталяції, зробленій
// установником. Порожнє значення тут коштувало б дорожче:
// сторінка журналу з коробки показувала б половину системи, і
// побачити другу половину змогли б лише ті, хто прочитав про
// цей прапорець. Там, де топологія інша, ім'я не розв'язується
// за мілісекунди, і сторінка чесно каже, що саме налаштувати.
//
// Ціна цієї зручності: у мережі, де DNS відповідає на БУДЬ-ЯКЕ
// ім'я (пошуковий домен, wildcard), запит із похідним від -dek
// токеном піде чужому хосту. Сам ключ із токена не дістати
// (HMAC), але тим, кого це стосується, лікується порожнім
// NETPULSE_COLLECTOR_LOG_URL.
collectorLog = flag.String("collector-log",
envOr("NETPULSE_COLLECTOR_LOG_URL", "http://collector:9444"),
"адреса внутрішньої ручки журналу колектора; порожньо — не питати")
)
flag.Parse()
if *dsn == "" {
return errors.New("не вказано -dsn (або NETPULSE_DSN)")
}
// Кільце останніх рядків власного журналу — щоб сторінку «Журнал
// сервера» можна було відкрити в браузері замість ssh. Створюється
// ДО логера й до всього іншого: рядки про невдалий старт (немає
// -jwt-secret, не розібрано ключі) цікаві найбільше, а їх пише
// найперший код.
//
// Секрети реєструються тут же й усі, які цей процес узагалі знає.
// Ідея не в тому, що вони точно потраплять у журнал, — відомі шляхи
// вже почищено в джерелі. Ідея в тому, що кільце робить журнал
// видимим у браузері й вивантажуваним у файл, і покладатись на
// повноту переліку місць, де хтось не забув почистити, тут не можна.
logRing := logbuf.New(logbuf.DefaultMaxBytes)
logRing.Mask(logbuf.DSNSecrets(*dsn)...)
logRing.Mask(logbuf.DSNSecrets(*dsnWorker)...)
logRing.Mask(logbuf.KeySecrets(*keysFlag)...)
logRing.Mask(*jwtKey)
log := newLogger(*logLevel, logRing)
// Секрет обовʼязковий і не генерується автоматично: випадковий
// ключ при кожному старті означав би, що будь-який перезапуск
// викидає всіх користувачів, а кілька екземплярів API за
// балансувальником не приймали б токени один одного.
signer, err := auth.NewSigner([]byte(*jwtKey))
if err != nil {
return fmt.Errorf("-jwt-secret: %w", err)
}
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
st, err2 := store.New(ctx, *dsn)
err = err2
if err != nil {
return fmt.Errorf("підключення до БД: %w", err)
}
defer st.Close()
// Друге з'єднання — роллю netpulse_worker, для запитів, які за
// побудовою ходять поверх усіх кабінетів (див. коментар до Store.bg
// і міграцію 0063). Порожня змінна лишає все як було: фонові запити
// підуть основним пулом. Саме тому цю зміну можна викотити наперед,
// а роль перемкнути окремим кроком.
if err := st.UseWorkerDSN(ctx, *dsnWorker); err != nil {
return fmt.Errorf("підключення воркера: %w", err)
}
ring, err := crypto.ParseKeyring(*keysFlag)
if err != nil {
return fmt.Errorf("-dek: %w", err)
}
// Сховище версій конфігів. Спільний каталог для API й колектора:
// колектор туди пише під час бекапу, API звідти читає для diff.
// Різні машини потребують спільного тому — інакше половина історії
// буде видна одному процесу, а половина другому.
if *gitRoot != "" {
st.UseGit(gitstore.New(*gitRoot))
}
api := httpapi.New(st, signer, log)
// Сторінка «Журнал сервера». Токен внутрішньої ручки колектора
// виводиться зі спеки ключів шифрування — єдиного секрету, який в
// обох процесів уже спільний; чому саме так — у logbuf/http.go.
// Без -dek токен порожній, і сторінка так і скаже.
api = api.WithServerLog(logRing, *collectorLog, logbuf.InternalToken(*keysFlag))
// Ключі шифрування потрібні не лише каналам сповіщень: секрет
// дзеркала конфігів лежить у тому самому core.secrets, а движок
// алертів на інсталяції може бути вимкнений.
api = api.WithKeyring(ring)
// Зібраний інтерфейс, якщо він є в цій збірці. Порожній dist —
// робочий стан: розробка йде проти vite, а API просто віддає API.
if fsys, ok := webui.FS(); ok {
httpapi.SetStatic(fsys)
log.Info("веб-інтерфейс вшито в бінарник")
}
// Ліцензія. Порожній ключ — робочий стан: збірка без вшитого
// відкритого ключа працює як працювала, ліміти беруться з тарифу, а
// сторінка тарифу чесно каже, що перевірити ключ нічим.
verifier, err := store.ParseLicenseVerifier(*licensePub)
if err != nil {
return fmt.Errorf("-license-pubkey: %w", err)
}
api = api.WithLicense(verifier)
go api.Hub().Run(ctx)
go pruneLoop(ctx, st, log, *pruneAge)
go licenseLoop(ctx, st, verifier, log, *licenseEvery)
// Движок алертів живе тут, а не в netpulse-server, бо саме цей
// процес уже читає БД для UI і має ключі для каналів доставки.
// Одночасний запуск кількох екземплярів безпечний: тік бере
// advisory-блокування, тож обчислює завжди рівно один.
if *alertEvery > 0 {
eng := alerting.New(st, ring, log, *alertEvery, *privateHooks)
api = api.WithNotifications(ring, eng.Notifier())
// Прогін відповідності запускають з UI, тобто з цього процесу —
// і саме він перетворює знахідку на алерт. Без цього тригер
// «порушено вимогу» лишався б тим, чим був: рядком у базі, який
// ніколи не спрацює.
api = api.WithEventAlerts(alerting.NewEventSink(st, log))
go eng.Run(ctx)
go eng.RunHousekeeping(ctx, *alertKeep)
// Приймач натискань кнопок під сповіщеннями Telegram.
//
// Тут же, де й доставка: кнопки малює notify.go, і розводити
// «надіслати» й «прийняти натиснуте» по різних процесах
// означало б інсталяцію, де кнопки є, а відповіді на них немає.
//
// Довге опитування, а не вебхук — розгортання за самопідписаним
// TLS на IP-адресі вебхука не приймає в принципі. Повне
// обґрунтування — у telegram_bot.go, поруч із самим кодом.
// Кілька екземплярів API безпечні: приймач тримає власне
// advisory-блокування, тож getUpdates робить рівно один.
if *telegramBot {
go alerting.NewBot(st, ring, log).Run(ctx)
}
}
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)
}
}
}
}
// ЧОМУ ПЕРЕВІРКА ЛІЦЕНЗІЇ — ЦИКЛ, А НЕ ОДИН КРОК ПРИ СТАРТІ
//
// Бо стан ліцензії міняється без жодної дії — просто від того, що минув
// час. Інсталяцію, яку не перезавантажували чотири місяці (а це норма
// для коробки в клієнта), перевірка при старті лишила б в active аж до
// наступного оновлення: попередження «лишилось п'ять діб» людина
// побачила б через півроку після того, як вони скінчились.
//
// ЧОМУ ПОМИЛКА ТУТ НЕ ЗУПИНЯЄ ПРОЦЕС
//
// Це головне рішення, і воно те саме, що й у решті ліцензійної частини.
// Моніторинг, який не піднявся через ліцензію, — аварія в мережі
// клієнта, спричинена нами: він не побачить падіння магістралі й
// дізнається про нього від абонентів. Тому будь-яка невдача перевірки —
// це рядок у журналі, а не код виходу. Найгірше, що з неї виходить, —
// стелі лишаються такими, якими були.
// licenseLoop перевіряє ліцензію при старті й далі за тактом.
func licenseLoop(ctx context.Context, st *store.Store, v *store.LicenseVerifier,
log *slog.Logger, every time.Duration) {
check := func() {
state, err := st.RefreshLicense(ctx, v)
if err != nil {
log.Warn("перевірка ліцензії", "err", err)
return
}
// Рівень залежить від стану, і це не косметика: рядок про
// пільговий період має бути помітним у потоці журналу, бо це
// єдине попередження, яке отримає той, хто на сторінку тарифу
// не заходить.
switch state.State {
case store.LicenseGrace:
log.Warn("ліцензія: пільговий період", "діб", state.DaysLeft,
"до", state.GraceUntil)
case store.LicenseExpired:
log.Warn("ліцензія протермінована — збір і сповіщення працюють, "+
"нові хости й зонди не заводяться", "кому", state.IssuedTo)
case store.LicenseInvalid:
log.Error("ліцензія не перевіряється", "причина", state.Reason)
default:
log.Info("ліцензія", "стан", state.State, "інсталяція", state.InstallID)
}
}
check()
if every <= 0 {
return
}
t := time.NewTicker(every)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
check()
}
}
}
func envOr(key, def string) string {
if v := os.Getenv(key); v != "" {
return v
}
return def
}
// newLogger збирає логер процесу.
//
// Кільце ДОПОВНЮЄ вивід, а не заміняє його: JSON у stderr лишається
// таким самим, як був, і його так само забирає docker. Це не дрібниця —
// stderr переживає падіння процесу, а кільце ні, тож заміна одного
// другим позбавила б інсталяцію єдиного джерела про причини аварії.
func newLogger(level string, ring *logbuf.Ring) *slog.Logger {
lv := slog.LevelInfo
switch level {
case "debug":
lv = slog.LevelDebug
case "warn":
lv = slog.LevelWarn
case "error":
lv = slog.LevelError
}
h := slog.NewJSONHandler(os.Stderr, &slog.HandlerOptions{Level: lv})
return slog.New(ring.Wrap(h))
}