Один коміт, а не десяток тематичних, свідомо: теми переплетені в
спільних файлах (store.go, docker-compose.yml, deploy/README.md), і
розділити їх можна було б лише індексуванням шматків. Коміти, які не
збираються, гірші за один великий — тим паче що це рівно той стан, який
перевірявся разом.
ЩО ПРАЦЮЄ НА СТЕНДІ Й ПЕРЕВІРЕНО ТАМ
0058 подієві алерти: syslog, ncm, compliance спрацьовують у мить
події; правило з нереалізованим джерелом більше не зберігається
мовчки
0059 snmp.walk і прототипи шаблонів — таблиці з динамічним індексом
описуються шаблоном, а не Go
0060 відкат конфігу: план як різниця, маскування паролів із підписом
плану, обов'язковий контрольний збір, verifying при обриві
0061 кнопки Telegram: довге опитування, авторизація не з callback_data
0062 аудит і архів хостів; тест на AST, що падає на ключі без назви
0063 RLS: три ролі, окремий пул для фонових тактів
0064 строки зберігання даних і сторінка сховища
0065 приймач SNMP-трапів; перевірено справжніми пакетами по дроту,
переклад v1→v2 за RFC 3584 дає правильний OID
0066 ескалації сповіщень
0067 алерт про вичерпання диска
0068 поля заливки конфігу переїхали в каталог профілів
Плюс: 137 тестів вебу з нуля (їх не було взагалі), одинадцять справжніх
вад, знайдених ними й виправлених, і виправлення двох інтеграційних
тестів grpcapi, які мовчки пропускались півтора року.
ЩО ЩЕ НЕ ЗАПУСКАЛОСЬ
netpulse установник: одна команда замість 18 змінних і
593 рядків інструкції
RLS з першого запуску нова інсталяція під політиками одразу;
RLS-EXISTING-INSTALL.md лишається тільки для
старих інсталяцій
.forgejo + CI раннер не зареєстрований
Ці три перевірені компіляцією й міркуванням, але не виконанням.
ГОЛОВНИЙ ВИСНОВОК ДВОХ СЕСІЙ
Зелена перевірка доводить рівно те, що вона перевіряє. Тест ізоляції RLS
був правильний і зелений — і пропустив зламаний вхід, бо перевіряв «чи
не видно чужого», коли зламалось «чи видно своє». Інтеграційні тести
grpcapi були зелені, бо не виконувались. Схема, довідник і протокол
описували те, чого в коді не існувало, і виглядало це як готове.
Тому в кожному завданні цих сесій стояла вимога назвати НЕПОКРИТЕ, а
чотири задачі закінчились не можливістю, а відмовою: правило з
нереалізованим джерелом не зберігається, профіль без команд заливки
каже про це замість мовчазної кнопки, міграція RLS валить сама себе на
таблиці без політики, тест словника аудиту падає на ключі без назви.
Подробиці — HISTORY.md, розділи за 26 і 27 серпня.
1159 lines
41 KiB
Go
1159 lines
41 KiB
Go
// Наскрізний тест серверної сторони: справжній PostgreSQL/TimescaleDB,
|
||
// справжній gRPC, справжня схема з Етапу 1.
|
||
//
|
||
// Запуск:
|
||
//
|
||
// NETPULSE_TEST_DSN="postgres://netpulse:netpulse@localhost/netpulse_it" go test ./...
|
||
//
|
||
// Без змінної тест пропускається — щоб `go test ./...` не падав там,
|
||
// де бази немає.
|
||
package grpcapi_test
|
||
|
||
import (
|
||
"context"
|
||
"crypto/sha256"
|
||
"encoding/hex"
|
||
"encoding/json"
|
||
"fmt"
|
||
"log/slog"
|
||
"net"
|
||
"os"
|
||
"testing"
|
||
"time"
|
||
|
||
"github.com/jackc/pgx/v5/pgxpool"
|
||
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
|
||
"github.com/netpulse/netpulse/server/internal/crypto"
|
||
"github.com/netpulse/netpulse/server/internal/grpcapi"
|
||
"github.com/netpulse/netpulse/server/internal/store"
|
||
"google.golang.org/grpc"
|
||
"google.golang.org/grpc/credentials/insecure"
|
||
"google.golang.org/grpc/metadata"
|
||
"google.golang.org/grpc/test/bufconn"
|
||
"google.golang.org/protobuf/types/known/durationpb"
|
||
"google.golang.org/protobuf/types/known/timestamppb"
|
||
)
|
||
|
||
type fixture struct {
|
||
pool *pgxpool.Pool
|
||
store *store.Store
|
||
ring *crypto.Keyring
|
||
client npv1.AgentServiceClient
|
||
ctx context.Context
|
||
tenantID string
|
||
agentID string
|
||
deviceID string
|
||
ifaceID string
|
||
peerID string
|
||
peerIfID string
|
||
checkID string
|
||
token string
|
||
}
|
||
|
||
func setup(t *testing.T) *fixture {
|
||
t.Helper()
|
||
|
||
dsn := os.Getenv("NETPULSE_TEST_DSN")
|
||
if dsn == "" {
|
||
t.Skip("NETPULSE_TEST_DSN не задано — інтеграційний тест пропущено")
|
||
}
|
||
|
||
ctx := context.Background()
|
||
|
||
st, err := store.New(ctx, dsn)
|
||
if err != nil {
|
||
t.Fatalf("підключення до БД: %v", err)
|
||
}
|
||
t.Cleanup(st.Close)
|
||
|
||
ring := crypto.NewKeyring()
|
||
key, err := crypto.GenerateKey()
|
||
if err != nil {
|
||
t.Fatalf("GenerateKey: %v", err)
|
||
}
|
||
if err := ring.Add("test-key", key); err != nil {
|
||
t.Fatalf("Add: %v", err)
|
||
}
|
||
|
||
f := &fixture{pool: st.Pool(), store: st, ring: ring, ctx: ctx}
|
||
f.seed(t)
|
||
|
||
// Сервер із тими самими перехоплювачами, що й у проді.
|
||
svc := grpcapi.New(st, ring,
|
||
slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelWarn})))
|
||
|
||
lis := bufconn.Listen(1 << 20)
|
||
srv := grpc.NewServer(
|
||
grpc.ChainUnaryInterceptor(svc.UnaryInterceptor),
|
||
grpc.ChainStreamInterceptor(svc.StreamInterceptor),
|
||
)
|
||
npv1.RegisterAgentServiceServer(srv, svc)
|
||
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()
|
||
})
|
||
|
||
f.client = npv1.NewAgentServiceClient(conn)
|
||
return f
|
||
}
|
||
|
||
// authCtx додає токен зонда так само, як це робить агент.
|
||
func (f *fixture) authCtx() context.Context {
|
||
return metadata.AppendToOutgoingContext(f.ctx, "authorization", "Bearer "+f.token)
|
||
}
|
||
|
||
func (f *fixture) seed(t *testing.T) {
|
||
t.Helper()
|
||
ctx := f.ctx
|
||
|
||
slug := fmt.Sprintf("it-%d", time.Now().UnixNano())
|
||
f.token = "np_test_" + slug
|
||
sum := sha256.Sum256([]byte(f.token))
|
||
|
||
must := func(q string, args ...any) {
|
||
t.Helper()
|
||
if _, err := f.pool.Exec(ctx, q, args...); err != nil {
|
||
t.Fatalf("seed %q: %v", q, err)
|
||
}
|
||
}
|
||
scan := func(dest *string, q string, args ...any) {
|
||
t.Helper()
|
||
if err := f.pool.QueryRow(ctx, q, args...).Scan(dest); err != nil {
|
||
t.Fatalf("seed %q: %v", q, err)
|
||
}
|
||
}
|
||
|
||
// slug має домен core.slug, name — text. Один і той самий $1 для
|
||
// обох Postgres вивести не може, тому передаємо окремо.
|
||
scan(&f.tenantID, `
|
||
INSERT INTO core.tenants (slug, name, status) VALUES ($1, $2, 'active')
|
||
RETURNING id::text`, slug, slug)
|
||
|
||
t.Cleanup(func() {
|
||
// Каскади прибирають усе, крім гіпертаблиць — у них немає FK.
|
||
_, _ = f.pool.Exec(context.Background(),
|
||
`DELETE FROM ts.icmp_samples WHERE tenant_id = $1`, f.tenantID)
|
||
_, _ = f.pool.Exec(context.Background(),
|
||
`DELETE FROM ts.if_counters WHERE tenant_id = $1`, f.tenantID)
|
||
_, _ = f.pool.Exec(context.Background(),
|
||
`DELETE FROM ts.device_status_history WHERE tenant_id = $1`, f.tenantID)
|
||
_, _ = f.pool.Exec(context.Background(),
|
||
`DELETE FROM ts.agent_health WHERE tenant_id = $1`, f.tenantID)
|
||
_, _ = f.pool.Exec(context.Background(),
|
||
`DELETE FROM ts.samples WHERE series_id IN (SELECT id FROM ts.series WHERE tenant_id = $1)`, f.tenantID)
|
||
_, _ = f.pool.Exec(context.Background(),
|
||
`DELETE FROM core.tenants WHERE id = $1`, f.tenantID)
|
||
})
|
||
|
||
scan(&f.agentID, `
|
||
INSERT INTO core.agents (tenant_id, name, token_hash, status, enabled_modules)
|
||
VALUES ($1, 'probe-it', $2, 'online', '{icmp,snmp,topology}')
|
||
RETURNING id::text`, f.tenantID, sum[:])
|
||
|
||
scan(&f.deviceID, `
|
||
INSERT INTO inv.devices (tenant_id, agent_id, name, address, kind, vendor, chassis_id, system_name)
|
||
VALUES ($1, $2, 'core-sw-01', '10.10.0.1', 'switch', 'cisco', '0011.2233.4455', 'core-sw-01')
|
||
RETURNING id::text`, f.tenantID, f.agentID)
|
||
|
||
// Сусід, якого автовиявлення має знайти.
|
||
scan(&f.peerID, `
|
||
INSERT INTO inv.devices (tenant_id, agent_id, name, address, kind, chassis_id, system_name)
|
||
VALUES ($1, $2, 'edge-rtr-01', '10.10.0.2', 'router', '0011.2233.6677', 'edge-rtr-01')
|
||
RETURNING id::text`, f.tenantID, f.agentID)
|
||
|
||
scan(&f.ifaceID, `
|
||
INSERT INTO inv.interfaces (tenant_id, device_id, if_index, name, speed_bps, oper_status)
|
||
VALUES ($1, $2, 1, 'GigabitEthernet0/1', 1000000000, 'up')
|
||
RETURNING id::text`, f.tenantID, f.deviceID)
|
||
|
||
scan(&f.peerIfID, `
|
||
INSERT INTO inv.interfaces (tenant_id, device_id, if_index, name, speed_bps, oper_status)
|
||
VALUES ($1, $2, 1, 'ether1', 1000000000, 'up')
|
||
RETURNING id::text`, f.tenantID, f.peerID)
|
||
|
||
scan(&f.checkID, `
|
||
INSERT INTO core.checks (tenant_id, device_id, check_type, params, interval_sec, timeout_ms, retries)
|
||
VALUES ($1, $2, 'icmp.ping', '{"count":3}'::jsonb, 60, 3000, 2)
|
||
RETURNING id::text`, f.tenantID, f.deviceID)
|
||
|
||
// Креденшел зі справжнім шифруванням: тест має пройти той самий
|
||
// шлях, що й прод, включно з AAD.
|
||
aad := f.tenantID + "|snmp_v2c|" + f.deviceID
|
||
sec, err := f.ring.Encrypt([]byte("public-it"), aad)
|
||
if err != nil {
|
||
t.Fatalf("Encrypt: %v", err)
|
||
}
|
||
|
||
var secretID string
|
||
scan(&secretID, `
|
||
INSERT INTO core.secrets (tenant_id, kind, key_id, nonce, ciphertext, auth_tag, aad)
|
||
VALUES ($1, 'snmp_v3', $2, $3, $4, $5, $6)
|
||
RETURNING id::text`,
|
||
f.tenantID, sec.KeyID, sec.Nonce, sec.Ciphertext, sec.AuthTag, aad)
|
||
|
||
var credID string
|
||
scan(&credID, `
|
||
INSERT INTO inv.credentials (tenant_id, name, proto, username, port, secret_id)
|
||
VALUES ($1, 'snmp-ro', 'snmp_v2c', '', 161, $2)
|
||
RETURNING id::text`, f.tenantID, secretID)
|
||
|
||
must(`INSERT INTO inv.device_credentials (device_id, credential_id, priority)
|
||
VALUES ($1, $2, 10)`, f.deviceID, credID)
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 1. Рукостискання, план задач, креденшели
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestControlHandshake(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := f.client.Control(ctx)
|
||
if err != nil {
|
||
t.Fatalf("Control: %v", err)
|
||
}
|
||
|
||
if err := stream.Send(&npv1.ControlUp{
|
||
Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{
|
||
AgentId: f.agentID,
|
||
Hostname: "probe-it",
|
||
Build: &npv1.AgentBuild{Version: "test", Os: "linux", Arch: "amd64"},
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("Hello: %v", err)
|
||
}
|
||
|
||
down, err := stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Welcome: %v", err)
|
||
}
|
||
w := down.GetWelcome()
|
||
if w == nil {
|
||
t.Fatalf("замість Welcome прийшло %T", down.Payload)
|
||
}
|
||
if !w.TaskPlanFollows {
|
||
t.Fatal("сервер не збирається слати план, хоча агент прийшов без хеша")
|
||
}
|
||
if w.TelemetryMaxBatchSize == 0 || w.MaxConcurrentChecks == 0 {
|
||
t.Fatalf("ліміти не заповнені: %+v", w)
|
||
}
|
||
|
||
var (
|
||
plan *npv1.TaskPlan
|
||
creds *npv1.CredentialBundle
|
||
mods *npv1.ModuleControl
|
||
)
|
||
deadline := time.Now().Add(10 * time.Second)
|
||
for (plan == nil || creds == nil || mods == nil) && time.Now().Before(deadline) {
|
||
msg, err := stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Recv: %v", err)
|
||
}
|
||
switch p := msg.Payload.(type) {
|
||
case *npv1.ControlDown_TaskPlan:
|
||
plan = p.TaskPlan
|
||
case *npv1.ControlDown_Credentials:
|
||
creds = p.Credentials
|
||
case *npv1.ControlDown_ModuleControl:
|
||
mods = p.ModuleControl
|
||
}
|
||
}
|
||
|
||
if plan == nil {
|
||
t.Fatal("план задач не надійшов")
|
||
}
|
||
// У плані не один чек, а два. Другий — `topology.identify`, і його
|
||
// заводить сам сервер при підключенні зонда (service.go,
|
||
// EnsureIdentifyChecks): хост зі SNMP-доступом отримує розпізнавання
|
||
// без жодного натискання.
|
||
//
|
||
// Тест писався до появи розпізнавання й перевіряв рівність одиниці.
|
||
// Півтора року він цього не помічав, бо мовчки пропускався без
|
||
// NETPULSE_TEST_DSN — перший же справжній прогін його завалив.
|
||
// Шукаємо СВІЙ чек серед решти, а не покладаємось на їхню кількість:
|
||
// наступний автоматичний чек інакше завалить його знову.
|
||
var task *npv1.Task
|
||
for _, tk := range plan.Tasks {
|
||
if tk.CheckId == f.checkID {
|
||
task = tk
|
||
}
|
||
}
|
||
if task == nil {
|
||
t.Fatalf("свого чека в плані немає: %+v", plan.Tasks)
|
||
}
|
||
if task.CheckType != "icmp.ping" {
|
||
t.Fatalf("check_type = %q", task.CheckType)
|
||
}
|
||
if task.Interval.AsDuration() != time.Minute {
|
||
t.Fatalf("interval = %v", task.Interval.AsDuration())
|
||
}
|
||
// Offset має бути в межах інтервалу — інакше задача або стартує
|
||
// одразу разом з усіма, або з'їде в наступний цикл.
|
||
off := task.ScheduleOffset.AsDuration()
|
||
if off < 0 || off >= time.Minute {
|
||
t.Fatalf("schedule_offset поза інтервалом: %v", off)
|
||
}
|
||
// І він має бути детермінованим: та сама функція на будь-якому вузлі.
|
||
if want := store.ScheduleOffset(f.checkID, time.Minute); off != want {
|
||
t.Fatalf("offset не детермінований: %v != %v", off, want)
|
||
}
|
||
|
||
if len(plan.Devices) != 1 || plan.Devices[0].Address != "10.10.0.1" {
|
||
t.Fatalf("пристрої в плані: %+v", plan.Devices)
|
||
}
|
||
if len(plan.PlanHash) == 0 {
|
||
t.Fatal("план без хеша — агент не зможе уникнути перезаливки")
|
||
}
|
||
|
||
if mods == nil || len(mods.Modules) == 0 {
|
||
t.Fatal("сервер не сказав, які модулі вмикати")
|
||
}
|
||
|
||
if creds == nil {
|
||
t.Fatal("креденшели не надійшли")
|
||
}
|
||
list := creds.ByDevice[f.deviceID]
|
||
if list == nil || len(list.Credentials) != 1 {
|
||
t.Fatalf("креденшели пристрою: %+v", creds.ByDevice)
|
||
}
|
||
// Найважливіше: секрет розшифрувався тим самим ключем і AAD.
|
||
if got := list.Credentials[0].GetCommunity(); got != "public-it" {
|
||
t.Fatalf("community = %q, очікували public-it", got)
|
||
}
|
||
if creds.ExpiresAt == nil || creds.ExpiresAt.AsTime().Before(time.Now()) {
|
||
t.Fatal("креденшели без дійсного TTL")
|
||
}
|
||
|
||
// Зонд позначений онлайн у БД.
|
||
var status string
|
||
if err := f.pool.QueryRow(f.ctx,
|
||
`SELECT status::text FROM core.agents WHERE id = $1`, f.agentID).Scan(&status); err != nil {
|
||
t.Fatalf("статус зонда: %v", err)
|
||
}
|
||
if status != "online" {
|
||
t.Fatalf("статус зонда = %q", status)
|
||
}
|
||
}
|
||
|
||
// Токен визначає зонда. Назватись чужим id не можна навіть із валідним
|
||
// токеном — інакше один зонд отримав би план іншого.
|
||
func TestControlRejectsMismatchedAgentID(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 10*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := f.client.Control(ctx)
|
||
if err != nil {
|
||
t.Fatalf("Control: %v", err)
|
||
}
|
||
if err := stream.Send(&npv1.ControlUp{
|
||
Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{
|
||
AgentId: "00000000-0000-0000-0000-0000000000ff",
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("Hello: %v", err)
|
||
}
|
||
if _, err := stream.Recv(); err == nil {
|
||
t.Fatal("сервер прийняв чужий agent_id")
|
||
}
|
||
}
|
||
|
||
func TestUnauthenticatedRejected(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.ctx, 10*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := f.client.Control(ctx)
|
||
if err != nil {
|
||
return // клієнт міг відмовити одразу
|
||
}
|
||
_ = stream.Send(&npv1.ControlUp{Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{}}})
|
||
if _, err := stream.Recv(); err == nil {
|
||
t.Fatal("сервер пустив сесію без токена")
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 2. Телеметрія доходить до гіпертаблиць
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestTelemetryPersisted(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := f.client.StreamTelemetry(ctx)
|
||
if err != nil {
|
||
t.Fatalf("StreamTelemetry: %v", err)
|
||
}
|
||
|
||
now := time.Now().Truncate(time.Millisecond)
|
||
|
||
batch := &npv1.TelemetryBatch{
|
||
BatchId: 1,
|
||
AgentId: f.agentID,
|
||
CreatedAt: timestamppb.New(now),
|
||
NewSeries: []*npv1.SeriesDescriptor{{
|
||
SeriesRef: 1,
|
||
DeviceId: f.deviceID,
|
||
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: f.deviceID, CheckId: f.checkID, Ts: timestamppb.New(now),
|
||
RttAvgMs: 1.25, RttMinMs: 1.1, RttMaxMs: 1.5, JitterMs: 0.2,
|
||
LossPct: 0, PacketsSent: 3, PacketsRecv: 3, Reachable: true,
|
||
}},
|
||
Interfaces: []*npv1.InterfaceCounters{{
|
||
DeviceId: f.deviceID, InterfaceId: f.ifaceID, Ts: timestamppb.New(now),
|
||
InOctets: 1 << 30, OutOctets: 1 << 29,
|
||
InBps: 420e6, OutBps: 780e6,
|
||
UtilInPct: 42, UtilOutPct: 78, OperUp: true,
|
||
Interval: durationpb.New(time.Minute),
|
||
}},
|
||
}
|
||
|
||
if err := stream.Send(batch); err != nil {
|
||
t.Fatalf("Send: %v", err)
|
||
}
|
||
ack, err := stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Recv ack: %v", err)
|
||
}
|
||
if ack.AckedThroughBatchId != 1 {
|
||
t.Fatalf("acked = %d, помилка: %+v", ack.AckedThroughBatchId, ack.Error)
|
||
}
|
||
|
||
// Семпл у ts.samples із розв'язаним series_id.
|
||
var value float64
|
||
var metricKey, unit string
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT s.value, se.metric_key, se.unit
|
||
FROM ts.samples s JOIN ts.series se ON se.id = s.series_id
|
||
WHERE se.tenant_id = $1 AND se.metric_key = 'cpu.util'
|
||
ORDER BY s.ts DESC LIMIT 1`, f.tenantID).Scan(&value, &metricKey, &unit); err != nil {
|
||
t.Fatalf("семпл не записався: %v", err)
|
||
}
|
||
if value != 37.5 || unit != "pct" {
|
||
t.Fatalf("семпл спотворений: value=%v unit=%q", value, unit)
|
||
}
|
||
|
||
// ICMP у широкій таблиці — те, що читає мапа.
|
||
var rtt float32
|
||
var reachable bool
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT rtt_avg_ms, reachable FROM ts.icmp_samples
|
||
WHERE device_id = $1 ORDER BY ts DESC LIMIT 1`, f.deviceID).Scan(&rtt, &reachable); err != nil {
|
||
t.Fatalf("ICMP не записався: %v", err)
|
||
}
|
||
if !reachable || rtt < 1.2 || rtt > 1.3 {
|
||
t.Fatalf("ICMP спотворений: rtt=%v reachable=%v", rtt, reachable)
|
||
}
|
||
|
||
// Завантаження інтерфейсу — джерело швидкості анімації на мапі.
|
||
var utilOut float32
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT util_out_pct FROM ts.if_counters
|
||
WHERE interface_id = $1 ORDER BY ts DESC LIMIT 1`, f.ifaceID).Scan(&utilOut); err != nil {
|
||
t.Fatalf("лічильники не записались: %v", err)
|
||
}
|
||
if utilOut != 78 {
|
||
t.Fatalf("util_out_pct = %v", utilOut)
|
||
}
|
||
|
||
// Статус пристрою піднявся з ICMP і лишив слід в історії.
|
||
var devStatus string
|
||
if err := f.pool.QueryRow(f.ctx,
|
||
`SELECT status::text FROM inv.devices WHERE id = $1`, f.deviceID).Scan(&devStatus); err != nil {
|
||
t.Fatalf("статус пристрою: %v", err)
|
||
}
|
||
if devStatus != "up" {
|
||
t.Fatalf("статус пристрою = %q, очікували up", devStatus)
|
||
}
|
||
|
||
var transitions int
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT count(*) FROM ts.device_status_history
|
||
WHERE device_id = $1 AND status = 'up'`, f.deviceID).Scan(&transitions); err != nil {
|
||
t.Fatalf("історія станів: %v", err)
|
||
}
|
||
if transitions != 1 {
|
||
t.Fatalf("переходів в історії: %d, очікували 1", transitions)
|
||
}
|
||
|
||
// Повтор того самого батчу (at-least-once після реконекту) не має
|
||
// ані падати, ані дублювати рядки.
|
||
batch.IsRetransmit = true
|
||
batch.BatchId = 2
|
||
batch.NewSeries = nil
|
||
if err := stream.Send(batch); err != nil {
|
||
t.Fatalf("повторний Send: %v", err)
|
||
}
|
||
if _, err := stream.Recv(); err != nil {
|
||
t.Fatalf("повторний ack: %v", err)
|
||
}
|
||
|
||
var icmpRows int
|
||
if err := f.pool.QueryRow(f.ctx,
|
||
`SELECT count(*) FROM ts.icmp_samples WHERE device_id = $1`, f.deviceID).Scan(&icmpRows); err != nil {
|
||
t.Fatalf("підрахунок: %v", err)
|
||
}
|
||
if icmpRows != 1 {
|
||
t.Fatalf("повтор батчу продублював рядки: %d", icmpRows)
|
||
}
|
||
|
||
// І жодного зайвого переходу в історії від повтору.
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT count(*) FROM ts.device_status_history WHERE device_id = $1`,
|
||
f.deviceID).Scan(&transitions); err != nil {
|
||
t.Fatalf("історія: %v", err)
|
||
}
|
||
if transitions != 1 {
|
||
t.Fatalf("повтор додав перехід в історію: %d", transitions)
|
||
}
|
||
}
|
||
|
||
// Невідомий series_ref → чесний запит на перереєстрацію.
|
||
func TestTelemetryUnknownSeriesRef(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := f.client.StreamTelemetry(ctx)
|
||
if err != nil {
|
||
t.Fatalf("StreamTelemetry: %v", err)
|
||
}
|
||
|
||
if err := stream.Send(&npv1.TelemetryBatch{
|
||
BatchId: 5, AgentId: f.agentID,
|
||
Samples: []*npv1.MetricSample{{SeriesRef: 404, 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)
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 3. Автовиявлення зводить сусідів у лінк
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestDiscoveryResolvesLink(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
ack, err := f.client.ReportDiscovery(ctx, &npv1.DiscoveryReport{
|
||
AgentId: f.agentID,
|
||
Final: true,
|
||
Neighbors: []*npv1.NeighborRecord{{
|
||
DeviceId: f.deviceID,
|
||
LocalInterfaceId: f.ifaceID,
|
||
Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_LLDP,
|
||
RemoteChassisId: "0011.2233.6677",
|
||
RemoteSystemName: "edge-rtr-01",
|
||
RemotePortId: "ether1",
|
||
RemoteCapabilities: []string{"bridge", "router"},
|
||
SeenAt: timestamppb.Now(),
|
||
}},
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("ReportDiscovery: %v", err)
|
||
}
|
||
if !ack.Accepted {
|
||
t.Fatalf("звіт відхилено: %+v", ack.Error)
|
||
}
|
||
if ack.NeighborsResolved != 1 {
|
||
t.Fatalf("зіставлено сусідів: %d", ack.NeighborsResolved)
|
||
}
|
||
if ack.LinksCreated != 1 {
|
||
t.Fatalf("створено лінків: %d", ack.LinksCreated)
|
||
}
|
||
|
||
// Сусід зберігся з резолвленими посиланнями й високою впевненістю:
|
||
// збіг за chassis-id надійніший за збіг за іменем.
|
||
var resolvedDev, resolvedIf string
|
||
var confidence int
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT resolved_device_id::text, resolved_interface_id::text, confidence
|
||
FROM topo.neighbors WHERE device_id = $1`, f.deviceID).
|
||
Scan(&resolvedDev, &resolvedIf, &confidence); err != nil {
|
||
t.Fatalf("сусід не записався: %v", err)
|
||
}
|
||
if resolvedDev != f.peerID {
|
||
t.Fatalf("сусід зіставлений не з тим пристроєм: %s", resolvedDev)
|
||
}
|
||
if resolvedIf != f.peerIfID {
|
||
t.Fatalf("порт сусіда не зіставлений: %s", resolvedIf)
|
||
}
|
||
if confidence < 90 {
|
||
t.Fatalf("впевненість збігу за chassis-id = %d, очікували >= 90", confidence)
|
||
}
|
||
|
||
// Лінк створився з портами й пропускною здатністю з інтерфейсу.
|
||
var capacity int64
|
||
var discoveredBy string
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT COALESCE(capacity_bps, 0), discovered_by::text
|
||
FROM topo.links WHERE tenant_id = $1`, f.tenantID).Scan(&capacity, &discoveredBy); err != nil {
|
||
t.Fatalf("лінк не створився: %v", err)
|
||
}
|
||
if capacity != 1000000000 {
|
||
t.Fatalf("capacity_bps = %d — анімація трафіку не матиме знаменника", capacity)
|
||
}
|
||
if discoveredBy != "lldp" {
|
||
t.Fatalf("discovered_by = %q", discoveredBy)
|
||
}
|
||
|
||
// Повторний звіт не має плодити другий лінк — навіть якщо сусід
|
||
// прийшов з іншого боку (B→A).
|
||
if _, err := f.client.ReportDiscovery(ctx, &npv1.DiscoveryReport{
|
||
AgentId: f.agentID, Final: true,
|
||
Neighbors: []*npv1.NeighborRecord{{
|
||
DeviceId: f.peerID,
|
||
LocalInterfaceId: f.peerIfID,
|
||
Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_LLDP,
|
||
RemoteChassisId: "0011.2233.4455",
|
||
RemotePortId: "GigabitEthernet0/1",
|
||
SeenAt: timestamppb.Now(),
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("повторний звіт: %v", err)
|
||
}
|
||
|
||
var links int
|
||
if err := f.pool.QueryRow(f.ctx,
|
||
`SELECT count(*) FROM topo.links WHERE tenant_id = $1`, f.tenantID).Scan(&links); err != nil {
|
||
t.Fatalf("підрахунок лінків: %v", err)
|
||
}
|
||
if links != 1 {
|
||
t.Fatalf("зустрічний звіт створив дублікат: лінків %d", links)
|
||
}
|
||
}
|
||
|
||
// Автовиявлення наповнює inv.interfaces — і сервер зобов'язаний одразу
|
||
// завести для них snmp.if-чек. Без цього кроку на мапі є лінки й немає
|
||
// трафіку: інтерфейси відомі, але їх ніхто не опитує.
|
||
func TestDiscoveryCreatesInterfaceChecks(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
report := func(ifs []*npv1.InterfaceRecord) {
|
||
t.Helper()
|
||
if _, err := f.client.ReportDiscovery(ctx, &npv1.DiscoveryReport{
|
||
AgentId: f.agentID, Final: true, Interfaces: ifs,
|
||
}); err != nil {
|
||
t.Fatalf("ReportDiscovery: %v", err)
|
||
}
|
||
}
|
||
|
||
report([]*npv1.InterfaceRecord{
|
||
{DeviceId: f.deviceID, IfIndex: 1, Name: "GigabitEthernet0/1",
|
||
Type: "ethernetCsmacd", SpeedBps: 1_000_000_000, OperStatus: "up", AdminStatus: "up"},
|
||
{DeviceId: f.deviceID, IfIndex: 2, Name: "GigabitEthernet0/2",
|
||
Type: "ethernetCsmacd", SpeedBps: 1_000_000_000, OperStatus: "up", AdminStatus: "up"},
|
||
// Loopback у чек потрапляти не має: графіка не дає, місце в PDU займає.
|
||
{DeviceId: f.deviceID, IfIndex: 99, Name: "Loopback0",
|
||
Type: "softwareLoopback", OperStatus: "up", AdminStatus: "up"},
|
||
})
|
||
|
||
var (
|
||
firstCheckID string
|
||
paramsJSON string
|
||
intervalSec int
|
||
)
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT id::text, params::text, interval_sec
|
||
FROM core.checks
|
||
WHERE device_id = $1 AND check_type = 'snmp.if'
|
||
`, f.deviceID).Scan(&firstCheckID, ¶msJSON, &intervalSec); err != nil {
|
||
t.Fatalf("snmp.if-чек не створено: %v", err)
|
||
}
|
||
if intervalSec != 60 {
|
||
t.Fatalf("interval_sec = %d", intervalSec)
|
||
}
|
||
|
||
var p struct {
|
||
UseHC bool `json:"use_hc_counters"`
|
||
Interfaces []struct {
|
||
IfIndex int64 `json:"if_index"`
|
||
InterfaceID string `json:"interface_id"`
|
||
SpeedBps uint64 `json:"speed_bps"`
|
||
} `json:"interfaces"`
|
||
}
|
||
if err := json.Unmarshal([]byte(paramsJSON), &p); err != nil {
|
||
t.Fatalf("params не розбираються: %v", err)
|
||
}
|
||
if !p.UseHC {
|
||
t.Fatal("use_hc_counters вимкнено — на гігабіті 32-бітні лічильники перевертаються за 34 с")
|
||
}
|
||
// Порт з фікстури (ifIndex 1) уже існував, два нові додались,
|
||
// loopback відсіяно.
|
||
if len(p.Interfaces) != 2 {
|
||
t.Fatalf("портів у чеку: %d, очікували 2 (без loopback): %s", len(p.Interfaces), paramsJSON)
|
||
}
|
||
for _, ifc := range p.Interfaces {
|
||
if ifc.InterfaceID == "" {
|
||
t.Fatalf("порт без interface_id: %+v", ifc)
|
||
}
|
||
if ifc.SpeedBps == 0 {
|
||
t.Fatalf("порт без speed_bps — не буде знаменника для util_pct: %+v", ifc)
|
||
}
|
||
}
|
||
|
||
// Повторний звіт із тим самим складом портів не має нічого міняти:
|
||
// інакше агент отримував би новий план на кожен обхід.
|
||
report([]*npv1.InterfaceRecord{
|
||
{DeviceId: f.deviceID, IfIndex: 1, Name: "GigabitEthernet0/1",
|
||
Type: "ethernetCsmacd", SpeedBps: 1_000_000_000, OperStatus: "up", AdminStatus: "up"},
|
||
{DeviceId: f.deviceID, IfIndex: 2, Name: "GigabitEthernet0/2",
|
||
Type: "ethernetCsmacd", SpeedBps: 1_000_000_000, OperStatus: "up", AdminStatus: "up"},
|
||
})
|
||
|
||
var sameParams string
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT params::text FROM core.checks WHERE id = $1
|
||
`, firstCheckID).Scan(&sameParams); err != nil {
|
||
t.Fatalf("чек зник: %v", err)
|
||
}
|
||
|
||
// Новий порт — чек має оновитись, а не задвоїтись. Унікальний індекс
|
||
// core.checks включає md5(params), тому наївний upsert плодив би
|
||
// новий рядок на кожну зміну складу портів.
|
||
report([]*npv1.InterfaceRecord{
|
||
{DeviceId: f.deviceID, IfIndex: 3, Name: "TenGigabitEthernet1/1",
|
||
Type: "ethernetCsmacd", SpeedBps: 10_000_000_000, OperStatus: "up", AdminStatus: "up"},
|
||
})
|
||
|
||
var checks int
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT count(*) FROM core.checks WHERE device_id = $1 AND check_type = 'snmp.if'
|
||
`, f.deviceID).Scan(&checks); err != nil {
|
||
t.Fatalf("підрахунок: %v", err)
|
||
}
|
||
if checks != 1 {
|
||
t.Fatalf("snmp.if-чеків: %d — зміна складу портів задвоїла чек", checks)
|
||
}
|
||
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT params::text FROM core.checks WHERE id = $1
|
||
`, firstCheckID).Scan(¶msJSON); err != nil {
|
||
t.Fatalf("чек зник: %v", err)
|
||
}
|
||
if err := json.Unmarshal([]byte(paramsJSON), &p); err != nil {
|
||
t.Fatalf("params: %v", err)
|
||
}
|
||
if len(p.Interfaces) != 3 {
|
||
t.Fatalf("після додавання порту в чеку %d портів", len(p.Interfaces))
|
||
}
|
||
}
|
||
|
||
// Без SNMP-креденшела чек створювати не можна: він лише щохвилини
|
||
// писав би помилку автентифікації.
|
||
func TestNoInterfaceCheckWithoutSnmpCredential(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
var bareDevice string
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
INSERT INTO inv.devices (tenant_id, agent_id, name, address, kind)
|
||
VALUES ($1, $2, 'no-creds', '10.10.0.9', 'switch')
|
||
RETURNING id::text`, f.tenantID, f.agentID).Scan(&bareDevice); err != nil {
|
||
t.Fatalf("seed: %v", err)
|
||
}
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
if _, err := f.client.ReportDiscovery(ctx, &npv1.DiscoveryReport{
|
||
AgentId: f.agentID, Final: true,
|
||
Interfaces: []*npv1.InterfaceRecord{{
|
||
DeviceId: bareDevice, IfIndex: 1, Name: "eth0",
|
||
Type: "ethernetCsmacd", SpeedBps: 1_000_000_000,
|
||
OperStatus: "up", AdminStatus: "up",
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("ReportDiscovery: %v", err)
|
||
}
|
||
|
||
var checks int
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT count(*) FROM core.checks WHERE device_id = $1 AND check_type = 'snmp.if'
|
||
`, bareDevice).Scan(&checks); err != nil {
|
||
t.Fatalf("підрахунок: %v", err)
|
||
}
|
||
if checks != 0 {
|
||
t.Fatalf("створено %d чеків для пристрою без SNMP-креденшела", checks)
|
||
}
|
||
|
||
// А інтерфейси при цьому все одно збережені: вони потрібні мапі
|
||
// незалежно від того, чи є чим їх опитувати.
|
||
var ifaces int
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT count(*) FROM inv.interfaces WHERE device_id = $1
|
||
`, bareDevice).Scan(&ifaces); err != nil {
|
||
t.Fatalf("підрахунок: %v", err)
|
||
}
|
||
if ifaces != 1 {
|
||
t.Fatalf("інтерфейсів збережено: %d", ifaces)
|
||
}
|
||
}
|
||
|
||
// Підтверджений людиною лінк не має перезаписуватись автовиявленням.
|
||
func TestDiscoveryRespectsPinnedLink(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
if _, err := f.pool.Exec(f.ctx, `
|
||
INSERT INTO topo.links
|
||
(tenant_id, a_device_id, a_interface_id, b_device_id, b_interface_id,
|
||
kind, discovered_by, confidence, is_pinned)
|
||
VALUES ($1,$2,$3,$4,$5,'physical','manual',100,true)
|
||
`, f.tenantID, f.deviceID, f.ifaceID, f.peerID, f.peerIfID); err != nil {
|
||
t.Fatalf("seed лінка: %v", err)
|
||
}
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
if _, err := f.client.ReportDiscovery(ctx, &npv1.DiscoveryReport{
|
||
AgentId: f.agentID, Final: true,
|
||
Neighbors: []*npv1.NeighborRecord{{
|
||
DeviceId: f.deviceID,
|
||
LocalInterfaceId: f.ifaceID,
|
||
Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_CDP,
|
||
RemoteChassisId: "0011.2233.6677",
|
||
RemotePortId: "ether1",
|
||
SeenAt: timestamppb.Now(),
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("ReportDiscovery: %v", err)
|
||
}
|
||
|
||
var discoveredBy string
|
||
var pinned bool
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT discovered_by::text, is_pinned FROM topo.links WHERE tenant_id = $1
|
||
`, f.tenantID).Scan(&discoveredBy, &pinned); err != nil {
|
||
t.Fatalf("лінк: %v", err)
|
||
}
|
||
if !pinned || discoveredBy != "manual" {
|
||
t.Fatalf("автовиявлення затерло ручний лінк: discovered_by=%q pinned=%v",
|
||
discoveredBy, pinned)
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 4. NCM: збір конфігу та дедуплікація
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestConfigUploadAndDedup(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
body := []byte("!\nhostname core-sw-01\ninterface Gi0/1\n description uplink\n!\n")
|
||
sum := sha256.Sum256(body)
|
||
|
||
upload := func(jobID string) *npv1.ConfigReceipt {
|
||
t.Helper()
|
||
stream, err := f.client.UploadConfig(ctx)
|
||
if err != nil {
|
||
t.Fatalf("UploadConfig: %v", err)
|
||
}
|
||
if err := stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Header{
|
||
Header: &npv1.ConfigHeader{
|
||
JobId: jobID, AgentId: f.agentID, DeviceId: f.deviceID,
|
||
ConfigType: "running", CollectedAt: timestamppb.Now(), Encoding: "none",
|
||
}}}); err != nil {
|
||
t.Fatalf("header: %v", err)
|
||
}
|
||
// Два чанки, щоб перевірити склеювання.
|
||
mid := len(body) / 2
|
||
for i, part := range [][]byte{body[:mid], body[mid:]} {
|
||
if err := stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Chunk{
|
||
Chunk: &npv1.ConfigChunk{Sequence: uint32(i), Data: part},
|
||
}}); err != nil {
|
||
t.Fatalf("chunk %d: %v", i, err)
|
||
}
|
||
}
|
||
if err := stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Trailer{
|
||
Trailer: &npv1.ConfigTrailer{
|
||
Success: true, ContentSha256: sum[:],
|
||
SizeBytes: uint64(len(body)), ChunkCount: 2,
|
||
}}}); err != nil {
|
||
t.Fatalf("trailer: %v", err)
|
||
}
|
||
receipt, err := stream.CloseAndRecv()
|
||
if err != nil {
|
||
t.Fatalf("CloseAndRecv: %v", err)
|
||
}
|
||
return receipt
|
||
}
|
||
|
||
first := upload("job-1")
|
||
if !first.Accepted || first.Unchanged {
|
||
t.Fatalf("перший конфіг: accepted=%v unchanged=%v err=%+v",
|
||
first.Accepted, first.Unchanged, first.Error)
|
||
}
|
||
if first.ConfigId == "" {
|
||
t.Fatal("немає config_id")
|
||
}
|
||
|
||
// Другий збір того самого конфігу — не нова версія.
|
||
second := upload("job-2")
|
||
if !second.Accepted {
|
||
t.Fatalf("другий конфіг відхилено: %+v", second.Error)
|
||
}
|
||
if !second.Unchanged {
|
||
t.Fatal("незмінений конфіг створив нову версію — історія засмітиться")
|
||
}
|
||
|
||
var versions int
|
||
if err := f.pool.QueryRow(f.ctx,
|
||
`SELECT count(*) FROM ncm.configs WHERE device_id = $1`, f.deviceID).Scan(&versions); err != nil {
|
||
t.Fatalf("підрахунок версій: %v", err)
|
||
}
|
||
if versions != 1 {
|
||
t.Fatalf("версій у ncm.configs: %d, очікували 1", versions)
|
||
}
|
||
|
||
// Тіло збережене зашифрованим і читається назад тим самим ключем.
|
||
var keyID, aad string
|
||
var nonce, ct, tag []byte
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT s.key_id, s.nonce, s.ciphertext, s.auth_tag, COALESCE(s.aad,'')
|
||
FROM ncm.configs c JOIN core.secrets s ON s.id = c.body_secret_id
|
||
WHERE c.device_id = $1`, f.deviceID).Scan(&keyID, &nonce, &ct, &tag, &aad); err != nil {
|
||
t.Fatalf("зашифроване тіло: %v", err)
|
||
}
|
||
plain, err := f.ring.Decrypt(&crypto.Secret{
|
||
KeyID: keyID, Nonce: nonce, Ciphertext: ct, AuthTag: tag,
|
||
}, aad)
|
||
if err != nil {
|
||
t.Fatalf("розшифрування тіла: %v", err)
|
||
}
|
||
if string(plain) != string(body) {
|
||
t.Fatal("тіло конфігу спотворилось при збереженні")
|
||
}
|
||
|
||
// І воно точно не лежить відкритим текстом.
|
||
if string(ct) == string(body) {
|
||
t.Fatal("конфіг збережено без шифрування")
|
||
}
|
||
}
|
||
|
||
func TestConfigUploadRejectsBadChecksum(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := f.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-x", DeviceId: f.deviceID, 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)
|
||
}
|
||
|
||
var versions int
|
||
_ = f.pool.QueryRow(f.ctx,
|
||
`SELECT count(*) FROM ncm.configs WHERE device_id = $1`, f.deviceID).Scan(&versions)
|
||
if versions != 0 {
|
||
t.Fatalf("зіпсований конфіг усе одно збережено: версій %d", versions)
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 5. Heartbeat лягає в самометрики
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestHeartbeatPersisted(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
stream, err := f.client.Control(ctx)
|
||
if err != nil {
|
||
t.Fatalf("Control: %v", err)
|
||
}
|
||
if err := stream.Send(&npv1.ControlUp{
|
||
Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{AgentId: f.agentID}},
|
||
}); err != nil {
|
||
t.Fatalf("Hello: %v", err)
|
||
}
|
||
if _, err := stream.Recv(); err != nil {
|
||
t.Fatalf("Welcome: %v", err)
|
||
}
|
||
|
||
if err := stream.Send(&npv1.ControlUp{
|
||
Payload: &npv1.ControlUp_Heartbeat{Heartbeat: &npv1.Heartbeat{
|
||
Ts: timestamppb.Now(),
|
||
Health: &npv1.AgentHealth{
|
||
CpuPct: 2.5, RssBytes: 12 << 20, Goroutines: 48,
|
||
QueueDepth: 3, DroppedSamples: 7,
|
||
},
|
||
TasksRunning: 1,
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("Heartbeat: %v", err)
|
||
}
|
||
|
||
// Дочекаємось Ping у відповідь — він приходить після запису.
|
||
deadline := time.Now().Add(10 * time.Second)
|
||
for time.Now().Before(deadline) {
|
||
msg, err := stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Recv: %v", err)
|
||
}
|
||
if msg.GetPing() != nil {
|
||
break
|
||
}
|
||
}
|
||
|
||
var rss int64
|
||
if err := f.pool.QueryRow(f.ctx, `
|
||
SELECT rss_bytes FROM ts.agent_health WHERE agent_id = $1
|
||
ORDER BY ts DESC LIMIT 1`, f.agentID).Scan(&rss); err != nil {
|
||
t.Fatalf("самометрики не записались: %v", err)
|
||
}
|
||
if rss != 12<<20 {
|
||
t.Fatalf("rss_bytes = %d", rss)
|
||
}
|
||
|
||
// Втрачені семпли мають бути видимі оператору у зведенні зонда,
|
||
// а не тільки в графіку.
|
||
var health []byte
|
||
if err := f.pool.QueryRow(f.ctx,
|
||
`SELECT health::text FROM core.agents WHERE id = $1`, f.agentID).Scan(&health); err != nil {
|
||
t.Fatalf("зведення: %v", err)
|
||
}
|
||
var h map[string]any
|
||
if err := json.Unmarshal(health, &h); err != nil {
|
||
t.Fatalf("health не JSON: %v", err)
|
||
}
|
||
if h["dropped_samples"] == nil || fmt.Sprint(h["dropped_samples"]) != "7" {
|
||
t.Fatalf("dropped_samples не потрапило у зведення: %v", h)
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// 6. Хеш плану дозволяє не переливати задачі після реконекту
|
||
// ---------------------------------------------------------------------
|
||
|
||
func TestPlanHashSkipsResend(t *testing.T) {
|
||
f := setup(t)
|
||
|
||
ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second)
|
||
defer cancel()
|
||
|
||
// Спершу даємо серверу завести автоматичні чеки, і лише потім
|
||
// рахуємо хеш.
|
||
//
|
||
// Інакше тест перевіряє не те, що збирався. Розпізнавання заводиться
|
||
// при підключенні зонда, тобто МІЖ нашим BuildPlan і Hello: хеш, з
|
||
// яким ми прийшли, застаріває дорогою, сервер чесно вирішує
|
||
// переслати план — і тест звинувачує його в тому, що зробив сам.
|
||
//
|
||
// Виклик ідемпотентний: другий раз чек не заводиться, тож після
|
||
// нього хеш уже стабільний — саме та властивість, яку тест і
|
||
// перевіряє.
|
||
if _, err := f.store.EnsureIdentifyChecks(f.ctx,
|
||
&store.Agent{ID: f.agentID, TenantID: f.tenantID}); err != nil {
|
||
t.Fatalf("EnsureIdentifyChecks: %v", err)
|
||
}
|
||
|
||
// Перше підключення: дізнаємось хеш.
|
||
plan, err := f.store.BuildPlan(f.ctx, &store.Agent{ID: f.agentID, TenantID: f.tenantID})
|
||
if err != nil {
|
||
t.Fatalf("BuildPlan: %v", err)
|
||
}
|
||
if len(plan.PlanHash) == 0 {
|
||
t.Fatal("план без хеша")
|
||
}
|
||
|
||
stream, err := f.client.Control(ctx)
|
||
if err != nil {
|
||
t.Fatalf("Control: %v", err)
|
||
}
|
||
if err := stream.Send(&npv1.ControlUp{
|
||
Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{
|
||
AgentId: f.agentID, TaskPlanHash: plan.PlanHash,
|
||
}},
|
||
}); err != nil {
|
||
t.Fatalf("Hello: %v", err)
|
||
}
|
||
|
||
down, err := stream.Recv()
|
||
if err != nil {
|
||
t.Fatalf("Welcome: %v", err)
|
||
}
|
||
if down.GetWelcome().GetTaskPlanFollows() {
|
||
t.Fatal("сервер збирається слати план, хоча хеш збігся")
|
||
}
|
||
|
||
// Хеш детермінований: та сама база — той самий хеш.
|
||
again, err := f.store.BuildPlan(f.ctx, &store.Agent{ID: f.agentID, TenantID: f.tenantID})
|
||
if err != nil {
|
||
t.Fatalf("BuildPlan: %v", err)
|
||
}
|
||
if hex.EncodeToString(again.PlanHash) != hex.EncodeToString(plan.PlanHash) {
|
||
t.Fatal("хеш плану не детермінований")
|
||
}
|
||
}
|