Один коміт, а не десяток тематичних, свідомо: теми переплетені в
спільних файлах (store.go, docker-compose.yml, deploy/README.md), і
розділити їх можна було б лише індексуванням шматків. Коміти, які не
збираються, гірші за один великий — тим паче що це рівно той стан, який
перевірявся разом.
ЩО ПРАЦЮЄ НА СТЕНДІ Й ПЕРЕВІРЕНО ТАМ
0058 подієві алерти: syslog, ncm, compliance спрацьовують у мить
події; правило з нереалізованим джерелом більше не зберігається
мовчки
0059 snmp.walk і прототипи шаблонів — таблиці з динамічним індексом
описуються шаблоном, а не Go
0060 відкат конфігу: план як різниця, маскування паролів із підписом
плану, обов'язковий контрольний збір, verifying при обриві
0061 кнопки Telegram: довге опитування, авторизація не з callback_data
0062 аудит і архів хостів; тест на AST, що падає на ключі без назви
0063 RLS: три ролі, окремий пул для фонових тактів
0064 строки зберігання даних і сторінка сховища
0065 приймач SNMP-трапів; перевірено справжніми пакетами по дроту,
переклад v1→v2 за RFC 3584 дає правильний OID
0066 ескалації сповіщень
0067 алерт про вичерпання диска
0068 поля заливки конфігу переїхали в каталог профілів
Плюс: 137 тестів вебу з нуля (їх не було взагалі), одинадцять справжніх
вад, знайдених ними й виправлених, і виправлення двох інтеграційних
тестів grpcapi, які мовчки пропускались півтора року.
ЩО ЩЕ НЕ ЗАПУСКАЛОСЬ
netpulse установник: одна команда замість 18 змінних і
593 рядків інструкції
RLS з першого запуску нова інсталяція під політиками одразу;
RLS-EXISTING-INSTALL.md лишається тільки для
старих інсталяцій
.forgejo + CI раннер не зареєстрований
Ці три перевірені компіляцією й міркуванням, але не виконанням.
ГОЛОВНИЙ ВИСНОВОК ДВОХ СЕСІЙ
Зелена перевірка доводить рівно те, що вона перевіряє. Тест ізоляції RLS
був правильний і зелений — і пропустив зламаний вхід, бо перевіряв «чи
не видно чужого», коли зламалось «чи видно своє». Інтеграційні тести
grpcapi були зелені, бо не виконувались. Схема, довідник і протокол
описували те, чого в коді не існувало, і виглядало це як готове.
Тому в кожному завданні цих сесій стояла вимога назвати НЕПОКРИТЕ, а
чотири задачі закінчились не можливістю, а відмовою: правило з
нереалізованим джерелом не зберігається, профіль без команд заливки
каже про це замість мовчазної кнопки, міграція RLS валить сама себе на
таблиці без політики, тест словника аудиту падає на ключі без назви.
Подробиці — HISTORY.md, розділи за 26 і 27 серпня.
864 lines
28 KiB
Go
864 lines
28 KiB
Go
// Package session тримає з'єднання з сервером і все, що з нього випливає:
|
||
// рукостискання, застосування плану задач, відправку телеметрії,
|
||
// підтвердження, реконект.
|
||
//
|
||
// Усе з'єднання ініціює агент — у мережі клієнта немає ані відкритих
|
||
// портів, ані прокидання NAT. Тому «команда з сервера» технічно є
|
||
// повідомленням у зустрічному напрямку вже відкритого стріму Control.
|
||
package session
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"fmt"
|
||
"log/slog"
|
||
"math/rand"
|
||
"runtime"
|
||
"sync"
|
||
"sync/atomic"
|
||
"time"
|
||
|
||
"github.com/netpulse/netpulse/agent/internal/module"
|
||
"github.com/netpulse/netpulse/agent/internal/modules/filecfg"
|
||
"github.com/netpulse/netpulse/agent/internal/modules/syslog"
|
||
"github.com/netpulse/netpulse/agent/internal/modules/traps"
|
||
"github.com/netpulse/netpulse/agent/internal/scheduler"
|
||
"github.com/netpulse/netpulse/agent/internal/telemetry"
|
||
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
|
||
"google.golang.org/grpc"
|
||
"google.golang.org/protobuf/types/known/durationpb"
|
||
"google.golang.org/protobuf/types/known/timestamppb"
|
||
)
|
||
|
||
// Conn — те, що вміє *grpc.ClientConn. Винесено в інтерфейс, щоб тести
|
||
// підставляли bufconn без справжньої мережі.
|
||
type Conn interface {
|
||
grpc.ClientConnInterface
|
||
Close() error
|
||
}
|
||
|
||
// DialFunc відкриває з'єднання до сервера.
|
||
type DialFunc func(ctx context.Context) (Conn, error)
|
||
|
||
type Config struct {
|
||
AgentID string
|
||
Hostname string
|
||
Build *npv1.AgentBuild
|
||
Registry *module.Registry
|
||
Buffer *telemetry.Buffer
|
||
Scheduler *scheduler.Scheduler
|
||
Dial DialFunc
|
||
Logger *slog.Logger
|
||
|
||
MinBackoff time.Duration
|
||
MaxBackoff time.Duration
|
||
|
||
// Модулі, увімкнені до першої відповіді сервера.
|
||
DefaultModules []string
|
||
|
||
// Приймач syslog. Порожній — зонд журнали не збирає, і стрім до
|
||
// сервера не відкривається взагалі: тримати порожній канал заради
|
||
// вимкненої можливості немає сенсу.
|
||
Syslog *syslog.Receiver
|
||
|
||
// Приймач SNMP-трапів. Порожній — зонд трапи не приймає. Окремо від
|
||
// Syslog, бо це окремий порт і окремий дозвіл у фаєрволі клієнта:
|
||
// вмикати обидва там, де просили лише один, означало б відкрити
|
||
// порт, про який ніхто не домовлявся.
|
||
Traps *traps.Receiver
|
||
|
||
// Збір конфіг-файлів власної машини. Порожній — зонд такі завдання
|
||
// відхиляє з поясненням.
|
||
FileCfg *filecfg.Module
|
||
}
|
||
|
||
type Session struct {
|
||
cfg Config
|
||
log *slog.Logger
|
||
|
||
// Статуси задач — діагностика. Канал з буфером і скиданням при
|
||
// переповненні: краще втратити рядок журналу, ніж загальмувати
|
||
// опитування через повільний контрольний канал.
|
||
statusCh chan *npv1.TaskStatusUpdate
|
||
|
||
// Звіти автовиявлення. Окремо від телеметрії: вони рідкі, великі
|
||
// й їдуть унарним ReportDiscovery, а не стрімом.
|
||
discoCh chan *npv1.DiscoveryReport
|
||
|
||
credMu sync.RWMutex
|
||
creds map[string][]*npv1.Credential
|
||
credExpiry time.Time
|
||
// Коли востаннє просили поновлення — щоб не питати щотіку.
|
||
credAsked time.Time
|
||
|
||
devMu sync.RWMutex
|
||
devices map[string]*npv1.DeviceTarget
|
||
|
||
// client живий лише в межах сесії. Збір конфігу триває хвилинами й
|
||
// не має тримати контрольний цикл, тому вивантаження йде окремою
|
||
// горутиною — а їй потрібен доступ до клієнта поточної сесії.
|
||
client atomic.Pointer[npv1.AgentServiceClient]
|
||
// jobs рахує незавершені збори: при обриві сесії їх треба дочекатись,
|
||
// інакше вивантаження піде в уже закритий канал.
|
||
jobs sync.WaitGroup
|
||
|
||
planHash atomic.Pointer[[]byte]
|
||
nextBatch atomic.Uint64
|
||
lastAcked atomic.Uint64
|
||
skewNanos atomic.Int64
|
||
connected atomic.Bool
|
||
seq atomic.Uint64
|
||
|
||
startedAt time.Time
|
||
}
|
||
|
||
func New(cfg Config) *Session {
|
||
if cfg.Logger == nil {
|
||
cfg.Logger = slog.Default()
|
||
}
|
||
if cfg.MinBackoff <= 0 {
|
||
cfg.MinBackoff = time.Second
|
||
}
|
||
if cfg.MaxBackoff <= 0 {
|
||
cfg.MaxBackoff = 2 * time.Minute
|
||
}
|
||
s := &Session{
|
||
cfg: cfg,
|
||
log: cfg.Logger,
|
||
statusCh: make(chan *npv1.TaskStatusUpdate, 512),
|
||
discoCh: make(chan *npv1.DiscoveryReport, 16),
|
||
creds: make(map[string][]*npv1.Credential),
|
||
devices: make(map[string]*npv1.DeviceTarget),
|
||
startedAt: time.Now(),
|
||
}
|
||
// Приймач слухає порт незалежно від сесії, але зіставити адресу з
|
||
// хостом може лише вона: список хостів приходить у плані.
|
||
if cfg.Syslog != nil {
|
||
cfg.Syslog.SetResolver(s.resolveDeviceByIP)
|
||
}
|
||
if cfg.Traps != nil {
|
||
cfg.Traps.SetResolver(s.resolveDeviceByIP)
|
||
}
|
||
return s
|
||
}
|
||
|
||
// SetScheduler замикає взаємну залежність: планувальник потребує
|
||
// колбеків сесії (статуси, креденшели), а сесія — планувальника, щоб
|
||
// застосовувати план. Конструктором це не виражається, тому окремий сетер.
|
||
// Викликати до Run.
|
||
func (s *Session) SetScheduler(sc *scheduler.Scheduler) { s.cfg.Scheduler = sc }
|
||
|
||
// Credentials — провайдер для планувальника.
|
||
func (s *Session) Credentials(deviceID string) []*npv1.Credential {
|
||
s.credMu.RLock()
|
||
defer s.credMu.RUnlock()
|
||
|
||
// Прострочені креденшели не віддаємо: краще явна помилка
|
||
// "немає креденшелів", ніж заблокований обліковий запис на
|
||
// половині комутаторів клієнта.
|
||
if !s.credExpiry.IsZero() && time.Now().After(s.credExpiry) {
|
||
return nil
|
||
}
|
||
return s.creds[deviceID]
|
||
}
|
||
|
||
// ReportDiscovery — колбек для планувальника.
|
||
//
|
||
// Черга навмисно коротка: звіт автовиявлення актуальний рівно доти,
|
||
// доки описує поточний стан мережі. Накопичувати десяток застарілих
|
||
// знімків, поки немає зв'язку, немає сенсу — наступний обхід віддасть
|
||
// свіжіші дані.
|
||
func (s *Session) ReportDiscovery(res module.Result) {
|
||
if len(res.Neighbors) == 0 && len(res.InterfaceRecords) == 0 && len(res.Devices) == 0 {
|
||
return
|
||
}
|
||
rep := &npv1.DiscoveryReport{
|
||
AgentId: s.cfg.AgentID,
|
||
Neighbors: res.Neighbors,
|
||
Interfaces: res.InterfaceRecords,
|
||
Devices: res.Devices,
|
||
Final: true,
|
||
StartedAt: timestamppb.Now(),
|
||
FinishedAt: timestamppb.Now(),
|
||
}
|
||
|
||
select {
|
||
case s.discoCh <- rep:
|
||
default:
|
||
// Витісняємо найстаріший, а не відкидаємо новий.
|
||
select {
|
||
case <-s.discoCh:
|
||
default:
|
||
}
|
||
select {
|
||
case s.discoCh <- rep:
|
||
default:
|
||
}
|
||
}
|
||
}
|
||
|
||
// ReportStatus — колбек для планувальника.
|
||
func (s *Session) ReportStatus(u *npv1.TaskStatusUpdate) {
|
||
select {
|
||
case s.statusCh <- u:
|
||
default:
|
||
}
|
||
}
|
||
|
||
// Connected — чи є жива сесія (для /healthz і самометрик).
|
||
func (s *Session) Connected() bool { return s.connected.Load() }
|
||
|
||
// ClockSkew — поправка годинника відносно сервера.
|
||
func (s *Session) ClockSkew() time.Duration {
|
||
return time.Duration(s.skewNanos.Load())
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Цикл підключення
|
||
// ---------------------------------------------------------------------
|
||
|
||
// Run тримає з'єднання, доки не скасують контекст.
|
||
func (s *Session) Run(ctx context.Context) error {
|
||
backoff := s.cfg.MinBackoff
|
||
|
||
for {
|
||
if ctx.Err() != nil {
|
||
return ctx.Err()
|
||
}
|
||
|
||
err := s.runOnce(ctx)
|
||
s.connected.Store(false)
|
||
|
||
if ctx.Err() != nil {
|
||
return ctx.Err()
|
||
}
|
||
if err != nil {
|
||
s.log.Warn("сесію розірвано", "err", err, "retry_in", backoff)
|
||
}
|
||
|
||
// Джитер — щоб сотня зондів не ломанулась перепідключатись
|
||
// одночасно після рестарту сервера.
|
||
jitter := time.Duration(rand.Int63n(int64(backoff/2 + 1)))
|
||
select {
|
||
case <-ctx.Done():
|
||
return ctx.Err()
|
||
case <-time.After(backoff + jitter):
|
||
}
|
||
|
||
backoff *= 2
|
||
if backoff > s.cfg.MaxBackoff {
|
||
backoff = s.cfg.MaxBackoff
|
||
}
|
||
}
|
||
}
|
||
|
||
func (s *Session) runOnce(ctx context.Context) error {
|
||
conn, err := s.cfg.Dial(ctx)
|
||
if err != nil {
|
||
return fmt.Errorf("dial: %w", err)
|
||
}
|
||
defer conn.Close()
|
||
|
||
client := npv1.NewAgentServiceClient(conn)
|
||
s.client.Store(&client)
|
||
defer s.client.Store(nil)
|
||
|
||
sctx, cancel := context.WithCancel(ctx)
|
||
defer cancel()
|
||
|
||
ctrl, err := client.Control(sctx)
|
||
if err != nil {
|
||
return fmt.Errorf("відкрити Control: %w", err)
|
||
}
|
||
|
||
if err := ctrl.Send(s.hello()); err != nil {
|
||
return fmt.Errorf("Hello: %w", err)
|
||
}
|
||
|
||
first, err := ctrl.Recv()
|
||
if err != nil {
|
||
return fmt.Errorf("очікування Welcome: %w", err)
|
||
}
|
||
welcome := first.GetWelcome()
|
||
if welcome == nil {
|
||
return fmt.Errorf("замість Welcome прийшло %T", first.Payload)
|
||
}
|
||
s.applyWelcome(welcome)
|
||
// Таблиця series_ref на сервері живе рівно стільки, скільки сесія.
|
||
// Перереєстровуємо всі відомі серії, замість обнуляти нумерацію:
|
||
// інакше довелося б викинути буфер, накопичений за час обриву, —
|
||
// саме ті дані, заради яких він і накопичувався.
|
||
s.cfg.Buffer.Interner().MarkAllPending()
|
||
s.connected.Store(true)
|
||
s.log.Info("сесію встановлено",
|
||
"session_id", welcome.SessionId,
|
||
"clock_skew", s.ClockSkew())
|
||
|
||
// Єдиний писар у контрольний стрім: gRPC не допускає паралельних
|
||
// Send, а слати треба і heartbeat, і статуси задач.
|
||
out := make(chan *npv1.ControlUp, 256)
|
||
|
||
var wg sync.WaitGroup
|
||
// Місткість дорівнює числу горутин: інакше та, що впала останньою,
|
||
// заблокується на записі й ніколи не дочекається wg.Wait().
|
||
errCh := make(chan error, 8)
|
||
|
||
spawn := func(name string, fn func() error) {
|
||
wg.Add(1)
|
||
go func() {
|
||
defer wg.Done()
|
||
if err := fn(); err != nil && !errors.Is(err, context.Canceled) {
|
||
errCh <- fmt.Errorf("%s: %w", name, err)
|
||
}
|
||
cancel()
|
||
}()
|
||
}
|
||
|
||
spawn("writer", func() error { return s.writerLoop(sctx, ctrl, out) })
|
||
spawn("heartbeat", func() error { return s.heartbeatLoop(sctx, out, welcome) })
|
||
spawn("status", func() error { return s.statusLoop(sctx, out) })
|
||
spawn("telemetry", func() error { return s.telemetryLoop(sctx, client, welcome) })
|
||
spawn("discovery", func() error { return s.discoveryLoop(sctx, client) })
|
||
if s.cfg.Syslog != nil || s.cfg.Traps != nil {
|
||
spawn("logs", func() error { return s.logsLoop(sctx, client) })
|
||
}
|
||
|
||
// Читання команд — у цій же горутині.
|
||
readErr := s.controlLoop(sctx, ctrl, out)
|
||
cancel()
|
||
wg.Wait()
|
||
// Збір конфігу переживає скасування контексту не довше за свій
|
||
// дедлайн, але закривати з'єднання під ним не можна: вивантаження
|
||
// впаде на середині, і сервер отримає обірваний стрім замість
|
||
// чесної помилки.
|
||
s.jobs.Wait()
|
||
close(errCh)
|
||
|
||
if readErr != nil && !errors.Is(readErr, context.Canceled) {
|
||
return readErr
|
||
}
|
||
for e := range errCh {
|
||
if e != nil {
|
||
return e
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func (s *Session) hello() *npv1.ControlUp {
|
||
var hash []byte
|
||
if p := s.planHash.Load(); p != nil {
|
||
hash = *p
|
||
}
|
||
return &npv1.ControlUp{
|
||
Seq: s.seq.Add(1),
|
||
Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{
|
||
AgentId: s.cfg.AgentID,
|
||
Build: s.cfg.Build,
|
||
Hostname: s.cfg.Hostname,
|
||
// Хеш плану дозволяє серверу не перезаливати 50 000 задач
|
||
// після кожного обриву зв'язку.
|
||
TaskPlanHash: hash,
|
||
// Продовжуємо з місця розриву, а не з нуля.
|
||
LastAckedBatchId: s.lastAcked.Load(),
|
||
StartedAt: timestamppb.New(s.startedAt),
|
||
}},
|
||
}
|
||
}
|
||
|
||
func (s *Session) applyWelcome(w *npv1.Welcome) {
|
||
if w.ServerTime != nil {
|
||
s.skewNanos.Store(int64(time.Until(w.ServerTime.AsTime())))
|
||
}
|
||
s.cfg.Registry.EnsureDefaults(s.cfg.DefaultModules...)
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Цикли
|
||
// ---------------------------------------------------------------------
|
||
|
||
func (s *Session) writerLoop(ctx context.Context, ctrl npv1.AgentService_ControlClient, out <-chan *npv1.ControlUp) error {
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return nil
|
||
case msg := <-out:
|
||
if err := ctrl.Send(msg); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
func (s *Session) enqueue(ctx context.Context, out chan<- *npv1.ControlUp, msg *npv1.ControlUp) {
|
||
msg.Seq = s.seq.Add(1)
|
||
select {
|
||
case out <- msg:
|
||
case <-ctx.Done():
|
||
default:
|
||
// Черга забита — контрольний канал не встигає. Втратити
|
||
// heartbeat не страшно: сервер помітить це за таймаутом.
|
||
}
|
||
}
|
||
|
||
func (s *Session) heartbeatLoop(ctx context.Context, out chan<- *npv1.ControlUp, w *npv1.Welcome) error {
|
||
interval := 30 * time.Second
|
||
if w.HeartbeatInterval != nil && w.HeartbeatInterval.AsDuration() > 0 {
|
||
interval = w.HeartbeatInterval.AsDuration()
|
||
}
|
||
|
||
t := time.NewTicker(interval)
|
||
defer t.Stop()
|
||
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return nil
|
||
case <-t.C:
|
||
hb := &npv1.Heartbeat{Ts: timestamppb.Now(), Health: s.health()}
|
||
if s.cfg.Scheduler != nil {
|
||
hb.TasksRunning = uint32(s.cfg.Scheduler.Running())
|
||
hb.TasksQueued = uint32(s.cfg.Scheduler.Len())
|
||
}
|
||
s.enqueue(ctx, out, &npv1.ControlUp{
|
||
Payload: &npv1.ControlUp_Heartbeat{Heartbeat: hb},
|
||
})
|
||
s.maybeRenewCredentials(ctx, out)
|
||
}
|
||
}
|
||
}
|
||
|
||
func (s *Session) health() *npv1.AgentHealth {
|
||
var ms runtime.MemStats
|
||
runtime.ReadMemStats(&ms)
|
||
|
||
stats := s.cfg.Buffer.Stats()
|
||
return &npv1.AgentHealth{
|
||
RssBytes: ms.Sys,
|
||
Goroutines: uint32(runtime.NumGoroutine()),
|
||
QueueDepth: uint32(stats.Items),
|
||
DroppedSamples: stats.Dropped,
|
||
Uptime: durationpb.New(time.Since(s.startedAt)),
|
||
ClockSkew: durationpb.New(s.ClockSkew()),
|
||
}
|
||
}
|
||
|
||
func (s *Session) statusLoop(ctx context.Context, out chan<- *npv1.ControlUp) error {
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return nil
|
||
case u := <-s.statusCh:
|
||
s.enqueue(ctx, out, &npv1.ControlUp{
|
||
Payload: &npv1.ControlUp_TaskStatus{TaskStatus: u},
|
||
})
|
||
}
|
||
}
|
||
}
|
||
|
||
// controlLoop читає команди сервера.
|
||
func (s *Session) controlLoop(ctx context.Context, ctrl npv1.AgentService_ControlClient, out chan<- *npv1.ControlUp) error {
|
||
for {
|
||
msg, err := ctrl.Recv()
|
||
if err != nil {
|
||
if ctx.Err() != nil {
|
||
return nil
|
||
}
|
||
return err
|
||
}
|
||
|
||
switch p := msg.Payload.(type) {
|
||
case *npv1.ControlDown_TaskPlan:
|
||
s.applyPlan(p.TaskPlan)
|
||
|
||
case *npv1.ControlDown_TaskDelta:
|
||
s.applyDelta(p.TaskDelta)
|
||
|
||
case *npv1.ControlDown_ModuleControl:
|
||
enabled := make(map[string]bool, len(p.ModuleControl.Modules))
|
||
for _, m := range p.ModuleControl.Modules {
|
||
enabled[m.Key] = m.Enabled
|
||
}
|
||
s.cfg.Registry.SetActive(enabled, p.ModuleControl.Exclusive)
|
||
|
||
case *npv1.ControlDown_Credentials:
|
||
s.applyCredentials(p.Credentials)
|
||
|
||
case *npv1.ControlDown_Ping:
|
||
s.enqueue(ctx, out, &npv1.ControlUp{
|
||
Payload: &npv1.ControlUp_Pong{Pong: &npv1.Pong{
|
||
PingId: p.Ping.PingId,
|
||
AgentTime: timestamppb.Now(),
|
||
}},
|
||
})
|
||
if p.Ping.ServerTime != nil {
|
||
s.skewNanos.Store(int64(time.Until(p.Ping.ServerTime.AsTime())))
|
||
}
|
||
|
||
case *npv1.ControlDown_DiscoveryRequest:
|
||
if s.cfg.Scheduler == nil {
|
||
continue
|
||
}
|
||
// Префікс — ключ модуля з крапкою, а не «topo.»: рядок
|
||
// «topology.identify» на «topo.» не починається, і поштовх
|
||
// не зрушував нічого. Помилку не помічали, бо до появи
|
||
// кнопки «Розпізнати зараз» DiscoveryRequest не слав ніхто.
|
||
n := s.cfg.Scheduler.TriggerNow(p.DiscoveryRequest.GetDeviceIds(), "topology.")
|
||
s.log.Info("сервер попросив запустити автовиявлення",
|
||
"run_id", p.DiscoveryRequest.GetRunId(), "задач_зрушено", n)
|
||
|
||
case *npv1.ControlDown_ConfigJob:
|
||
job := p.ConfigJob
|
||
s.jobs.Add(1)
|
||
go func() {
|
||
defer s.jobs.Done()
|
||
s.runConfigJob(ctx, job)
|
||
}()
|
||
|
||
case *npv1.ControlDown_ConfigApplyJob:
|
||
// Заливка конфігу на пристрій. Окрема гілка, а не ще одне
|
||
// значення config_type у ConfigJob, як зроблено для
|
||
// масових команд: там різниця була лише в тому, що робити
|
||
// з виводом, а тут інша сама природа завдання — ми пишемо
|
||
// на залізо, і звіт їде назад іншим повідомленням.
|
||
apply := p.ConfigApplyJob
|
||
s.jobs.Add(1)
|
||
go func() {
|
||
defer s.jobs.Done()
|
||
s.runApplyJob(ctx, out, apply)
|
||
}()
|
||
|
||
case *npv1.ControlDown_Directive:
|
||
if stop := s.applyDirective(p.Directive); stop {
|
||
return nil
|
||
}
|
||
|
||
case *npv1.ControlDown_Welcome:
|
||
// Повторний Welcome у межах сесії — сервер щось наплутав.
|
||
s.log.Warn("повторний Welcome у живій сесії — ігноруємо")
|
||
}
|
||
}
|
||
}
|
||
|
||
func (s *Session) applyDirective(d *npv1.Directive) (stop bool) {
|
||
switch d.Action {
|
||
case npv1.Directive_ACTION_PAUSE:
|
||
s.log.Info("опитування призупинено сервером", "reason", d.Reason)
|
||
s.cfg.Scheduler.SetPaused(true)
|
||
case npv1.Directive_ACTION_RESUME:
|
||
s.log.Info("опитування відновлено")
|
||
s.cfg.Scheduler.SetPaused(false)
|
||
case npv1.Directive_ACTION_RESET_SERIES_TABLE:
|
||
s.cfg.Buffer.ResetSeries()
|
||
case npv1.Directive_ACTION_RECONNECT, npv1.Directive_ACTION_DRAIN:
|
||
s.log.Info("сервер попросив перепідключитись", "reason", d.Reason)
|
||
return true
|
||
case npv1.Directive_ACTION_UPDATE:
|
||
// Застосування оновлення — окремий крок; без валідного підпису
|
||
// Ed25519 воно не має відбутись узагалі, тому тут лише журнал.
|
||
s.log.Info("доступне оновлення", "version", d.Update.GetVersion())
|
||
}
|
||
return false
|
||
}
|
||
|
||
// credRenewBefore — за скільки до кінця терміну просити нові.
|
||
//
|
||
// Комплект живе годину, heartbeat ходить раз на 30 секунд: десяти
|
||
// хвилин вистачає на кількадесят спроб навіть на поганому каналі.
|
||
const credRenewBefore = 10 * time.Minute
|
||
|
||
// credAskEvery — не частіше, ніж раз на хвилину.
|
||
//
|
||
// Якщо сервер мовчить, питати щотіку означає засипати його однаковими
|
||
// запитами саме тоді, коли йому й так погано.
|
||
const credAskEvery = time.Minute
|
||
|
||
// maybeRenewCredentials просить нову пачку доступів до того, як стара
|
||
// протухне.
|
||
//
|
||
// Без цього зонд працює рівно годину: `Credentials()` свідомо не віддає
|
||
// прострочені, щоб не блокувати облікові записи на пристроях, — і після
|
||
// цього кожна перевірка падає з «немає креденшелів», доки хтось не
|
||
// перезапустить зв'язок. Знайдено живим прогоном: SNMP замовк через
|
||
// годину після підключення, а в журналі не було ані слова про помилку.
|
||
func (s *Session) maybeRenewCredentials(ctx context.Context, out chan<- *npv1.ControlUp) {
|
||
s.credMu.Lock()
|
||
expiry, asked := s.credExpiry, s.credAsked
|
||
now := time.Now()
|
||
// Нульовий термін — комплект без TTL, поновлювати нічого.
|
||
if expiry.IsZero() || now.Add(credRenewBefore).Before(expiry) {
|
||
s.credMu.Unlock()
|
||
return
|
||
}
|
||
if !asked.IsZero() && now.Sub(asked) < credAskEvery {
|
||
s.credMu.Unlock()
|
||
return
|
||
}
|
||
s.credAsked = now
|
||
reason := "expiring"
|
||
if now.After(expiry) {
|
||
reason = "expired"
|
||
}
|
||
s.credMu.Unlock()
|
||
|
||
s.enqueue(ctx, out, &npv1.ControlUp{
|
||
Payload: &npv1.ControlUp_CredentialRequest{
|
||
CredentialRequest: &npv1.CredentialRequest{Reason: reason},
|
||
},
|
||
})
|
||
}
|
||
|
||
func (s *Session) applyCredentials(b *npv1.CredentialBundle) {
|
||
s.credMu.Lock()
|
||
defer s.credMu.Unlock()
|
||
|
||
s.creds = make(map[string][]*npv1.Credential, len(b.ByDevice))
|
||
for devID, list := range b.ByDevice {
|
||
s.creds[devID] = list.Credentials
|
||
}
|
||
if b.ExpiresAt != nil {
|
||
s.credExpiry = b.ExpiresAt.AsTime()
|
||
} else {
|
||
s.credExpiry = time.Time{}
|
||
}
|
||
// Пачка прийшла — наступний запит рахується від нуля.
|
||
s.credAsked = time.Time{}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// План задач
|
||
// ---------------------------------------------------------------------
|
||
|
||
func (s *Session) applyPlan(plan *npv1.TaskPlan) {
|
||
// Часткові плани не застосовуємо: половина розкладу гірша за старий
|
||
// цілий. Сервер має надіслати final = true.
|
||
if !plan.Final {
|
||
s.stashDevices(plan.Devices)
|
||
return
|
||
}
|
||
s.stashDevices(plan.Devices)
|
||
|
||
tasks, meta := s.buildTasks(plan.Tasks)
|
||
s.cfg.Scheduler.Apply(tasks, meta)
|
||
|
||
h := plan.PlanHash
|
||
s.planHash.Store(&h)
|
||
s.log.Info("план задач застосовано", "tasks", len(tasks))
|
||
}
|
||
|
||
func (s *Session) applyDelta(d *npv1.TaskDelta) {
|
||
s.stashDevices(d.UpsertDevices)
|
||
|
||
if len(d.RemoveCheckIds) > 0 {
|
||
s.cfg.Scheduler.Remove(d.RemoveCheckIds)
|
||
}
|
||
if len(d.Upsert) > 0 {
|
||
tasks, meta := s.buildTasks(d.Upsert)
|
||
s.cfg.Scheduler.Upsert(tasks, meta)
|
||
}
|
||
if len(d.RemoveDeviceIds) > 0 {
|
||
s.devMu.Lock()
|
||
for _, id := range d.RemoveDeviceIds {
|
||
delete(s.devices, id)
|
||
}
|
||
s.devMu.Unlock()
|
||
}
|
||
|
||
h := d.PlanHash
|
||
s.planHash.Store(&h)
|
||
}
|
||
|
||
func (s *Session) stashDevices(devs []*npv1.DeviceTarget) {
|
||
if len(devs) == 0 {
|
||
return
|
||
}
|
||
s.devMu.Lock()
|
||
defer s.devMu.Unlock()
|
||
for _, d := range devs {
|
||
s.devices[d.DeviceId] = d
|
||
}
|
||
}
|
||
|
||
func (s *Session) buildTasks(in []*npv1.Task) ([]module.Task, []scheduler.TaskMeta) {
|
||
s.devMu.RLock()
|
||
defer s.devMu.RUnlock()
|
||
|
||
tasks := make([]module.Task, 0, len(in))
|
||
meta := make([]scheduler.TaskMeta, 0, len(in))
|
||
|
||
for _, t := range in {
|
||
if !t.Enabled {
|
||
continue
|
||
}
|
||
dev := s.devices[t.DeviceId]
|
||
if dev == nil {
|
||
// Задача без пристрою — сервер надіслав неповний план.
|
||
s.log.Warn("задача посилається на невідомий пристрій",
|
||
"check_id", t.CheckId, "device_id", t.DeviceId)
|
||
continue
|
||
}
|
||
|
||
tasks = append(tasks, module.Task{
|
||
CheckID: t.CheckId,
|
||
DeviceID: t.DeviceId,
|
||
InterfaceID: t.InterfaceId,
|
||
CheckType: t.CheckType,
|
||
Params: t.ParamsJson,
|
||
Target: module.Target{
|
||
DeviceID: dev.DeviceId,
|
||
Name: dev.Name,
|
||
Address: dev.Address,
|
||
},
|
||
Timeout: t.Timeout.AsDuration(),
|
||
})
|
||
meta = append(meta, scheduler.TaskMeta{
|
||
Interval: t.Interval.AsDuration(),
|
||
Offset: t.ScheduleOffset.AsDuration(),
|
||
})
|
||
}
|
||
return tasks, meta
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Автовиявлення
|
||
// ---------------------------------------------------------------------
|
||
|
||
func (s *Session) discoveryLoop(ctx context.Context, client npv1.AgentServiceClient) error {
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return nil
|
||
case rep := <-s.discoCh:
|
||
ack, err := client.ReportDiscovery(ctx, rep)
|
||
if err != nil {
|
||
if ctx.Err() != nil {
|
||
return nil
|
||
}
|
||
// Звіт втрачено свідомо: наступний обхід за розкладом
|
||
// віддасть свіжіші дані, а ретрай застарілого знімка
|
||
// топології нічого не додає.
|
||
s.log.Warn("звіт автовиявлення не доставлено", "err", err)
|
||
continue
|
||
}
|
||
s.log.Info("автовиявлення прийнято сервером",
|
||
"сусідів", len(rep.Neighbors),
|
||
"зіставлено", ack.GetNeighborsResolved(),
|
||
"лінків", ack.GetLinksCreated())
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Телеметрія
|
||
// ---------------------------------------------------------------------
|
||
|
||
func (s *Session) telemetryLoop(ctx context.Context, client npv1.AgentServiceClient, w *npv1.Welcome) error {
|
||
stream, err := client.StreamTelemetry(ctx)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
maxBatch := int(w.TelemetryMaxBatchSize)
|
||
if maxBatch <= 0 {
|
||
maxBatch = 500
|
||
}
|
||
flush := 5 * time.Second
|
||
if w.TelemetryMaxBatchInterval != nil && w.TelemetryMaxBatchInterval.AsDuration() > 0 {
|
||
flush = w.TelemetryMaxBatchInterval.AsDuration()
|
||
}
|
||
maxInFlight := int(w.TelemetryMaxInFlight)
|
||
if maxInFlight <= 0 {
|
||
maxInFlight = 4
|
||
}
|
||
|
||
var (
|
||
mu sync.Mutex
|
||
inFlight = make(map[uint64]*npv1.TelemetryBatch)
|
||
)
|
||
|
||
// Читач підтверджень. Send і Recv на різних горутинах — це те
|
||
// єдине, що gRPC дозволяє робити паралельно на одному стрімі.
|
||
ackErr := make(chan error, 1)
|
||
go func() {
|
||
for {
|
||
ack, err := stream.Recv()
|
||
if err != nil {
|
||
ackErr <- err
|
||
return
|
||
}
|
||
if ack.ResetSeriesTable {
|
||
s.log.Warn("сервер попросив перереєструвати серії")
|
||
s.cfg.Buffer.ResetSeries()
|
||
}
|
||
mu.Lock()
|
||
for id := range inFlight {
|
||
if id <= ack.AckedThroughBatchId {
|
||
delete(inFlight, id)
|
||
}
|
||
}
|
||
mu.Unlock()
|
||
if ack.AckedThroughBatchId > s.lastAcked.Load() {
|
||
s.lastAcked.Store(ack.AckedThroughBatchId)
|
||
}
|
||
if ack.RetryAfter != nil && ack.RetryAfter.AsDuration() > 0 {
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
case <-time.After(ack.RetryAfter.AsDuration()):
|
||
}
|
||
}
|
||
}
|
||
}()
|
||
|
||
// Усе, що не встигло підтвердитись, повертаємо в буфер:
|
||
// схема БД робить повторний запис безпечним (at-least-once).
|
||
defer func() {
|
||
mu.Lock()
|
||
for _, b := range inFlight {
|
||
b.IsRetransmit = true
|
||
s.cfg.Buffer.Requeue(b)
|
||
}
|
||
mu.Unlock()
|
||
}()
|
||
|
||
ticker := time.NewTicker(flush)
|
||
defer ticker.Stop()
|
||
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return nil
|
||
case err := <-ackErr:
|
||
return err
|
||
case <-ticker.C:
|
||
case <-s.cfg.Buffer.Notify():
|
||
}
|
||
|
||
mu.Lock()
|
||
busy := len(inFlight) >= maxInFlight
|
||
mu.Unlock()
|
||
if busy {
|
||
// Зворотний тиск: сервер не встигає — не додаємо йому роботи.
|
||
continue
|
||
}
|
||
|
||
batchID := s.nextBatch.Add(1)
|
||
batch := s.cfg.Buffer.Drain(batchID, s.cfg.AgentID, maxBatch)
|
||
if batch == nil {
|
||
s.nextBatch.Add(^uint64(0)) // відкотити невикористаний номер
|
||
continue
|
||
}
|
||
|
||
mu.Lock()
|
||
inFlight[batchID] = batch
|
||
mu.Unlock()
|
||
|
||
if err := stream.Send(batch); err != nil {
|
||
mu.Lock()
|
||
delete(inFlight, batchID)
|
||
mu.Unlock()
|
||
s.cfg.Buffer.Requeue(batch)
|
||
return err
|
||
}
|
||
}
|
||
}
|