Автентифікація зондів за токеном, побудова 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>
146 lines
5.3 KiB
Go
146 lines
5.3 KiB
Go
// Package telemetry — накопичення результатів опитування та підготовка
|
||
// їх до відправки: інтернування серій і буфер із бюджетом пам'яті.
|
||
package telemetry
|
||
|
||
import (
|
||
"sync"
|
||
|
||
"github.com/netpulse/netpulse/agent/internal/module"
|
||
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
|
||
)
|
||
|
||
// Interner присвоює серіям короткі номери в межах сесії.
|
||
//
|
||
// Навіщо: повторювати device_id (36 байт) + metric_key + labels у кожному
|
||
// семплі задорого. Реєструємо серію раз, далі шлемо лише uint32.
|
||
// Заміряно в test/contract: 66 → 25 байт на семпл.
|
||
//
|
||
// Номер дійсний рівно в межах сесії. Після реконекту або на вимогу
|
||
// сервера (TelemetryAck.reset_series_table) таблиця обнуляється й усі
|
||
// серії реєструються заново.
|
||
type Interner struct {
|
||
mu sync.Mutex
|
||
next uint32
|
||
// канонічний ключ серії → присвоєний ref
|
||
refs map[string]uint32
|
||
// усі створені дескриптори за ref — потрібні, щоб після
|
||
// реконекту перереєструвати серії, не втрачаючи буфер
|
||
descs map[uint32]*npv1.SeriesDescriptor
|
||
// дескриптори, які ще не поїхали на сервер
|
||
pending []*npv1.SeriesDescriptor
|
||
}
|
||
|
||
func NewInterner() *Interner {
|
||
return &Interner{
|
||
refs: make(map[string]uint32),
|
||
descs: make(map[uint32]*npv1.SeriesDescriptor),
|
||
}
|
||
}
|
||
|
||
// Ref повертає номер серії, реєструючи її за потреби.
|
||
//
|
||
// Другим значенням — чи серія нова. Викликач має покласти дескриптор
|
||
// (через TakePending) у ТОЙ САМИЙ батч, що й перший семпл: інакше
|
||
// сервер отримає посилання на невідомий ref і попросить повний скид.
|
||
func (i *Interner) Ref(deviceID, pluginKey string, m module.Metric) (uint32, bool) {
|
||
key := m.SeriesKey(deviceID)
|
||
|
||
i.mu.Lock()
|
||
defer i.mu.Unlock()
|
||
|
||
if ref, ok := i.refs[key]; ok {
|
||
return ref, false
|
||
}
|
||
|
||
i.next++
|
||
ref := i.next
|
||
i.refs[key] = ref
|
||
|
||
// Копіюємо labels: модуль може перевикористати мапу між викликами.
|
||
var labels map[string]string
|
||
if len(m.Labels) > 0 {
|
||
labels = make(map[string]string, len(m.Labels))
|
||
for k, v := range m.Labels {
|
||
labels[k] = v
|
||
}
|
||
}
|
||
|
||
desc := &npv1.SeriesDescriptor{
|
||
SeriesRef: ref,
|
||
DeviceId: deviceID,
|
||
InterfaceId: m.InterfaceID,
|
||
PluginKey: pluginKey,
|
||
MetricKey: m.MetricKey,
|
||
Unit: m.Unit,
|
||
Labels: labels,
|
||
}
|
||
i.descs[ref] = desc
|
||
i.pending = append(i.pending, desc)
|
||
return ref, true
|
||
}
|
||
|
||
// TakePending забирає дескриптори, які ще не відправлені.
|
||
func (i *Interner) TakePending() []*npv1.SeriesDescriptor {
|
||
i.mu.Lock()
|
||
defer i.mu.Unlock()
|
||
|
||
if len(i.pending) == 0 {
|
||
return nil
|
||
}
|
||
out := i.pending
|
||
i.pending = nil
|
||
return out
|
||
}
|
||
|
||
// ReturnPending повертає дескриптори назад, якщо батч не вдалося
|
||
// відправити. Без цього серія лишилась би зареєстрованою локально,
|
||
// але невідомою серверу — і всі наступні семпли летіли б у нікуди.
|
||
func (i *Interner) ReturnPending(descs []*npv1.SeriesDescriptor) {
|
||
if len(descs) == 0 {
|
||
return
|
||
}
|
||
i.mu.Lock()
|
||
defer i.mu.Unlock()
|
||
i.pending = append(descs, i.pending...)
|
||
}
|
||
|
||
// MarkAllPending повертає ВСІ відомі серії в чергу на відправку.
|
||
//
|
||
// Викликається на початку кожної нової сесії. Таблиця series_ref живе
|
||
// на сервері рівно стільки, скільки сесія, тому після реконекту він
|
||
// нічого не пам'ятає. Альтернатива — обнулити нумерацію на агенті —
|
||
// означала б викинути весь накопичений за час обриву буфер: саме ті
|
||
// дані, заради яких він і накопичувався.
|
||
func (i *Interner) MarkAllPending() {
|
||
i.mu.Lock()
|
||
defer i.mu.Unlock()
|
||
|
||
if len(i.descs) == 0 {
|
||
return
|
||
}
|
||
pending := make([]*npv1.SeriesDescriptor, 0, len(i.descs))
|
||
for _, d := range i.descs {
|
||
pending = append(pending, d)
|
||
}
|
||
i.pending = pending
|
||
}
|
||
|
||
// Reset обнуляє таблицю: сервер попросив перереєстрацію з нуля.
|
||
func (i *Interner) Reset() {
|
||
i.mu.Lock()
|
||
defer i.mu.Unlock()
|
||
|
||
i.next = 0
|
||
i.refs = make(map[string]uint32)
|
||
i.descs = make(map[uint32]*npv1.SeriesDescriptor)
|
||
i.pending = nil
|
||
}
|
||
|
||
// Len — скільки серій зареєстровано. Використовується в самометриках:
|
||
// нескінченне зростання означає витік кардинальності (напр. плагін
|
||
// пхає timestamp у labels).
|
||
func (i *Interner) Len() int {
|
||
i.mu.Lock()
|
||
defer i.mu.Unlock()
|
||
return len(i.refs)
|
||
}
|