Netpulse_SasS/agent/internal/session/session.go
byrsapty ed8fc831bf Дві сесії роботи: 0058–0068, розгортання однією командою, тести
Один коміт, а не десяток тематичних, свідомо: теми переплетені в
спільних файлах (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 серпня.
2026-08-27 17:32:49 +03:00

864 lines
28 KiB
Go
Raw Permalink 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.

// 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
}
}
}