// Команда 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/icmp" "github.com/netpulse/netpulse/agent/internal/modules/snmp" "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, } // Реєстрація перед усім іншим. // // Посвідчення зберігається на диск одразу: одноразовий токен згорає // на сервері, і другої спроби не буде — впасти після успішного // 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, err := config.LoadIdentity(cfg.IdentityPath) if err != nil { return err } 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 } reg.EnsureDefaults(cfg.DefaultModules...) 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, }) 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) }() 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 }