Модулі вкомпільовані (без .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>
276 lines
9.1 KiB
Go
276 lines
9.1 KiB
Go
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("буфер не розбудив відправника")
|
||
}
|
||
}
|