Netpulse_SasS/agent/cmd/netpulse-agent/main.go
byrsapty 32b0b003df
Some checks are pending
CI / web (push) Waiting to run
CI / server (push) Waiting to run
CI / agent (push) Waiting to run
Дрібний борг: модуль http/ssl і зрозумілі помилки на мапі
Плагін http стояв у сіді як базовий — тобто обіцяний усім одразу, — а
модуля не існувало. Тепер є: код відповіді, час, збіг слова в тілі й
залишок днів до кінця сертифіката.

Недоступність повертається нулем, а не помилкою чека: на графіку це
читається як провал, і саме за цим ставлять тригер. Довіру до ланцюга
сертифікатів свідомо не перевіряємо — питають строк, а самопідписаний
теж має дату.

Повторна лінія між вузлами й повторно доданий хост давали однакове
«такий запис уже існує». Тепер кожен випадок каже, що робити далі.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-25 18:37:58 +03:00

330 lines
11 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-agent — легкий зонд збору телеметрії.
//
// Єдиний статичний бінарник. Усі з'єднання вихідні: у мережі клієнта
// не потрібно відкривати жодного порту.
package main
import (
"context"
"crypto/tls"
"errors"
"fmt"
"log/slog"
"os"
"os/signal"
"runtime"
"runtime/debug"
"sync"
"syscall"
"time"
"github.com/netpulse/netpulse/agent/internal/config"
"github.com/netpulse/netpulse/agent/internal/module"
"github.com/netpulse/netpulse/agent/internal/modules/httpx"
"github.com/netpulse/netpulse/agent/internal/modules/icmp"
"github.com/netpulse/netpulse/agent/internal/modules/snmp"
"github.com/netpulse/netpulse/agent/internal/modules/syslog"
"github.com/netpulse/netpulse/agent/internal/modules/topology"
"github.com/netpulse/netpulse/agent/internal/scheduler"
"github.com/netpulse/netpulse/agent/internal/session"
"github.com/netpulse/netpulse/agent/internal/telemetry"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/keepalive"
)
// version підставляється при збірці: -ldflags "-X main.version=1.2.3"
var (
version = "dev"
commit = "none"
)
func main() {
if err := run(); err != nil {
fmt.Fprintln(os.Stderr, "netpulse-agent:", err)
os.Exit(1)
}
}
func run() error {
cfg, err := config.Parse(os.Args[1:])
if err != nil {
return err
}
// Для реєстрації досить версії й платформи: перелік скомпільованих
// модулів збирається нижче, разом із реєстром, і чекати на нього
// заради Enroll немає сенсу.
enrollBuild := &npv1.AgentBuild{
Version: version,
Commit: commit,
Os: runtime.GOOS,
Arch: runtime.GOARCH,
}
// Наявне посвідчення важить більше за запрошення.
//
// Запрошення одноразове, а живе воно в змінних оточення чи в unit-
// файлі — тобто лишається на місці назавжди. Кожен перезапуск зонда
// намагався б зареєструватись повторно й падав із «запрошення
// недійсне», хоча посвідчення поруч і цілком робоче.
//
// Знайдено живим розгортанням: контейнер із NETPULSE_ENROLL у .env
// піднявся один раз, а після першого ж restart уже не вставав.
saved, err := config.LoadIdentity(cfg.IdentityPath)
if err != nil {
return err
}
if saved != nil && cfg.EnrollToken != "" {
fmt.Fprintf(os.Stderr,
"netpulse-agent: посвідчення вже є (%s), запрошення проігноровано\n",
cfg.IdentityPath)
cfg.EnrollToken = ""
}
// Реєстрація перед усім іншим.
//
// Посвідчення зберігається на диск одразу: одноразовий токен згорає
// на сервері, і другої спроби не буде — впасти після успішного
// Enroll означало б залишити людину без агента й без запрошення.
if cfg.EnrollToken != "" {
id, err := enroll(cfg, enrollBuild)
if err != nil {
return err
}
cfg.AgentID, cfg.Token, cfg.Endpoint = id.AgentID, id.Token, id.Endpoint
fmt.Fprintf(os.Stderr, "netpulse-agent: зареєстровано як %q, посвідчення в %s\n",
id.AgentName, cfg.IdentityPath)
} else if cfg.AgentID == "" || cfg.Token == "" {
id := saved
if id == nil {
return fmt.Errorf(
"зонд не зареєстрований: посвідчення не знайдено (%s). "+
"Створіть запрошення в UI і запустіть із -enroll <токен>",
cfg.IdentityPath)
}
cfg.AgentID, cfg.Token = id.AgentID, id.Token
if cfg.Endpoint == "" {
cfg.Endpoint = id.Endpoint
}
}
if cfg.AgentID == "" || cfg.Token == "" {
return errors.New("зонд не зареєстрований: потрібні -enroll або -agent-id з -token")
}
log := newLogger(cfg)
// Бюджет пам'яті — вимога, а не побажання: зонд часто живе на
// роутері або в контейнері зі 64 МБ. GOMEMLIMIT змушує збирач
// працювати агресивніше замість того, щоб дати OOM killer'у
// вбити процес і осліпити моніторинг саме тоді, коли він потрібен.
debug.SetMemoryLimit(48 << 20)
reg := module.NewRegistry()
if err := reg.Register(icmp.New()); err != nil {
return err
}
if err := reg.Register(snmp.New()); err != nil {
return err
}
if err := reg.Register(topology.New()); err != nil {
return err
}
if err := reg.Register(httpx.New()); err != nil {
return err
}
reg.EnsureDefaults(cfg.DefaultModules...)
// Приймач syslog не модуль реєстру: у нього немає задач і розкладу,
// він просто слухає порт. Вмикається тим самим переліком -modules,
// щоб людині не треба було знати про цю різницю.
var syslogRecv *syslog.Receiver
for _, m := range cfg.DefaultModules {
if m == "syslog" {
syslogRecv = syslog.New(cfg.SyslogListen, log)
break
}
}
buf := telemetry.NewBuffer(telemetry.Options{
MaxItems: cfg.BufferMaxItems,
MaxBytes: cfg.BufferMaxBytes,
})
build := &npv1.AgentBuild{
Version: version,
Commit: commit,
Os: runtime.GOOS,
Arch: runtime.GOARCH,
GoVersion: runtime.Version(),
CompiledModules: reg.Compiled(),
}
dial, err := dialer(cfg)
if err != nil {
return err
}
sess := session.New(session.Config{
AgentID: cfg.AgentID,
Hostname: cfg.Hostname,
Build: build,
Registry: reg,
Buffer: buf,
Dial: dial,
Logger: log,
MinBackoff: cfg.MinBackoff,
MaxBackoff: cfg.MaxBackoff,
DefaultModules: cfg.DefaultModules,
Syslog: syslogRecv,
})
sched := scheduler.New(scheduler.Config{
Registry: reg,
Sink: buf,
OnStatus: sess.ReportStatus,
Credentials: sess.Credentials,
OnDiscovery: sess.ReportDiscovery,
MaxConcurrency: cfg.MaxConcurrency,
})
sess.SetScheduler(sched)
ctx, stop := signal.NotifyContext(context.Background(),
os.Interrupt, syscall.SIGTERM)
defer stop()
log.Info("запуск",
"version", version,
"server", cfg.Endpoint,
"agent_id", cfg.AgentID,
"modules", reg.Compiled())
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
sched.Run(ctx)
}()
// Приймач живе поза сесією: обрив зв'язку з сервером не має
// зупиняти збір журналу. Події лягають у чергу й доїдуть, коли
// сесія відновиться — саме заради цього черга й існує.
if syslogRecv != nil {
wg.Add(1)
go func() {
defer wg.Done()
if err := syslogRecv.Run(ctx); err != nil {
// Зайнятий порт або брак прав — не привід зупиняти
// зонд: решта модулів працює, а причину видно в журналі.
log.Error("приймач syslog не запустився",
"адреса", cfg.SyslogListen, "err", err)
}
}()
}
runErr := sess.Run(ctx)
wg.Wait()
for _, e := range reg.CloseAll() {
log.Warn("помилка при закритті модуля", "err", e)
}
if runErr != nil && !errors.Is(runErr, context.Canceled) {
return runErr
}
log.Info("зупинено")
return nil
}
func newLogger(cfg *config.Config) *slog.Logger {
level := slog.LevelInfo
switch cfg.LogLevel {
case "debug":
level = slog.LevelDebug
case "warn":
level = slog.LevelWarn
case "error":
level = slog.LevelError
}
opts := &slog.HandlerOptions{Level: level}
if cfg.LogJSON {
return slog.New(slog.NewJSONHandler(os.Stderr, opts))
}
return slog.New(slog.NewTextHandler(os.Stderr, opts))
}
// tokenCreds додає токен зонда в метадані кожного виклику.
type tokenCreds struct {
token string
allowInsecure bool
}
func (t tokenCreds) GetRequestMetadata(ctx context.Context, uri ...string) (map[string]string, error) {
return map[string]string{"authorization": "Bearer " + t.token}, nil
}
func (t tokenCreds) RequireTransportSecurity() bool { return !t.allowInsecure }
func dialer(cfg *config.Config) (session.DialFunc, error) {
tc, err := cfg.TLSConfig()
if err != nil {
return nil, err
}
creds := insecure.NewCredentials()
if tc != nil {
creds = credentials.NewTLS(tc)
}
opts := []grpc.DialOption{
grpc.WithTransportCredentials(creds),
// Токен їде в метаданих кожного виклику. gRPC відмовиться
// слати його по незашифрованому каналу, якщо ми явно не
// дозволимо це для локального стенду.
grpc.WithPerRPCCredentials(tokenCreds{token: cfg.Token, allowInsecure: cfg.Insecure}),
// Keepalive потрібен через NAT: без нього проміжний
// маршрутизатор тихо викидає сесію після кількох хвилин
// мовчання, і сервер бачить агента живим, коли той уже глухий.
grpc.WithKeepaliveParams(keepalive.ClientParameters{
Time: 30 * time.Second,
Timeout: 10 * time.Second,
PermitWithoutStream: true,
}),
}
return func(ctx context.Context) (session.Conn, error) {
return grpc.NewClient(cfg.Endpoint, opts...)
}, nil
}
// enroll реєструє зонд і зберігає посвідчення.
func enroll(cfg *config.Config, build *npv1.AgentBuild) (*config.Identity, error) {
var tlsCfg *tls.Config
if !cfg.Insecure {
var err error
tlsCfg, err = cfg.TLSConfig()
if err != nil {
return nil, err
}
}
id, err := config.Enroll(context.Background(), cfg.Endpoint,
cfg.EnrollToken, cfg.EnrollName, tlsCfg, build)
if err != nil {
return nil, err
}
if err := config.SaveIdentity(cfg.IdentityPath, id); err != nil {
// Токен уже згорів на сервері, тож повідомлення має нести сам
// токен: інакше людині доведеться створювати нове запрошення
// лише через те, що каталог виявився недоступним для запису.
return nil, fmt.Errorf(
"зонд зареєстровано (id=%s), але посвідчення не збереглося в %s: %w\n"+
"збережіть вручну: {\"agent_id\":%q,\"token\":%q}",
id.AgentID, cfg.IdentityPath, err, id.AgentID, id.Token)
}
return id, nil
}