Схема PostgreSQL 16+/TimescaleDB: 11 міграцій, 7 схем, топологія (neighbors -> links -> maps -> nodes/edges), time-series з CAGG, NCM, alerting, білінг з entitlements, RLS. Контракт agent<->server: 6 proto-файлів, gRPC, інтернування серій, at-least-once з ack, чанкування конфігів. Перевірено на стенді Debian 13 / PG 17.11 / TimescaleDB 2.29.1: міграції + 8 функціональних перевірок схеми, buf lint + 5 наскрізних gRPC-тестів контракту. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
805 lines
25 KiB
Go
805 lines
25 KiB
Go
// =====================================================================
|
||
// NetPulse :: contract_test.go
|
||
//
|
||
// Наскрізна перевірка контракту агент↔сервер на справжньому gRPC
|
||
// (bufconn, без мережі). Мета — не покрити код тестами, а довести,
|
||
// що модель із .proto реально працює:
|
||
//
|
||
// 1. Рукостискання: агент відкриває стрім, шле Hello, отримує Welcome.
|
||
// 2. Сервер віддає план задач у зустрічному напрямку вже відкритого
|
||
// агентом стріму — тобто "команда з сервера" працює без жодного
|
||
// вхідного з'єднання в мережу клієнта.
|
||
// 3. Інтернування серій: агент реєструє серію один раз, далі шле
|
||
// лише номер; сервер розв'язує його назад у метрику.
|
||
// 4. Зворотний тиск: сервер підтверджує батчі й диктує max_in_flight.
|
||
// 5. Ping/Pong для вимірювання розсинхронізації годинника.
|
||
// 6. Вивантаження конфігу чанками зі звіркою sha256.
|
||
//
|
||
// Запуск: go test ./test/contract/... -v
|
||
// =====================================================================
|
||
|
||
package contract_test
|
||
|
||
import (
|
||
"bytes"
|
||
"context"
|
||
"crypto/sha256"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"io"
|
||
"net"
|
||
"strings"
|
||
"testing"
|
||
"time"
|
||
|
||
"google.golang.org/grpc"
|
||
"google.golang.org/grpc/credentials/insecure"
|
||
"google.golang.org/grpc/test/bufconn"
|
||
"google.golang.org/protobuf/proto"
|
||
"google.golang.org/protobuf/types/known/durationpb"
|
||
"google.golang.org/protobuf/types/known/timestamppb"
|
||
|
||
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
|
||
)
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Мінімальна реалізація серверного боку
|
||
// ---------------------------------------------------------------------
|
||
|
||
// resolvedSample — семпл, уже розв'язаний із series_ref у повний опис.
|
||
// Саме це сервер записав би в ts.series + ts.samples.
|
||
type resolvedSample struct {
|
||
deviceID string
|
||
metricKey string
|
||
unit string
|
||
labels map[string]string
|
||
value float64
|
||
}
|
||
|
||
type storedConfig struct {
|
||
jobID string
|
||
body []byte
|
||
}
|
||
|
||
type testServer struct {
|
||
npv1.UnimplementedAgentServiceServer
|
||
|
||
hello chan *npv1.Hello
|
||
pong chan *npv1.Pong
|
||
samples chan resolvedSample
|
||
icmp chan *npv1.IcmpResult
|
||
ifc chan *npv1.InterfaceCounters
|
||
acks chan uint64
|
||
config chan storedConfig
|
||
planSent chan struct{}
|
||
}
|
||
|
||
func newTestServer() *testServer {
|
||
return &testServer{
|
||
hello: make(chan *npv1.Hello, 1),
|
||
pong: make(chan *npv1.Pong, 1),
|
||
samples: make(chan resolvedSample, 16),
|
||
icmp: make(chan *npv1.IcmpResult, 16),
|
||
ifc: make(chan *npv1.InterfaceCounters, 16),
|
||
acks: make(chan uint64, 16),
|
||
config: make(chan storedConfig, 1),
|
||
planSent: make(chan struct{}, 1),
|
||
}
|
||
}
|
||
|
||
func (s *testServer) Control(stream npv1.AgentService_ControlServer) error {
|
||
// Перше повідомлення зобов'язане бути Hello — інакше сесії немає.
|
||
first, err := stream.Recv()
|
||
if err != nil {
|
||
return err
|
||
}
|
||
hello := first.GetHello()
|
||
if hello == nil {
|
||
return errors.New("перше повідомлення в Control має бути Hello")
|
||
}
|
||
s.hello <- hello
|
||
|
||
if err := stream.Send(&npv1.ControlDown{
|
||
Seq: 1,
|
||
Payload: &npv1.ControlDown_Welcome{Welcome: &npv1.Welcome{
|
||
SessionId: "sess-0001",
|
||
ServerTime: timestamppb.Now(),
|
||
HeartbeatInterval: durationpb.New(30 * time.Second),
|
||
TelemetryMaxBatchSize: 500,
|
||
TelemetryMaxBatchInterval: durationpb.New(5 * time.Second),
|
||
TelemetryMaxInFlight: 4,
|
||
MaxConcurrentChecks: 256,
|
||
IcmpRatePps: 500,
|
||
TaskPlanFollows: true,
|
||
}},
|
||
}); err != nil {
|
||
return err
|
||
}
|
||
|
||
// Параметри чека — непрозорий JSON: ядро їх не тлумачить.
|
||
params, err := json.Marshal(map[string]any{"count": 3, "packet_size": 56})
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
if err := stream.Send(&npv1.ControlDown{
|
||
Seq: 2,
|
||
Payload: &npv1.ControlDown_TaskPlan{TaskPlan: &npv1.TaskPlan{
|
||
PlanHash: []byte("plan-v1"),
|
||
Devices: []*npv1.DeviceTarget{{
|
||
DeviceId: "dev-1", Name: "core-sw-01", Address: "10.0.0.1",
|
||
}},
|
||
Tasks: []*npv1.Task{{
|
||
CheckId: "chk-1",
|
||
DeviceId: "dev-1",
|
||
CheckType: "icmp.ping",
|
||
ParamsJson: params,
|
||
Interval: durationpb.New(60 * time.Second),
|
||
Timeout: durationpb.New(3 * time.Second),
|
||
Retries: 2,
|
||
Enabled: true,
|
||
ScheduleOffset: durationpb.New(7 * time.Second),
|
||
}},
|
||
Final: true,
|
||
}},
|
||
}); err != nil {
|
||
return err
|
||
}
|
||
s.planSent <- struct{}{}
|
||
|
||
for {
|
||
msg, err := stream.Recv()
|
||
if errors.Is(err, io.EOF) {
|
||
return nil
|
||
}
|
||
if err != nil {
|
||
return err
|
||
}
|
||
switch p := msg.Payload.(type) {
|
||
case *npv1.ControlUp_Heartbeat:
|
||
// На heartbeat відповідаємо Ping — так міряється clock skew.
|
||
if err := stream.Send(&npv1.ControlDown{
|
||
Seq: 3,
|
||
Payload: &npv1.ControlDown_Ping{Ping: &npv1.Ping{
|
||
PingId: 42, ServerTime: timestamppb.Now(),
|
||
}},
|
||
}); err != nil {
|
||
return err
|
||
}
|
||
case *npv1.ControlUp_Pong:
|
||
s.pong <- p.Pong
|
||
}
|
||
}
|
||
}
|
||
|
||
func (s *testServer) StreamTelemetry(stream npv1.AgentService_StreamTelemetryServer) error {
|
||
// Таблиця серій жива рівно стільки, скільки сесія. Після реконекту
|
||
// агент нумерує заново — тому вона локальна для цього виклику.
|
||
table := make(map[uint32]*npv1.SeriesDescriptor)
|
||
|
||
for {
|
||
batch, err := stream.Recv()
|
||
if errors.Is(err, io.EOF) {
|
||
return nil
|
||
}
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
for _, d := range batch.NewSeries {
|
||
table[d.SeriesRef] = d
|
||
}
|
||
|
||
for _, smp := range batch.Samples {
|
||
d, ok := table[smp.SeriesRef]
|
||
if !ok {
|
||
// Ми загубили стан — чесно просимо перереєструвати все,
|
||
// замість тихо викидати дані.
|
||
return stream.Send(&npv1.TelemetryAck{
|
||
AckedThroughBatchId: batch.BatchId - 1,
|
||
ResetSeriesTable: true,
|
||
Error: &npv1.Error{
|
||
Code: "unknown_series_ref",
|
||
Message: fmt.Sprintf("series_ref %d не зареєстровано", smp.SeriesRef),
|
||
Retryable: true,
|
||
},
|
||
})
|
||
}
|
||
s.samples <- resolvedSample{
|
||
deviceID: d.DeviceId,
|
||
metricKey: d.MetricKey,
|
||
unit: d.Unit,
|
||
labels: d.Labels,
|
||
value: smp.Value,
|
||
}
|
||
}
|
||
|
||
for _, r := range batch.Icmp {
|
||
s.icmp <- r
|
||
}
|
||
for _, r := range batch.Interfaces {
|
||
s.ifc <- r
|
||
}
|
||
|
||
if err := stream.Send(&npv1.TelemetryAck{
|
||
AckedThroughBatchId: batch.BatchId,
|
||
MaxInFlight: 4,
|
||
}); err != nil {
|
||
return err
|
||
}
|
||
s.acks <- batch.BatchId
|
||
}
|
||
}
|
||
|
||
func (s *testServer) UploadConfig(stream npv1.AgentService_UploadConfigServer) error {
|
||
var hdr *npv1.ConfigHeader
|
||
var body []byte
|
||
var wantSeq uint32
|
||
|
||
for {
|
||
msg, err := stream.Recv()
|
||
if errors.Is(err, io.EOF) {
|
||
return errors.New("стрім завершився без trailer")
|
||
}
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
switch p := msg.Part.(type) {
|
||
case *npv1.ConfigUpload_Header:
|
||
hdr = p.Header
|
||
|
||
case *npv1.ConfigUpload_Chunk:
|
||
if hdr == nil {
|
||
return errors.New("chunk прийшов раніше за header")
|
||
}
|
||
if p.Chunk.Sequence != wantSeq {
|
||
return fmt.Errorf("чанки не по порядку: чекали %d, отримали %d",
|
||
wantSeq, p.Chunk.Sequence)
|
||
}
|
||
wantSeq++
|
||
body = append(body, p.Chunk.Data...)
|
||
|
||
case *npv1.ConfigUpload_Trailer:
|
||
if hdr == nil {
|
||
return errors.New("trailer без header")
|
||
}
|
||
sum := sha256.Sum256(body)
|
||
if !bytes.Equal(sum[:], p.Trailer.ContentSha256) {
|
||
return stream.SendAndClose(&npv1.ConfigReceipt{
|
||
JobId: hdr.JobId,
|
||
Accepted: false,
|
||
Error: &npv1.Error{Code: "checksum_mismatch", Retryable: true},
|
||
})
|
||
}
|
||
if uint64(len(body)) != p.Trailer.SizeBytes {
|
||
return stream.SendAndClose(&npv1.ConfigReceipt{
|
||
JobId: hdr.JobId, Accepted: false,
|
||
Error: &npv1.Error{Code: "size_mismatch"},
|
||
})
|
||
}
|
||
s.config <- storedConfig{jobID: hdr.JobId, body: body}
|
||
return stream.SendAndClose(&npv1.ConfigReceipt{
|
||
JobId: hdr.JobId,
|
||
Accepted: true,
|
||
CommitSha: "3f1a2b8c9d0e",
|
||
})
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Обв'язка
|
||
// ---------------------------------------------------------------------
|
||
|
||
func dial(t *testing.T, s *testServer) npv1.AgentServiceClient {
|
||
t.Helper()
|
||
|
||
lis := bufconn.Listen(1 << 20)
|
||
srv := grpc.NewServer()
|
||
npv1.RegisterAgentServiceServer(srv, s)
|
||
go func() { _ = srv.Serve(lis) }()
|
||
|
||
conn, err := grpc.NewClient("passthrough:///bufnet",
|
||
grpc.WithContextDialer(func(ctx context.Context, _ string) (net.Conn, error) {
|
||
return lis.DialContext(ctx)
|
||
}),
|
||
grpc.WithTransportCredentials(insecure.NewCredentials()),
|
||
)
|
||
if err != nil {
|
||
t.Fatalf("не вдалося створити клієнт: %v", err)
|
||
}
|
||
|
||
t.Cleanup(func() {
|
||
_ = conn.Close()
|
||
srv.Stop()
|
||
_ = lis.Close()
|
||
})
|
||
return npv1.NewAgentServiceClient(conn)
|
||
}
|
||
|
||
// proto1Size — скільки байтів займе на дроті один семпл разом із
|
||
// дескриптором серії (якщо той ще не інтернований).
|
||
func proto1Size(t *testing.T, sample *npv1.MetricSample, desc *npv1.SeriesDescriptor) int {
|
||
t.Helper()
|
||
b, err := proto.Marshal(sample)
|
||
if err != nil {
|
||
t.Fatalf("Marshal(sample): %v", err)
|
||
}
|
||
n := len(b)
|
||
if desc != nil {
|
||
d, err := proto.Marshal(desc)
|
||
if err != nil {
|
||
t.Fatalf("Marshal(descriptor): %v", err)
|
||
}
|
||
n += len(d)
|
||
}
|
||
return n
|
||
}
|
||
|
||
func recvWithin[T any](t *testing.T, ch <-chan T, what string) T {
|
||
t.Helper()
|
||
select {
|
||
case v := <-ch:
|
||
return v
|
||
case <-time.After(5 * time.Second):
|
||
var zero T
|
||
t.Fatalf("не дочекались: %s", what)
|
||
return zero
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 1. Контрольний канал: Hello → Welcome → TaskPlan → Heartbeat → Ping/Pong
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestControlHandshakeAndTaskPush(t *testing.T) {
|
||
srv := newTestServer()
|
||
client := dial(t, srv)
|
||
|
||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := client.Control(ctx)
|
||
if err != nil {
|
||
t.Fatalf("Control: %v", err)
|
||
}
|
||
|
||
// Агент представляється.
|
||
agentStart := time.Now()
|
||
if err := stream.Send(&npv1.ControlUp{
|
||
Seq: 1,
|
||
Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{
|
||
AgentId: "agent-1",
|
||
Hostname: "probe-kyv-01",
|
||
Build: &npv1.AgentBuild{
|
||
Version: "1.0.0", Os: "linux", Arch: "amd64",
|
||
CompiledModules: []string{"icmp", "snmp", "topology"},
|
||
},
|
||
LocalAddresses: []string{"10.0.0.250/24"},
|
||
StartedAt: timestamppb.New(agentStart),
|
||
TaskPlanHash: nil, // плану ще немає — чекаємо повний
|
||
LastAckedBatchId: 0,
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("Send(Hello): %v", err)
|
||
}
|
||
|
||
hello := recvWithin(t, srv.hello, "Hello на сервері")
|
||
if hello.AgentId != "agent-1" {
|
||
t.Fatalf("agent_id = %q, очікували agent-1", hello.AgentId)
|
||
}
|
||
|
||
// Welcome
|
||
down, err := stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Recv(Welcome): %v", err)
|
||
}
|
||
welcome := down.GetWelcome()
|
||
if welcome == nil {
|
||
t.Fatalf("перше повідомлення від сервера — не Welcome, а %T", down.Payload)
|
||
}
|
||
if welcome.SessionId == "" {
|
||
t.Fatal("Welcome без session_id")
|
||
}
|
||
if !welcome.TaskPlanFollows {
|
||
t.Fatal("сервер не попередив, що надішле план")
|
||
}
|
||
if got := welcome.TelemetryMaxBatchSize; got != 500 {
|
||
t.Fatalf("telemetry_max_batch_size = %d, очікували 500", got)
|
||
}
|
||
|
||
// Розсинхронізація годинника — те, заради чого в Welcome є server_time.
|
||
skew := welcome.ServerTime.AsTime().Sub(agentStart)
|
||
if skew > time.Minute {
|
||
t.Fatalf("підозріло великий clock skew: %v", skew)
|
||
}
|
||
|
||
// План задач приходить у зустрічному напрямку вже відкритого агентом
|
||
// стріму — жодного вхідного з'єднання в мережу клієнта.
|
||
down, err = stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Recv(TaskPlan): %v", err)
|
||
}
|
||
plan := down.GetTaskPlan()
|
||
if plan == nil {
|
||
t.Fatalf("очікували TaskPlan, отримали %T", down.Payload)
|
||
}
|
||
if !plan.Final {
|
||
t.Fatal("план не позначений як final — агент не має права його застосовувати")
|
||
}
|
||
if len(plan.Tasks) != 1 || len(plan.Devices) != 1 {
|
||
t.Fatalf("план: %d задач, %d пристроїв", len(plan.Tasks), len(plan.Devices))
|
||
}
|
||
|
||
task := plan.Tasks[0]
|
||
if task.CheckType != "icmp.ping" {
|
||
t.Fatalf("check_type = %q", task.CheckType)
|
||
}
|
||
// Префікс до крапки визначає модуль-виконавця.
|
||
if module, _, ok := strings.Cut(task.CheckType, "."); !ok || module != "icmp" {
|
||
t.Fatalf("з check_type %q не виводиться модуль", task.CheckType)
|
||
}
|
||
// Непрозорі параметри розбирає сам модуль.
|
||
var p struct {
|
||
Count int `json:"count"`
|
||
PacketSize int `json:"packet_size"`
|
||
}
|
||
if err := json.Unmarshal(task.ParamsJson, &p); err != nil {
|
||
t.Fatalf("params_json не розбирається: %v", err)
|
||
}
|
||
if p.Count != 3 || p.PacketSize != 56 {
|
||
t.Fatalf("params_json спотворені: %+v", p)
|
||
}
|
||
if task.Interval.AsDuration() != time.Minute {
|
||
t.Fatalf("interval = %v", task.Interval.AsDuration())
|
||
}
|
||
if task.ScheduleOffset.AsDuration() != 7*time.Second {
|
||
t.Fatalf("schedule_offset = %v", task.ScheduleOffset.AsDuration())
|
||
}
|
||
|
||
<-srv.planSent
|
||
|
||
// Heartbeat із самометриками агента.
|
||
if err := stream.Send(&npv1.ControlUp{
|
||
Seq: 2,
|
||
Payload: &npv1.ControlUp_Heartbeat{Heartbeat: &npv1.Heartbeat{
|
||
Ts: timestamppb.Now(),
|
||
Health: &npv1.AgentHealth{
|
||
CpuPct: 1.5, RssBytes: 22 * 1024 * 1024,
|
||
Goroutines: 48, QueueDepth: 0, ChecksPerSec: 12.5,
|
||
Uptime: durationpb.New(time.Hour),
|
||
},
|
||
TasksRunning: 1,
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("Send(Heartbeat): %v", err)
|
||
}
|
||
|
||
// Сервер відповідає Ping — агент має відповісти Pong.
|
||
down, err = stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Recv(Ping): %v", err)
|
||
}
|
||
ping := down.GetPing()
|
||
if ping == nil {
|
||
t.Fatalf("очікували Ping, отримали %T", down.Payload)
|
||
}
|
||
if err := stream.Send(&npv1.ControlUp{
|
||
Seq: 3,
|
||
Payload: &npv1.ControlUp_Pong{Pong: &npv1.Pong{PingId: ping.PingId, AgentTime: timestamppb.Now()}},
|
||
}); err != nil {
|
||
t.Fatalf("Send(Pong): %v", err)
|
||
}
|
||
|
||
pong := recvWithin(t, srv.pong, "Pong на сервері")
|
||
if pong.PingId != ping.PingId {
|
||
t.Fatalf("ping_id не збігся: %d != %d", pong.PingId, ping.PingId)
|
||
}
|
||
|
||
if err := stream.CloseSend(); err != nil {
|
||
t.Fatalf("CloseSend: %v", err)
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 2. Телеметрія: інтернування серій, гарячі шляхи, ack
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestTelemetrySeriesInterning(t *testing.T) {
|
||
srv := newTestServer()
|
||
client := dial(t, srv)
|
||
|
||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := client.StreamTelemetry(ctx)
|
||
if err != nil {
|
||
t.Fatalf("StreamTelemetry: %v", err)
|
||
}
|
||
|
||
now := time.Now()
|
||
|
||
// Батч 1: реєструємо серію разом із першим семплом.
|
||
batch1 := &npv1.TelemetryBatch{
|
||
BatchId: 1,
|
||
AgentId: "agent-1",
|
||
CreatedAt: timestamppb.New(now),
|
||
NewSeries: []*npv1.SeriesDescriptor{{
|
||
SeriesRef: 1,
|
||
DeviceId: "dev-1",
|
||
PluginKey: "snmp",
|
||
MetricKey: "cpu.util",
|
||
Unit: "pct",
|
||
Labels: map[string]string{"core": "0"},
|
||
}},
|
||
Samples: []*npv1.MetricSample{
|
||
{SeriesRef: 1, Ts: timestamppb.New(now), Value: 37.5},
|
||
},
|
||
Icmp: []*npv1.IcmpResult{{
|
||
DeviceId: "dev-1", CheckId: "chk-1", Ts: timestamppb.New(now),
|
||
RttAvgMs: 1.24, RttMinMs: 1.10, RttMaxMs: 1.55, JitterMs: 0.2,
|
||
LossPct: 0, PacketsSent: 3, PacketsRecv: 3, Reachable: true,
|
||
}},
|
||
Interfaces: []*npv1.InterfaceCounters{{
|
||
DeviceId: "dev-1", InterfaceId: "if-1", Ts: timestamppb.New(now),
|
||
InOctets: 1 << 40, OutOctets: 1 << 39,
|
||
InBps: 420e6, OutBps: 780e6,
|
||
UtilInPct: 42, UtilOutPct: 78,
|
||
OperUp: true, AdminUp: true,
|
||
Interval: durationpb.New(60 * time.Second),
|
||
}},
|
||
}
|
||
if err := stream.Send(batch1); err != nil {
|
||
t.Fatalf("Send(batch1): %v", err)
|
||
}
|
||
|
||
ack, err := stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Recv(ack1): %v", err)
|
||
}
|
||
if ack.AckedThroughBatchId != 1 {
|
||
t.Fatalf("acked_through = %d, очікували 1", ack.AckedThroughBatchId)
|
||
}
|
||
if ack.ResetSeriesTable {
|
||
t.Fatal("сервер попросив скинути таблицю серій одразу після реєстрації")
|
||
}
|
||
if ack.MaxInFlight == 0 {
|
||
t.Fatal("сервер не повідомив max_in_flight — агент не знатиме межі зворотного тиску")
|
||
}
|
||
|
||
got := recvWithin(t, srv.samples, "розв'язаний семпл")
|
||
if got.metricKey != "cpu.util" || got.deviceID != "dev-1" || got.unit != "pct" {
|
||
t.Fatalf("семпл розв'язався неправильно: %+v", got)
|
||
}
|
||
if got.labels["core"] != "0" {
|
||
t.Fatalf("labels загубились: %+v", got.labels)
|
||
}
|
||
if got.value != 37.5 {
|
||
t.Fatalf("value = %v", got.value)
|
||
}
|
||
|
||
icmp := recvWithin(t, srv.icmp, "ICMP-результат")
|
||
if !icmp.Reachable || icmp.PacketsRecv != 3 {
|
||
t.Fatalf("icmp спотворений: %+v", icmp)
|
||
}
|
||
|
||
ifc := recvWithin(t, srv.ifc, "лічильники інтерфейсу")
|
||
if ifc.UtilOutPct != 78 {
|
||
t.Fatalf("util_out_pct = %v — це джерело швидкості анімації, воно має доїхати точно", ifc.UtilOutPct)
|
||
}
|
||
if ifc.Interval.AsDuration() != time.Minute {
|
||
t.Fatalf("interval = %v", ifc.Interval.AsDuration())
|
||
}
|
||
|
||
// Батч 2: серія вже відома — шлемо лише номер.
|
||
// Саме заради цього економія й затівалась.
|
||
batch2 := &npv1.TelemetryBatch{
|
||
BatchId: 2,
|
||
AgentId: "agent-1",
|
||
CreatedAt: timestamppb.New(now.Add(time.Minute)),
|
||
Samples: []*npv1.MetricSample{
|
||
{SeriesRef: 1, Ts: timestamppb.New(now.Add(time.Minute)), Value: 41.0},
|
||
},
|
||
}
|
||
if err := stream.Send(batch2); err != nil {
|
||
t.Fatalf("Send(batch2): %v", err)
|
||
}
|
||
|
||
ack, err = stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Recv(ack2): %v", err)
|
||
}
|
||
if ack.AckedThroughBatchId != 2 {
|
||
t.Fatalf("acked_through = %d, очікували 2", ack.AckedThroughBatchId)
|
||
}
|
||
|
||
got = recvWithin(t, srv.samples, "семпл із батчу без дескриптора")
|
||
if got.metricKey != "cpu.util" || got.value != 41.0 {
|
||
t.Fatalf("серія не розв'язалась із таблиці сесії: %+v", got)
|
||
}
|
||
|
||
// Розмір на дроті: другий батч має бути помітно компактнішим,
|
||
// хоч і несе стільки ж корисних чисел.
|
||
size1 := proto1Size(t, batch1.Samples[0], batch1.NewSeries[0])
|
||
size2 := proto1Size(t, batch2.Samples[0], nil)
|
||
t.Logf("семпл із дескриптором: %d байт, семпл із series_ref: %d байт", size1, size2)
|
||
if size2 >= size1 {
|
||
t.Fatalf("інтернування не дає економії: %d >= %d", size2, size1)
|
||
}
|
||
|
||
if err := stream.CloseSend(); err != nil {
|
||
t.Fatalf("CloseSend: %v", err)
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 3. Невідомий series_ref → сервер чесно просить перереєстрацію
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestTelemetryUnknownSeriesRefTriggersReset(t *testing.T) {
|
||
srv := newTestServer()
|
||
client := dial(t, srv)
|
||
|
||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := client.StreamTelemetry(ctx)
|
||
if err != nil {
|
||
t.Fatalf("StreamTelemetry: %v", err)
|
||
}
|
||
|
||
// Шлемо семпл на серію, яку ніколи не реєстрували.
|
||
if err := stream.Send(&npv1.TelemetryBatch{
|
||
BatchId: 7,
|
||
AgentId: "agent-1",
|
||
Samples: []*npv1.MetricSample{
|
||
{SeriesRef: 999, Ts: timestamppb.Now(), Value: 1},
|
||
},
|
||
}); err != nil {
|
||
t.Fatalf("Send: %v", err)
|
||
}
|
||
|
||
ack, err := stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Recv: %v", err)
|
||
}
|
||
if !ack.ResetSeriesTable {
|
||
t.Fatal("сервер мовчки проковтнув невідомий series_ref замість запиту на перереєстрацію")
|
||
}
|
||
if ack.Error == nil || ack.Error.Code != "unknown_series_ref" {
|
||
t.Fatalf("немає діагностики помилки: %+v", ack.Error)
|
||
}
|
||
if !ack.Error.Retryable {
|
||
t.Fatal("помилка має бути позначена як retryable — дані не втрачені, їх треба переслати")
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 4. NCM: вивантаження конфігу чанками зі звіркою sha256
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestConfigUploadChunked(t *testing.T) {
|
||
srv := newTestServer()
|
||
client := dial(t, srv)
|
||
|
||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := client.UploadConfig(ctx)
|
||
if err != nil {
|
||
t.Fatalf("UploadConfig: %v", err)
|
||
}
|
||
|
||
// Конфіг, більший за один чанк.
|
||
var sb strings.Builder
|
||
sb.WriteString("!\nversion 15.2\nhostname core-sw-01\n!\n")
|
||
for i := 0; i < 500; i++ {
|
||
fmt.Fprintf(&sb, "interface GigabitEthernet0/%d\n description port-%d\n!\n", i, i)
|
||
}
|
||
body := []byte(sb.String())
|
||
want := sha256.Sum256(body)
|
||
|
||
if err := stream.Send(&npv1.ConfigUpload{
|
||
Part: &npv1.ConfigUpload_Header{Header: &npv1.ConfigHeader{
|
||
JobId: "job-1",
|
||
AgentId: "agent-1",
|
||
DeviceId: "dev-1",
|
||
ConfigType: "running",
|
||
CollectedAt: timestamppb.Now(),
|
||
Encoding: "none",
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("Send(header): %v", err)
|
||
}
|
||
|
||
const chunkSize = 4096
|
||
var seq uint32
|
||
for off := 0; off < len(body); off += chunkSize {
|
||
end := min(off+chunkSize, len(body))
|
||
if err := stream.Send(&npv1.ConfigUpload{
|
||
Part: &npv1.ConfigUpload_Chunk{Chunk: &npv1.ConfigChunk{
|
||
Sequence: seq, Data: body[off:end],
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("Send(chunk %d): %v", seq, err)
|
||
}
|
||
seq++
|
||
}
|
||
if seq < 2 {
|
||
t.Fatalf("тест не перевіряє чанкування: вийшов %d чанк", seq)
|
||
}
|
||
|
||
if err := stream.Send(&npv1.ConfigUpload{
|
||
Part: &npv1.ConfigUpload_Trailer{Trailer: &npv1.ConfigTrailer{
|
||
Success: true,
|
||
ContentSha256: want[:],
|
||
SizeBytes: uint64(len(body)),
|
||
LineCount: uint32(bytes.Count(body, []byte("\n"))),
|
||
ChunkCount: seq,
|
||
Duration: durationpb.New(2 * time.Second),
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("Send(trailer): %v", err)
|
||
}
|
||
|
||
receipt, err := stream.CloseAndRecv()
|
||
if err != nil {
|
||
t.Fatalf("CloseAndRecv: %v", err)
|
||
}
|
||
if !receipt.Accepted {
|
||
t.Fatalf("сервер відхилив конфіг: %+v", receipt.Error)
|
||
}
|
||
if receipt.CommitSha == "" {
|
||
t.Fatal("немає commit_sha — конфіг не потрапив у Git")
|
||
}
|
||
|
||
stored := recvWithin(t, srv.config, "збережений конфіг")
|
||
if !bytes.Equal(stored.body, body) {
|
||
t.Fatalf("тіло спотворилось: %d байт замість %d", len(stored.body), len(body))
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 5. Пошкоджений конфіг має бути відхилений, а не тихо збережений
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestConfigUploadRejectsBadChecksum(t *testing.T) {
|
||
srv := newTestServer()
|
||
client := dial(t, srv)
|
||
|
||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := client.UploadConfig(ctx)
|
||
if err != nil {
|
||
t.Fatalf("UploadConfig: %v", err)
|
||
}
|
||
|
||
body := []byte("hostname core-sw-01\n")
|
||
bogus := sha256.Sum256([]byte("щось зовсім інше"))
|
||
|
||
_ = stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Header{
|
||
Header: &npv1.ConfigHeader{JobId: "job-2", DeviceId: "dev-1", ConfigType: "running"},
|
||
}})
|
||
_ = stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Chunk{
|
||
Chunk: &npv1.ConfigChunk{Sequence: 0, Data: body},
|
||
}})
|
||
_ = stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Trailer{
|
||
Trailer: &npv1.ConfigTrailer{
|
||
Success: true, ContentSha256: bogus[:], SizeBytes: uint64(len(body)), ChunkCount: 1,
|
||
},
|
||
}})
|
||
|
||
receipt, err := stream.CloseAndRecv()
|
||
if err != nil {
|
||
t.Fatalf("CloseAndRecv: %v", err)
|
||
}
|
||
if receipt.Accepted {
|
||
t.Fatal("сервер прийняв конфіг із неправильною контрольною сумою")
|
||
}
|
||
if receipt.Error == nil || receipt.Error.Code != "checksum_mismatch" {
|
||
t.Fatalf("немає зрозумілої причини відмови: %+v", receipt.Error)
|
||
}
|
||
}
|