Netpulse_SasS/agent/internal/session/session.go
zotac 8a92ef8a45 Етап 2: серверна сторона AgentService
Автентифікація зондів за токеном, побудова TaskPlan із детермінованим
schedule_offset, видача розшифрованих креденшелів із TTL, запис телеметрії
в гіпертаблиці, резолвер сусідів LLDP/CDP у topo.links, прийом конфігів.

Ізоляція тенантів робиться двічі — RLS плюс явний предикат tenant_id, бо
RLS не працює на гіпертаблях, а саме туди йде вся телеметрія.

Агент: -token і передача його в метаданих; MarkAllPending() перереєстровує
серії на початку сесії замість обнуляти нумерацію й губити буфер.

Перевірено на Debian 13 / PG 17.11 / TimescaleDB 2.29.1: 11 інтеграційних
тестів проти живої БД (-race), плюс живий прогін справжнього агента проти
справжнього сервера — телеметрія, статус пристрою, heartbeat у базі.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-14 04:13:40 +03:00

662 lines
19 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.

// 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/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
}
type Session struct {
cfg Config
log *slog.Logger
// Статуси задач — діагностика. Канал з буфером і скиданням при
// переповненні: краще втратити рядок журналу, ніж загальмувати
// опитування через повільний контрольний канал.
statusCh chan *npv1.TaskStatusUpdate
credMu sync.RWMutex
creds map[string][]*npv1.Credential
credExpiry time.Time
devMu sync.RWMutex
devices map[string]*npv1.DeviceTarget
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
}
return &Session{
cfg: cfg,
log: cfg.Logger,
statusCh: make(chan *npv1.TaskStatusUpdate, 512),
creds: make(map[string][]*npv1.Credential),
devices: make(map[string]*npv1.DeviceTarget),
startedAt: time.Now(),
}
}
// 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]
}
// 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)
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
errCh := make(chan error, 4)
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) })
// Читання команд — у цій же горутині.
readErr := s.controlLoop(sctx, ctrl, out)
cancel()
wg.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},
})
}
}
}
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_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
}
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{}
}
}
// ---------------------------------------------------------------------
// План задач
// ---------------------------------------------------------------------
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) 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
}
}
}