Netpulse_SasS/agent/internal/telemetry/interner.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

146 lines
5.3 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 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)
}