// Команда netpulse-server — приймальна сторона AgentService. // // Зонди самі підключаються сюди; сервер до них не ходить. package main import ( "context" "crypto/tls" "crypto/x509" "errors" "flag" "fmt" "log/slog" "net" "net/http" "os" "os/signal" "syscall" "time" npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" "github.com/netpulse/netpulse/server/internal/alerting" "github.com/netpulse/netpulse/server/internal/crypto" "github.com/netpulse/netpulse/server/internal/gitstore" "github.com/netpulse/netpulse/server/internal/grpcapi" "github.com/netpulse/netpulse/server/internal/logbuf" "github.com/netpulse/netpulse/server/internal/store" "google.golang.org/grpc" "google.golang.org/grpc/credentials" "google.golang.org/grpc/keepalive" ) var version = "dev" func main() { if err := run(); err != nil { fmt.Fprintln(os.Stderr, "netpulse-server:", err) os.Exit(1) } } func run() error { var ( listen = flag.String("listen", envOr("NETPULSE_LISTEN", ":9443"), "адреса прослуховування gRPC") 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_CERT"), "сертифікат сервера") keyFile = flag.String("key", os.Getenv("NETPULSE_KEY"), "приватний ключ") caFile = flag.String("client-ca", os.Getenv("NETPULSE_CLIENT_CA"), "CA для перевірки сертифікатів зондів (mTLS)") insecure = flag.Bool("insecure", os.Getenv("NETPULSE_INSECURE") == "1", "без TLS — лише локальний стенд") keysFlag = flag.String("dek", os.Getenv("NETPULSE_DEK"), "ключі шифрування: key_id=[,...]") logLevel = flag.String("log-level", envOr("NETPULSE_LOG_LEVEL", "info"), "debug|info|warn|error") gitRoot = flag.String("git-root", envOr("NETPULSE_GIT_ROOT", "/var/lib/netpulse/git"), "корінь сховища версій конфігів; порожньо — без Git") controlEndpoint = flag.String("control-endpoint", os.Getenv("NETPULSE_CONTROL_ENDPOINT"), "адреса, за якою зонди мають підключатись; порожньо — та, якою вони прийшли") // Внутрішня ручка журналу. Порт сусідній із gRPC і назовні не // публікується — так само, як 9443: до нього ходить лише // REST-процес усередині мережі, і лише з токеном. internalListen = flag.String("internal-listen", envOr("NETPULSE_INTERNAL_LISTEN", ":9444"), "адреса внутрішньої ручки журналу для REST-процесу; порожньо — не слухати") ) flag.Parse() if *dsn == "" { return errors.New("не вказано -dsn (або NETPULSE_DSN)") } // Кільце останніх рядків власного журналу — джерело для сторінки // «Журнал сервера». Саме в цьому процесі крутяться такти, які й // лаються вночі: диспетчер збору конфігів, розклад бекапів, обидва // прибиральники, закриття періодів SLA й дзеркалення на Git. Без // цієї половини сторінка показувала б тишу там, де щогодини падає // push у Forgejo. // // Секрети реєструються всі, які цей процес знає: кільце робить // журнал видимим у браузері й вивантажуваним у файл, і другий рубіж // маскування тут не зайвий — див. logbuf.Ring.Mask. logRing := logbuf.New(logbuf.DefaultMaxBytes) logRing.Mask(logbuf.DSNSecrets(*dsn)...) logRing.Mask(logbuf.DSNSecrets(*dsnWorker)...) logRing.Mask(logbuf.KeySecrets(*keysFlag)...) log := newLogger(*logLevel, logRing) ring, err := buildKeyring(*keysFlag) if err != nil { return err } 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() // Друге з'єднання — роллю netpulse_worker, для запитів, які за // побудовою ходять поверх усіх кабінетів (див. коментар до Store.bg // і міграцію 0063). Порожня змінна лишає все як було: фонові запити // підуть основним пулом. Саме тому цю зміну можна викотити наперед, // а роль перемкнути окремим кроком. if err := st.UseWorkerDSN(ctx, *dsnWorker); err != nil { return fmt.Errorf("підключення воркера: %w", err) } // Сховище версій конфігів. Спільний каталог для API й колектора: // колектор туди пише під час бекапу, API звідти читає для diff. // Різні машини потребують спільного тому — інакше половина історії // буде видна одному процесу, а половина другому. if *gitRoot != "" { st.UseGit(gitstore.New(*gitRoot)) } // Подієві алерти на журналі й конфігах. // // Живуть у цьому процесі, бо саме сюди приходять і те, і те. Сам // приймач нічого нікуди не шле — він лише піднімає алерт із // позначкою «розіслати»; розсилає netpulse-api, де є ключі каналів, // маршрути й тихі години. svc := grpcapi.New(st, ring, log). WithEventAlerts(alerting.NewEventSink(st, log)) opts := []grpc.ServerOption{ grpc.ChainUnaryInterceptor(svc.UnaryInterceptor), grpc.ChainStreamInterceptor(svc.StreamInterceptor), // Зонди сидять за NAT: без keepalive проміжний маршрутизатор // тихо викидає сесію, і сервер вважає мертвого агента живим. grpc.KeepaliveParams(keepalive.ServerParameters{ Time: 60 * time.Second, Timeout: 20 * time.Second, }), grpc.KeepaliveEnforcementPolicy(keepalive.EnforcementPolicy{ MinTime: 20 * time.Second, PermitWithoutStream: true, }), // Конфіги великих маршрутизаторів бувають на кілька МБ, // але їх шлють чанками — ліміт лишається захистом. grpc.MaxRecvMsgSize(16 << 20), } if !*insecure { tc, err := serverTLS(*certFile, *keyFile, *caFile) if err != nil { return err } opts = append(opts, grpc.Creds(credentials.NewTLS(tc))) } else { log.Warn("запуск без TLS — припустимо лише для локального стенду") } srv := grpc.NewServer(opts...) npv1.RegisterAgentServiceServer(srv, svc) // Реєстрація йде БЕЗ токена зонда, тож окремим сервісом: інакше // довелося б робити виняток усередині інтерсептора автентифікації. // Адреса для зондів — не та, на якій ми слухаємо. // // Раніше сюди йшов -listen, і зонд зберігав у посвідченні ":9443" // або "0.0.0.0:9443" — адресу, за якою до сервера не достукатись // нізвідки. Помітно це стає лише після перезапуску зонда: реєстрація // проходить, а підключитись він більше не може ніколи. // // Порожнє значення означає «лишись на тій адресі, якою прийшов» — і // це правильна поведінка за замовчуванням: зонд дійшов, отже, вона // робоча. Прапорець потрібен лише там, де реєстрацію й роботу // свідомо розводять на різні адреси. npv1.RegisterEnrollmentServiceServer(srv, grpcapi.NewEnrollService(st, *controlEndpoint)) lis, err := net.Listen("tcp", *listen) if err != nil { return fmt.Errorf("listen %s: %w", *listen, err) } log.Info("запуск", "version", version, "listen", *listen, "tls", !*insecure) // Внутрішня ручка журналу — окремим слухачем і НЕ фатально. // // Зайнятий порт, відсутній -dek, будь-яка інша причина — усе це має // коштувати рівно одного попередження й порожньої половини на // сторінці. Збір даних із мережі не повинен не піднятись через // зручність перегляду власного журналу; той самий принцип, що й у // перевірці ліцензії в REST-процесі. serveInternalLog(ctx, logRing, *internalListen, logbuf.InternalToken(*keysFlag), log) // Диспетчер збору конфігів: черга наповнюється REST-процесом, а // живі сесії зондів тримає саме цей. go svc.DispatchConfigJobs(ctx, 5*time.Second) // Розклад бекапів. Кілька екземплярів безпечні: тік бере // advisory-блокування, тож розкручує його рівно один. go svc.ScheduleBackups(ctx) // Звірка планів: чеки міняє REST-процес, а перезалити план може // лише той, хто тримає сесію зонда. go svc.SyncPlans(ctx) // Прибиральник старих версій конфігів. Тут, а не в REST-процесі: // поруч із тим, хто версії створює, і подалі від шляху запитів // людини — див. ncm_retention.go. go svc.SweepRetention(ctx) // Прибиральник телеметрії та журналів за строками зберігання, він // же — спостерігач за розміром бази. Тут із тих самих міркувань, що // й попередній: поруч із тим, хто ці дані створює, і подалі від // шляху запитів людини — див. storage_retention.go. go svc.SweepDataRetention(ctx) // Закриття періодів SLA. Тут із тієї ж причини, що й обидва // прибиральники, і ще з однієї: закривати період треба ВЧАСНО, поки // дані під ним ще є, а не тоді, коли хтось відкриє сторінку — див. // sla_close.go. go svc.CloseSLAPeriods(ctx) // Дзеркалення архіву конфігів на зовнішній Git. Окремий такт, а не // push після коміту: недоступний Forgejo не має коштувати жодного // бекапу — див. ncm_mirror.go. go svc.MirrorGit(ctx) errCh := make(chan error, 1) go func() { errCh <- srv.Serve(lis) }() select { case <-ctx.Done(): log.Info("зупинка", "сесій_онлайн", svc.SessionCount()) // GracefulStop дає активним стрімам догратись: обірваний // посеред батчу агент однаково перешле його, але зайвих // ретраїв краще уникнути. done := make(chan struct{}) go func() { srv.GracefulStop(); close(done) }() select { case <-done: case <-time.After(15 * time.Second): log.Warn("м'яка зупинка не встигла — примусова") srv.Stop() } return nil case err := <-errCh: return err } } 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)) } // serveInternalLog піднімає внутрішню ручку журналу. // // Мовчазна відмова тут неприпустима, тому кожна причина «не піднялось» // їде в журнал окремим рядком: сторінка покаже «джерело недоступне», і // пояснення до цього має бути хоч десь. А от зупиняти процес жодна з // них не має права. func serveInternalLog(ctx context.Context, ring *logbuf.Ring, addr, token string, log *slog.Logger) { if addr == "" { log.Info("внутрішня ручка журналу вимкнена (-internal-listen порожній)") return } if token == "" { log.Warn("внутрішня ручка журналу не піднята: немає -dek, " + "а токен виводиться саме з нього") return } lis, err := net.Listen("tcp", addr) if err != nil { log.Warn("внутрішня ручка журналу не піднята", "addr", addr, "err", err) return } mux := http.NewServeMux() mux.Handle(logbuf.InternalPath, ring.HTTPHandler("collector", token)) srv := &http.Server{ Handler: mux, ReadHeaderTimeout: 5 * time.Second, } go func() { if err := srv.Serve(lis); err != nil && !errors.Is(err, http.ErrServerClosed) { log.Warn("внутрішня ручка журналу зупинилась", "err", err) } }() go func() { <-ctx.Done() _ = srv.Close() }() log.Info("внутрішня ручка журналу", "addr", addr, "path", logbuf.InternalPath) } // buildKeyring розбирає ключі шифрування. // // Ключі приходять через оточення або секрет-менеджер і ніколи не // лежать поруч із дампом БД: інакше шифрування конфігів і паролів // від обладнання не давало б нічого. На відміну від REST-процесу, // тут ключ обов'язковий: без нього AgentService не може ані видати // креденшели, ані зберегти конфіг. func buildKeyring(spec string) (*crypto.Keyring, error) { ring, err := crypto.ParseKeyring(spec) if err != nil { return nil, err } if ring == nil { return nil, errors.New("не вказано -dek: без ключа неможливо ані видати креденшели, ані зберегти конфіг") } return ring, nil } func serverTLS(certFile, keyFile, clientCA string) (*tls.Config, error) { if certFile == "" || keyFile == "" { return nil, errors.New("потрібні -cert і -key (або -insecure для локального стенду)") } cert, err := tls.LoadX509KeyPair(certFile, keyFile) if err != nil { return nil, fmt.Errorf("сертифікат сервера: %w", err) } tc := &tls.Config{ Certificates: []tls.Certificate{cert}, MinVersion: tls.VersionTLS13, } if clientCA != "" { pem, err := os.ReadFile(clientCA) if err != nil { return nil, fmt.Errorf("CA зондів: %w", err) } pool := x509.NewCertPool() if !pool.AppendCertsFromPEM(pem) { return nil, errors.New("CA зондів: не вдалося розібрати PEM") } tc.ClientCAs = pool // Токен каже, ЯКИЙ це зонд; сертифікат — що він узагалі має // право говорити з сервером. Обидва обов'язкові. tc.ClientAuth = tls.RequireAndVerifyClientCert } return tc, nil }