Netpulse_SasS/server/internal/grpcapi/integration_test.go
byrsapty ed8fc831bf Дві сесії роботи: 0058–0068, розгортання однією командою, тести
Один коміт, а не десяток тематичних, свідомо: теми переплетені в
спільних файлах (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 серпня.
2026-08-27 17:32:49 +03:00

1159 lines
41 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.

// Наскрізний тест серверної сторони: справжній 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, &paramsJSON, &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(&paramsJSON); 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("хеш плану не детермінований")
}
}