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

276 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 (
"testing"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/protobuf/types/known/timestamppb"
)
func metric(key string, v float64, labels map[string]string) module.Metric {
return module.Metric{MetricKey: key, Unit: "pct", Value: v, Labels: labels, Ts: time.Now()}
}
// ---------------------------------------------------------------------
// Interner
// ---------------------------------------------------------------------
func TestInternerAssignsStableRefs(t *testing.T) {
in := NewInterner()
ref1, isNew1 := in.Ref("dev-1", "snmp", metric("cpu.util", 1, map[string]string{"core": "0"}))
ref2, isNew2 := in.Ref("dev-1", "snmp", metric("cpu.util", 2, map[string]string{"core": "0"}))
if !isNew1 {
t.Fatal("перша реєстрація не позначена як нова")
}
if isNew2 {
t.Fatal("та сама серія зареєстрована двічі")
}
if ref1 != ref2 {
t.Fatalf("та сама серія отримала різні ref: %d != %d", ref1, ref2)
}
// Інший label — інша серія.
ref3, isNew3 := in.Ref("dev-1", "snmp", metric("cpu.util", 3, map[string]string{"core": "1"}))
if !isNew3 || ref3 == ref1 {
t.Fatalf("різні labels дали той самий ref: %d", ref3)
}
// Інший пристрій — інша серія.
ref4, _ := in.Ref("dev-2", "snmp", metric("cpu.util", 4, map[string]string{"core": "0"}))
if ref4 == ref1 {
t.Fatal("серії різних пристроїв злилися")
}
}
// Порядок ключів у map недетермінований — канонічний ключ має це пережити.
func TestInternerLabelOrderDoesNotMatter(t *testing.T) {
in := NewInterner()
a := metric("sensor.temp", 1, map[string]string{"a": "1", "b": "2", "c": "3"})
b := metric("sensor.temp", 2, map[string]string{"c": "3", "a": "1", "b": "2"})
refA, _ := in.Ref("dev-1", "snmp", a)
refB, isNew := in.Ref("dev-1", "snmp", b)
if isNew || refA != refB {
t.Fatalf("порядок labels зробив нову серію: %d vs %d", refA, refB)
}
}
func TestInternerPendingRoundTrip(t *testing.T) {
in := NewInterner()
in.Ref("dev-1", "snmp", metric("cpu.util", 1, nil))
pending := in.TakePending()
if len(pending) != 1 {
t.Fatalf("pending = %d", len(pending))
}
if len(in.TakePending()) != 0 {
t.Fatal("TakePending віддав ті самі дескриптори двічі")
}
// Батч не поїхав — дескриптори мають повернутись, інакше серія
// лишиться відомою агенту й невідомою серверу.
in.ReturnPending(pending)
if len(in.TakePending()) != 1 {
t.Fatal("ReturnPending загубив дескриптор")
}
}
func TestInternerResetRenumbers(t *testing.T) {
in := NewInterner()
in.Ref("dev-1", "snmp", metric("cpu.util", 1, nil))
in.Ref("dev-2", "snmp", metric("cpu.util", 1, nil))
if in.Len() != 2 {
t.Fatalf("Len = %d", in.Len())
}
in.Reset()
if in.Len() != 0 {
t.Fatal("Reset не очистив таблицю")
}
ref, isNew := in.Ref("dev-1", "snmp", metric("cpu.util", 1, nil))
if !isNew || ref != 1 {
t.Fatalf("після Reset нумерація не почалась з 1: ref=%d new=%v", ref, isNew)
}
}
// ---------------------------------------------------------------------
// Buffer
// ---------------------------------------------------------------------
func TestBufferDrainCarriesDescriptorsWithSamples(t *testing.T) {
b := NewBuffer(Options{})
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", 50, nil)},
})
batch := b.Drain(1, "agent-1", 100)
if batch == nil {
t.Fatal("Drain повернув nil при непорожньому буфері")
}
if len(batch.Samples) != 1 {
t.Fatalf("семплів у батчі: %d", len(batch.Samples))
}
// Дескриптор мусить бути в ТОМУ САМОМУ батчі, інакше сервер отримає
// посилання на невідомий ref і попросить повний скид.
if len(batch.NewSeries) != 1 {
t.Fatalf("дескрипторів у батчі: %d", len(batch.NewSeries))
}
if batch.NewSeries[0].SeriesRef != batch.Samples[0].SeriesRef {
t.Fatal("ref дескриптора не збігається з ref семпла")
}
if b.Drain(2, "agent-1", 100) != nil {
t.Fatal("буфер не спорожнів після Drain")
}
}
func TestBufferPrioritisesStatusChanges(t *testing.T) {
b := NewBuffer(Options{})
for i := 0; i < 20; i++ {
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", float64(i), nil)},
})
}
b.AddStatusChange(&npv1.StatusChange{
DeviceId: "dev-1", Ts: timestamppb.Now(),
Status: npv1.Status_STATUS_DOWN, PreviousStatus: npv1.Status_STATUS_UP,
})
// Місця лише на 2 записи — зміна стану має пролізти першою:
// без неї мапа показуватиме пристрій живим, поки він лежить.
batch := b.Drain(1, "agent-1", 2)
if len(batch.StatusChanges) != 1 {
t.Fatalf("зміна стану не потрапила в перший батч: %+v", batch.StatusChanges)
}
}
func TestBufferDropsOldestOnOverflow(t *testing.T) {
b := NewBuffer(Options{MaxItems: 10, MaxBytes: 1 << 20})
for i := 0; i < 25; i++ {
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", float64(i), nil)},
})
}
st := b.Stats()
if st.Items > 10 {
t.Fatalf("буфер переріс ліміт: %d записів", st.Items)
}
if st.Dropped == 0 {
t.Fatal("викидання не пораховано — діра в даних буде невидимою")
}
// Лишитись мають НАЙНОВІШІ: після відновлення зв'язку оператору
// потрібен поточний стан, а не півгодинної давнини.
batch := b.Drain(1, "agent-1", 100)
if len(batch.Samples) == 0 {
t.Fatal("порожній батч")
}
last := batch.Samples[len(batch.Samples)-1]
if last.Value != 24 {
t.Fatalf("останній семпл = %v, очікували 24 (найновіший)", last.Value)
}
}
func TestBufferRespectsByteBudget(t *testing.T) {
// Бюджету вистачає приблизно на 3 семпли.
b := NewBuffer(Options{MaxItems: 1000, MaxBytes: seriesWeight + 3*sampleWeight})
for i := 0; i < 50; i++ {
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", float64(i), nil)},
})
}
st := b.Stats()
if st.Bytes > seriesWeight+3*sampleWeight {
t.Fatalf("буфер перевищив бюджет пам'яті: %d байт", st.Bytes)
}
if st.Items > 4 {
t.Fatalf("у буфері %d записів — бюджет не тримається", st.Items)
}
}
func TestBufferRequeuePreservesData(t *testing.T) {
b := NewBuffer(Options{})
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", 7, nil)},
Icmp: &npv1.IcmpResult{DeviceId: "dev-1", Reachable: true},
})
batch := b.Drain(1, "agent-1", 100)
if batch == nil {
t.Fatal("Drain повернув nil")
}
// Стрім обірвався між Send і Ack.
b.Requeue(batch)
again := b.Drain(2, "agent-1", 100)
if again == nil {
t.Fatal("після Requeue буфер порожній — дані втрачено")
}
if len(again.Samples) != 1 || again.Samples[0].Value != 7 {
t.Fatalf("семпл не повернувся: %+v", again.Samples)
}
if len(again.Icmp) != 1 {
t.Fatal("ICMP-результат не повернувся")
}
// Дескриптор теж має повернутись, інакше семпл поїде з ref,
// якого сервер ніколи не бачив.
if len(again.NewSeries) != 1 {
t.Fatalf("дескриптор серії не повернувся: %+v", again.NewSeries)
}
}
func TestBufferResetSeriesDropsOrphanSamples(t *testing.T) {
b := NewBuffer(Options{})
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", 1, nil)},
})
before := b.Stats()
b.ResetSeries()
after := b.Stats()
if after.Series != 0 {
t.Fatal("таблиця серій не очищена")
}
// Семпли посилаються на номери, яких сервер більше не знає —
// чесніше викинути й показати це в dropped, ніж слати в нікуди.
if after.Dropped <= before.Dropped {
t.Fatal("викинуті семпли не пораховані")
}
if b.Drain(1, "agent-1", 100) != nil {
t.Fatal("осиротілі семпли лишились у буфері")
}
}
func TestBufferNotifySignalsProducer(t *testing.T) {
b := NewBuffer(Options{})
select {
case <-b.Notify():
t.Fatal("порожній буфер сигналить про дані")
default:
}
b.Add("dev-1", "snmp", module.Result{Metrics: []module.Metric{metric("cpu.util", 1, nil)}})
select {
case <-b.Notify():
case <-time.After(time.Second):
t.Fatal("буфер не розбудив відправника")
}
}