Netpulse_SasS/agent/internal/telemetry/buffer.go
zotac c8075ee2a6 Етап 2: Go-агент — планувальник, буфер, сесія, модулі ICMP і SNMP
Модулі вкомпільовані (без .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>
2026-08-14 03:46:46 +03:00

329 lines
9.1 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
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
}