Модулі вкомпільовані (без .so): один бінарник на Alpine, Windows і роутер. Креденшели беруться на момент виконання, бо мають TTL. Розклад вирівняний по сітці інтервалу, тому переживає рестарт. Буфер обмежений і за кількістю, і за пам'яттю; при переповненні викидає найстаріше, зміни статусу — останніми. Перевірено на Debian 13 / Go 1.25: go vet чисто, go test -race усі пакети ok. Релізний бінарник 12 МБ, базовий RSS 11.6 МБ при бюджеті 30 МБ. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
329 lines
9.1 KiB
Go
329 lines
9.1 KiB
Go
package telemetry
|
||
|
||
import (
|
||
"sync"
|
||
"time"
|
||
|
||
"github.com/netpulse/netpulse/agent/internal/module"
|
||
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
|
||
"google.golang.org/protobuf/types/known/timestamppb"
|
||
)
|
||
|
||
// Приблизна вага одного запису на дроті й у пам'яті. Точний proto.Size()
|
||
// на кожен запис — O(n) на виклик і з'їдає більше CPU, ніж економить
|
||
// пам'яті. Для бюджету RSS вистачає оцінки з запасом.
|
||
const (
|
||
sampleWeight = 48
|
||
icmpWeight = 96
|
||
ifcWeight = 192
|
||
seriesWeight = 160
|
||
)
|
||
|
||
// Buffer — обмежений буфер результатів між модулями й відправником.
|
||
//
|
||
// Головна вимога: агент не має права рости в пам'яті, коли зв'язок із
|
||
// сервером пропав. Тому буфер обмежений і за кількістю записів, і за
|
||
// оціночним обсягом; при переповненні викидаються НАЙСТАРІШІ дані.
|
||
//
|
||
// Чому найстаріші, а не найновіші: коли зв'язок відновиться, оператору
|
||
// потрібен передусім поточний стан мережі — свіжий семпл цінніший за
|
||
// півгодинної давнини діру, яку однаково вже ніхто не побачить у
|
||
// реальному часі. Кількість викинутого їде в AgentHealth.dropped_samples,
|
||
// щоб діра була видима, а не мовчазна.
|
||
type Buffer struct {
|
||
mu sync.Mutex
|
||
|
||
maxItems int
|
||
maxBytes int
|
||
|
||
samples []*npv1.MetricSample
|
||
icmp []*npv1.IcmpResult
|
||
interfaces []*npv1.InterfaceCounters
|
||
statuses []*npv1.StatusChange
|
||
checks []*npv1.CheckResult
|
||
|
||
bytes int
|
||
dropped uint64
|
||
|
||
interner *Interner
|
||
notify chan struct{}
|
||
}
|
||
|
||
type Options struct {
|
||
MaxItems int
|
||
MaxBytes int
|
||
Interner *Interner
|
||
}
|
||
|
||
func NewBuffer(opts Options) *Buffer {
|
||
if opts.MaxItems <= 0 {
|
||
opts.MaxItems = 50_000
|
||
}
|
||
if opts.MaxBytes <= 0 {
|
||
// ~8 МБ на дані при бюджеті RSS 30 МБ: решта йде на рантайм,
|
||
// пули з'єднань і gRPC-фрейми.
|
||
opts.MaxBytes = 8 << 20
|
||
}
|
||
if opts.Interner == nil {
|
||
opts.Interner = NewInterner()
|
||
}
|
||
return &Buffer{
|
||
maxItems: opts.MaxItems,
|
||
maxBytes: opts.MaxBytes,
|
||
interner: opts.Interner,
|
||
notify: make(chan struct{}, 1),
|
||
}
|
||
}
|
||
|
||
func (b *Buffer) Interner() *Interner { return b.interner }
|
||
|
||
// Notify сигналить відправнику, що є дані. Канал з буфером 1 —
|
||
// зайві сигнали не накопичуються.
|
||
func (b *Buffer) Notify() <-chan struct{} { return b.notify }
|
||
|
||
func (b *Buffer) signal() {
|
||
select {
|
||
case b.notify <- struct{}{}:
|
||
default:
|
||
}
|
||
}
|
||
|
||
// Add кладе результат виконання задачі в буфер.
|
||
func (b *Buffer) Add(deviceID, pluginKey string, res module.Result) {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
|
||
for _, m := range res.Metrics {
|
||
ref, isNew := b.interner.Ref(deviceID, pluginKey, m)
|
||
if isNew {
|
||
b.bytes += seriesWeight
|
||
}
|
||
ts := m.Ts
|
||
if ts.IsZero() {
|
||
ts = time.Now()
|
||
}
|
||
b.samples = append(b.samples, &npv1.MetricSample{
|
||
SeriesRef: ref,
|
||
Ts: timestamppb.New(ts),
|
||
Value: m.Value,
|
||
})
|
||
b.bytes += sampleWeight
|
||
}
|
||
|
||
if res.Icmp != nil {
|
||
b.icmp = append(b.icmp, res.Icmp)
|
||
b.bytes += icmpWeight
|
||
}
|
||
for _, ifc := range res.Interfaces {
|
||
b.interfaces = append(b.interfaces, ifc)
|
||
b.bytes += ifcWeight
|
||
}
|
||
|
||
b.enforceLimitsLocked()
|
||
b.signal()
|
||
}
|
||
|
||
// AddStatusChange — зміна стану пристрою. Йде окремо, бо має потрапити
|
||
// на сервер якнайшвидше: від неї залежить колір вузла на мапі.
|
||
func (b *Buffer) AddStatusChange(sc *npv1.StatusChange) {
|
||
b.mu.Lock()
|
||
b.statuses = append(b.statuses, sc)
|
||
b.bytes += icmpWeight
|
||
b.enforceLimitsLocked()
|
||
b.mu.Unlock()
|
||
b.signal()
|
||
}
|
||
|
||
// AddCheckResult — службовий результат чека (успіх/помилка/тривалість).
|
||
func (b *Buffer) AddCheckResult(cr *npv1.CheckResult) {
|
||
b.mu.Lock()
|
||
b.checks = append(b.checks, cr)
|
||
b.bytes += icmpWeight
|
||
b.enforceLimitsLocked()
|
||
b.mu.Unlock()
|
||
b.signal()
|
||
}
|
||
|
||
func (b *Buffer) count() int {
|
||
return len(b.samples) + len(b.icmp) + len(b.interfaces) + len(b.statuses) + len(b.checks)
|
||
}
|
||
|
||
// enforceLimitsLocked викидає найстаріші записи, доки буфер не влізе
|
||
// в обидва ліміти.
|
||
func (b *Buffer) enforceLimitsLocked() {
|
||
for (b.count() > b.maxItems || b.bytes > b.maxBytes) && b.count() > 0 {
|
||
switch {
|
||
case len(b.samples) > 0:
|
||
b.samples = b.samples[1:]
|
||
b.bytes -= sampleWeight
|
||
case len(b.interfaces) > 0:
|
||
b.interfaces = b.interfaces[1:]
|
||
b.bytes -= ifcWeight
|
||
case len(b.icmp) > 0:
|
||
b.icmp = b.icmp[1:]
|
||
b.bytes -= icmpWeight
|
||
case len(b.checks) > 0:
|
||
b.checks = b.checks[1:]
|
||
b.bytes -= icmpWeight
|
||
default:
|
||
// Зміни стану викидаємо останніми: без них мапа
|
||
// показуватиме пристрій живим, поки він лежить.
|
||
b.statuses = b.statuses[1:]
|
||
b.bytes -= icmpWeight
|
||
}
|
||
b.dropped++
|
||
}
|
||
if b.bytes < 0 {
|
||
b.bytes = 0
|
||
}
|
||
}
|
||
|
||
// Stats — те, що йде в AgentHealth.
|
||
type Stats struct {
|
||
Items int
|
||
Bytes int
|
||
Dropped uint64
|
||
Series int
|
||
}
|
||
|
||
func (b *Buffer) Stats() Stats {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
return Stats{
|
||
Items: b.count(),
|
||
Bytes: b.bytes,
|
||
Dropped: b.dropped,
|
||
Series: b.interner.Len(),
|
||
}
|
||
}
|
||
|
||
// Drain формує батч не більший за maxItems записів.
|
||
//
|
||
// Повертає nil, якщо даних немає, — відправнику нема чого слати.
|
||
// Дескриптори нових серій потрапляють у той самий батч, що й семпли,
|
||
// які на них посилаються.
|
||
func (b *Buffer) Drain(batchID uint64, agentID string, maxItems int) *npv1.TelemetryBatch {
|
||
if maxItems <= 0 {
|
||
maxItems = 500
|
||
}
|
||
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
|
||
if b.count() == 0 {
|
||
return nil
|
||
}
|
||
|
||
batch := &npv1.TelemetryBatch{
|
||
BatchId: batchID,
|
||
AgentId: agentID,
|
||
CreatedAt: timestamppb.Now(),
|
||
NewSeries: b.interner.TakePending(),
|
||
}
|
||
|
||
budget := maxItems
|
||
|
||
// Зміни стану — першими: вони найцінніші для мапи.
|
||
take := minInt(budget, len(b.statuses))
|
||
if take > 0 {
|
||
batch.StatusChanges = b.statuses[:take:take]
|
||
b.statuses = b.statuses[take:]
|
||
b.bytes -= take * icmpWeight
|
||
budget -= take
|
||
}
|
||
|
||
take = minInt(budget, len(b.icmp))
|
||
if take > 0 {
|
||
batch.Icmp = b.icmp[:take:take]
|
||
b.icmp = b.icmp[take:]
|
||
b.bytes -= take * icmpWeight
|
||
budget -= take
|
||
}
|
||
|
||
take = minInt(budget, len(b.interfaces))
|
||
if take > 0 {
|
||
batch.Interfaces = b.interfaces[:take:take]
|
||
b.interfaces = b.interfaces[take:]
|
||
b.bytes -= take * ifcWeight
|
||
budget -= take
|
||
}
|
||
|
||
take = minInt(budget, len(b.samples))
|
||
if take > 0 {
|
||
batch.Samples = b.samples[:take:take]
|
||
b.samples = b.samples[take:]
|
||
b.bytes -= take * sampleWeight
|
||
budget -= take
|
||
}
|
||
|
||
take = minInt(budget, len(b.checks))
|
||
if take > 0 {
|
||
batch.CheckResults = b.checks[:take:take]
|
||
b.checks = b.checks[take:]
|
||
b.bytes -= take * icmpWeight
|
||
}
|
||
|
||
if b.bytes < 0 {
|
||
b.bytes = 0
|
||
}
|
||
// Лишилось ще — не даємо відправнику заснути.
|
||
if b.count() > 0 {
|
||
b.signal()
|
||
}
|
||
return batch
|
||
}
|
||
|
||
// Requeue повертає невідправлений батч у голову буфера.
|
||
//
|
||
// Викликається, коли стрім обірвався між Send і Ack. Дескриптори серій
|
||
// теж повертаються — інакше сервер отримав би семпли з невідомим ref.
|
||
func (b *Buffer) Requeue(batch *npv1.TelemetryBatch) {
|
||
if batch == nil {
|
||
return
|
||
}
|
||
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
|
||
b.interner.ReturnPending(batch.NewSeries)
|
||
|
||
b.statuses = append(batch.StatusChanges, b.statuses...)
|
||
b.icmp = append(batch.Icmp, b.icmp...)
|
||
b.interfaces = append(batch.Interfaces, b.interfaces...)
|
||
b.samples = append(batch.Samples, b.samples...)
|
||
b.checks = append(batch.CheckResults, b.checks...)
|
||
|
||
b.bytes += len(batch.StatusChanges)*icmpWeight +
|
||
len(batch.Icmp)*icmpWeight +
|
||
len(batch.Interfaces)*ifcWeight +
|
||
len(batch.Samples)*sampleWeight +
|
||
len(batch.CheckResults)*icmpWeight
|
||
|
||
b.enforceLimitsLocked()
|
||
b.signal()
|
||
}
|
||
|
||
// ResetSeries — сервер попросив перереєструвати серії.
|
||
//
|
||
// Разом із таблицею доводиться викинути й накопичені семпли: вони
|
||
// посилаються на номери, яких сервер більше не знає. Це чесніше, ніж
|
||
// відправити їх у нікуди, — і воно видиме в dropped.
|
||
func (b *Buffer) ResetSeries() {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
|
||
b.dropped += uint64(len(b.samples))
|
||
b.bytes -= len(b.samples) * sampleWeight
|
||
b.samples = nil
|
||
b.interner.Reset()
|
||
if b.bytes < 0 {
|
||
b.bytes = 0
|
||
}
|
||
}
|
||
|
||
func minInt(a, b int) int {
|
||
if a < b {
|
||
return a
|
||
}
|
||
return b
|
||
}
|