From aaf067a46b921faea40e9743e9f43e93d1f9e674 Mon Sep 17 00:00:00 2001 From: zotac Date: Fri, 14 Aug 2026 15:10:02 +0300 Subject: [PATCH] =?UTF-8?q?=D0=95=D1=82=D0=B0=D0=BF=202:=20=D1=81=D0=B5?= =?UTF-8?q?=D1=80=D0=B2=D0=B5=D1=80=20=D1=81=D0=B0=D0=BC=20=D0=B7=D0=B0?= =?UTF-8?q?=D0=B2=D0=BE=D0=B4=D0=B8=D1=82=D1=8C=20snmp.if-=D1=87=D0=B5?= =?UTF-8?q?=D0=BA=D0=B8=20=D0=B7=20=D0=B2=D0=B8=D1=8F=D0=B2=D0=BB=D0=B5?= =?UTF-8?q?=D0=BD=D0=B8=D1=85=20=D1=96=D0=BD=D1=82=D0=B5=D1=80=D1=84=D0=B5?= =?UTF-8?q?=D0=B9=D1=81=D1=96=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Автовиявлення наповнювало inv.interfaces, але їх ніхто не опитував: на мапі були лінки й не було трафіку. Тепер сервер формує snmp.if-чек зі складу портів і штовхає його живій сесії як TaskDelta — без цього після кожного нового комутатора були б години порожніх графіків. Чек оновлюється, а не задвоюється: унікальний індекс core.checks включає md5(params). Склад портів порівнюється як множина, бо порядок ключів у jsonb не гарантований. Без SNMP-креденшела чек не створюється. Живий прогін знайшов помилку: OID у запиті йшов без провідної крапки, а pdu.Name повертається з нею — пошук у мапі мовчки не знаходив нічого, і чек виглядав як "жоден інтерфейс не відповів". Канонізація тепер у snmpx.Normalize, застосована з обох боків. Перевірено наскрізь: виявлення -> автостворення чека -> справжні HC-лічильники -> ts.if_counters, без жодного ручного кроку. Co-Authored-By: Claude Opus 5 --- HISTORY.md | 54 ++++- agent/go.mod | 12 +- agent/internal/modules/snmp/snmp.go | 31 ++- agent/internal/snmpx/client.go | 15 +- agent/internal/snmpx/oid.go | 13 ++ server/README.md | 31 +++ server/go.mod | 16 +- server/internal/grpcapi/integration_test.go | 169 ++++++++++++++ server/internal/grpcapi/streams.go | 66 ++++++ server/internal/store/autochecks.go | 230 ++++++++++++++++++++ server/internal/store/plan.go | 36 +++ test/contract/go.mod | 11 +- 12 files changed, 651 insertions(+), 33 deletions(-) create mode 100644 server/internal/store/autochecks.go diff --git a/HISTORY.md b/HISTORY.md index 6d5912e..351258c 100644 --- a/HISTORY.md +++ b/HISTORY.md @@ -326,10 +326,56 @@ server/cmd/netpulse-secret/ заведення шифрованих с - **`snmp.if`**: сервер поки не генерує для нього перелік інтерфейсів, тому чек нікуди не призначається. +--- + +## 2026-08-14 — Етап 2 (частина 5): автостворення snmp.if-чеків + +### Створено + +`server/internal/store/autochecks.go` — з виявлених інтерфейсів формується +`snmp.if`-чек і одразу штовхається живій сесії як `TaskDelta`. + +### Прийняті рішення + +1. **Чек оновлюється, а не задвоюється.** Унікальний індекс `core.checks` + включає `md5(params)`, тому наївний upsert плодив би новий рядок на кожну + зміну складу портів. Шукаємо існуючий чек за `(device_id, 'snmp.if')`. +2. **Склад портів порівнюється як множина.** Postgres не гарантує порядок + ключів у `jsonb` — пряме порівняння рядків давало б хибну зміну на кожному + обході, і агент отримував би новий план щоразу. +3. **Без SNMP-креденшела чек не створюється:** він лише щохвилини писав би + помилку автентифікації. Інтерфейси при цьому все одно зберігаються — вони + потрібні мапі незалежно від того, чи є чим їх опитувати. +4. **Loopback і `notPresent` відсіюються** — графіка не дають, місце в PDU + займають. Ліміт 256 портів на чек: далі опитування не вкладається у власний + таймаут. +5. **Зміна штовхається живому зонду.** Чекати наступного перепідключення — це + години порожніх графіків після кожного нового комутатора. + +### Знайдено живим прогоном + +**OID у запиті без провідної крапки, а `pdu.Name` — з нею.** Пошук у мапі +результатів мовчки нічого не знаходив, і чек виглядав як «жоден інтерфейс не +відповів». `snmp.get` працював, бо там уже була нормалізація, а `snmp.if` — ні. +Канонізацію винесено в `snmpx.Normalize` і застосовано з обох боків у `GetUints`. + +### Перевірено наскрізь проти справжнього SNMP + +Агент і сервер запущені як є, без жодного ручного кроку між ними: + +- `topology.discover` знайшов `lo` і `eth0` → `inv.interfaces`; +- сервер створив `snmp.if`-чек на 1 порт (loopback відсіяно), 10 Гбіт/с; +- дельта доїхала до живої сесії (`надіслано_наживо: true`); +- агент опитав справжні HC-лічильники: `in_octets` 625 246 266, `in_bps` 11 938; +- `util_out_pct` 0.000001 % — знаменник із `ifHighSpeed`; +- помилок чеків немає. + +Це повний шлях даних для анімації трафіку на мапі. Усі три набори тестів +проходять з `-race`. + ### Далі -- Сервер має сам створювати `snmp.if`-чеки з `inv.interfaces` після - автовиявлення — тоді запрацює анімація трафіку на мапі. -- `EnrollmentService` + видача сертифікатів. +- `EnrollmentService` + видача сертифікатів зондам. - Планувальник NCM: бекап за cron і за Syslog-подією. -- REST/WebSocket API для UI поверх тієї ж БД. +- REST/WebSocket API для UI: у базі вже є все для мапи — вузли, лінки, + статуси й завантаження каналів. diff --git a/agent/go.mod b/agent/go.mod index 37ae71a..f624d44 100644 --- a/agent/go.mod +++ b/agent/go.mod @@ -1,14 +1,20 @@ module github.com/netpulse/netpulse/agent -go 1.24 +go 1.25.0 require ( github.com/gosnmp/gosnmp v1.42.1 github.com/netpulse/netpulse/gen/go v0.0.0 - golang.org/x/net v0.46.0 - google.golang.org/grpc v1.76.0 + golang.org/x/net v0.55.0 + google.golang.org/grpc v1.83.0 google.golang.org/protobuf v1.36.12 ) +require ( + golang.org/x/sys v0.45.0 // indirect + golang.org/x/text v0.37.0 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect +) + // Згенерований із proto код лежить у репозиторії поруч. replace github.com/netpulse/netpulse/gen/go => ../gen/go diff --git a/agent/internal/modules/snmp/snmp.go b/agent/internal/modules/snmp/snmp.go index 83673b0..bdbc6ae 100644 --- a/agent/internal/modules/snmp/snmp.go +++ b/agent/internal/modules/snmp/snmp.go @@ -30,18 +30,18 @@ import ( // Базові OID. HC — 64-бітні лічильники з ifXTable; на гігабіті // 32-бітні перевертаються за 34 секунди, тому вони лише запасний варіант. const ( - oidIfInOctets = "1.3.6.1.2.1.2.2.1.10" - oidIfOutOctets = "1.3.6.1.2.1.2.2.1.16" - oidIfHCInOctets = "1.3.6.1.2.1.31.1.1.1.6" - oidIfHCOutOctets = "1.3.6.1.2.1.31.1.1.1.10" - oidIfHCInUcast = "1.3.6.1.2.1.31.1.1.1.7" - oidIfHCOutUcast = "1.3.6.1.2.1.31.1.1.1.11" - oidIfInErrors = "1.3.6.1.2.1.2.2.1.14" - oidIfOutErrors = "1.3.6.1.2.1.2.2.1.20" - oidIfInDiscards = "1.3.6.1.2.1.2.2.1.13" - oidIfOutDiscards = "1.3.6.1.2.1.2.2.1.19" - oidIfAdminStatus = "1.3.6.1.2.1.2.2.1.7" - oidIfOperStatus = "1.3.6.1.2.1.2.2.1.8" + oidIfInOctets = ".1.3.6.1.2.1.2.2.1.10" + oidIfOutOctets = ".1.3.6.1.2.1.2.2.1.16" + oidIfHCInOctets = ".1.3.6.1.2.1.31.1.1.1.6" + oidIfHCOutOctets = ".1.3.6.1.2.1.31.1.1.1.10" + oidIfHCInUcast = ".1.3.6.1.2.1.31.1.1.1.7" + oidIfHCOutUcast = ".1.3.6.1.2.1.31.1.1.1.11" + oidIfInErrors = ".1.3.6.1.2.1.2.2.1.14" + oidIfOutErrors = ".1.3.6.1.2.1.2.2.1.20" + oidIfInDiscards = ".1.3.6.1.2.1.2.2.1.13" + oidIfOutDiscards = ".1.3.6.1.2.1.2.2.1.19" + oidIfAdminStatus = ".1.3.6.1.2.1.2.2.1.7" + oidIfOperStatus = ".1.3.6.1.2.1.2.2.1.8" ) // --------------------------------------------------------------------- @@ -339,12 +339,7 @@ func (m *Module) runGet(ctx context.Context, client *gosnmp.GoSNMP, task module. return out, nil } -func normalizeOID(oid string) string { - if oid == "" || oid[0] == '.' { - return oid - } - return "." + oid -} +func normalizeOID(oid string) string { return snmpx.Normalize(oid) } // --------------------------------------------------------------------- // Транспорт diff --git a/agent/internal/snmpx/client.go b/agent/internal/snmpx/client.go index 73fc4e2..fbad9cc 100644 --- a/agent/internal/snmpx/client.go +++ b/agent/internal/snmpx/client.go @@ -117,19 +117,26 @@ func PickCredential(creds []*npv1.Credential) *npv1.Credential { func GetUints(ctx context.Context, c *gosnmp.GoSNMP, oids []string) (map[string]uint64, error) { out := make(map[string]uint64, len(oids)) - for start := 0; start < len(oids); start += MaxVarsPerPDU { + // Ключі канонізуємо з обох боків: інакше запит без провідної + // крапки ніколи не знайде свою відповідь. + req := make([]string, len(oids)) + for i, o := range oids { + req[i] = Normalize(o) + } + + for start := 0; start < len(req); start += MaxVarsPerPDU { if err := ctx.Err(); err != nil { return nil, err } - end := min(start+MaxVarsPerPDU, len(oids)) + end := min(start+MaxVarsPerPDU, len(req)) - pkt, err := c.Get(oids[start:end]) + pkt, err := c.Get(req[start:end]) if err != nil { return nil, fmt.Errorf("snmp get: %w", err) } for _, pdu := range pkt.Variables { if v, ok := AsUint(pdu); ok { - out[pdu.Name] = v + out[Normalize(pdu.Name)] = v } } } diff --git a/agent/internal/snmpx/oid.go b/agent/internal/snmpx/oid.go index efe14dc..d87df4a 100644 --- a/agent/internal/snmpx/oid.go +++ b/agent/internal/snmpx/oid.go @@ -7,6 +7,19 @@ import ( "strings" ) +// Normalize зводить OID до канонічного вигляду з провідною крапкою. +// +// Потрібне тому, що gosnmp повертає pdu.Name ЗАВЖДИ з крапкою, а +// запитувати дозволяє і без неї. Без канонізації пошук у мапі +// результатів мовчки не знаходить нічого, і чек виглядає як «пристрій +// не відповів» — саме так це й проявилось на живому стенді. +func Normalize(oid string) string { + if oid == "" || oid[0] == '.' { + return oid + } + return "." + oid +} + // Index — суфікс OID після базового піддерева, розібраний на числа. // // Уся робота з SNMP-таблицями зводиться до одного: узяти OID змінної, diff --git a/server/README.md b/server/README.md index 75bb02f..9fe7203 100644 --- a/server/README.md +++ b/server/README.md @@ -23,6 +23,7 @@ internal/ credentials.go розшифровка секретів і видача комплекту з TTL telemetry.go резолвер серій + запис у гіпертаблиці discovery.go сусіди → зіставлення з інвентарем → topo.links + autochecks.go snmp.if-чеки з виявлених інтерфейсів ncm.go прийом конфігів, дедуплікація, шифроване тіло grpcapi/ AgentService: Control, StreamTelemetry, StreamLogs, ReportDiscovery, UploadConfig + перехоплювачі автентифікації @@ -69,6 +70,17 @@ internal/ read-modify-write: два воркери, що обробляють сусідні батчі, наввипередки писали б неіснуючі переходи `up→up`. +**Виявлені інтерфейси одразу отримують snmp.if-чек.** Без цього кроку +автовиявлення наповнює `inv.interfaces`, але ніхто їх не опитує: на мапі є лінки й +немає трафіку. Чекати наступного перепідключення агента, щоб він забрав новий +план, — це години порожніх графіків після кожного нового комутатора, тому зміна +штовхається живій сесії як `TaskDelta`. Чек не створюється для пристрою без +SNMP-креденшела: він лише щохвилини писав би помилку автентифікації. + +**Склад портів порівнюється як множина, а не як рядок JSON.** Postgres не +гарантує порядок ключів у `jsonb`, тому пряме порівняння давало б хибну зміну на +кожному обході автовиявлення — і агент отримував би новий план щоразу. + **Автовиявлення інтерпретує сервер.** Агент доповідає лише «на порту X бачу chassis Y». Впевненість зіставлення спадає за надійністю ознаки: chassis-id (95) → MAC (90) → IP керування (80) → sysName (60). Останнє низьке навмисно: `sysName` вводить людина, і @@ -107,6 +119,8 @@ gRPC (`NETPULSE_TEST_DSN`; без змінної пропускаються): | `TestConfigUploadRejectsBadChecksum` | зіпсований конфіг не потрапляє в базу | | `TestHeartbeatPersisted` | самометрики в `ts.agent_health`, `dropped_samples` видно у зведенні зонда | | `TestPlanHashSkipsResend` | збіг хеша → сервер не шле план | +| `TestDiscoveryCreatesInterfaceChecks` | виявлені порти дають `snmp.if`-чек із `interface_id` і `speed_bps`; loopback відсіяно; зміна складу портів **оновлює** чек, а не задвоює | +| `TestNoInterfaceCheckWithoutSnmpCredential` | без SNMP-креденшела чек не створюється, але інтерфейси все одно збережені | ### Живий наскрізний прогін @@ -124,6 +138,23 @@ gRPC (`NETPULSE_TEST_DSN`; без змінної пропускаються): | Розрив сесії | зонд позначений `offline` | | Розсинхронізація годинника | −0.9 мс | +### Наскрізний ланцюг проти справжнього SNMP + +Стенд із `snmpd` + `lldpd`. Агент і сервер запущені як є, без жодного ручного +кроку між ними: + +| Крок | Результат | +|------|-----------| +| `topology.discover` знайшов порти | `lo`, `eth0` у `inv.interfaces` | +| Сервер створив `snmp.if`-чек | 1 порт (loopback відсіяно), 10 Гбіт/с | +| Дельта доїхала до живої сесії | `надіслано_наживо: true` | +| Агент опитав справжні HC-лічильники | `in_octets` 625 246 266, `in_bps` 11 938 | +| `util_out_pct` пораховано | 0.000001 % — знаменник 10 Гбіт/с із `ifHighSpeed` | +| Помилок чеків | немає | + +Це повний шлях даних для анімації трафіку на мапі: від виявлення порту до +`ts.if_counters`, без жодного ручного налаштування. + ### Знайдено під час перевірки Приведення типу на місці (`$1::text`) **не розв'язує** конфлікт виведення типів у diff --git a/server/go.mod b/server/go.mod index fd922d1..6cd95f0 100644 --- a/server/go.mod +++ b/server/go.mod @@ -1,12 +1,24 @@ module github.com/netpulse/netpulse/server -go 1.24 +go 1.25.0 require ( github.com/jackc/pgx/v5 v5.7.6 github.com/netpulse/netpulse/gen/go v0.0.0 - google.golang.org/grpc v1.76.0 + google.golang.org/grpc v1.83.0 google.golang.org/protobuf v1.36.12 ) +require ( + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect + golang.org/x/crypto v0.51.0 // indirect + golang.org/x/net v0.55.0 // indirect + golang.org/x/sync v0.20.0 // indirect + golang.org/x/sys v0.45.0 // indirect + golang.org/x/text v0.37.0 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect +) + replace github.com/netpulse/netpulse/gen/go => ../gen/go diff --git a/server/internal/grpcapi/integration_test.go b/server/internal/grpcapi/integration_test.go index ce5b38e..d69959f 100644 --- a/server/internal/grpcapi/integration_test.go +++ b/server/internal/grpcapi/integration_test.go @@ -650,6 +650,175 @@ func TestDiscoveryResolvesLink(t *testing.T) { } } +// Автовиявлення наповнює 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) diff --git a/server/internal/grpcapi/streams.go b/server/internal/grpcapi/streams.go index ceb152b..65204aa 100644 --- a/server/internal/grpcapi/streams.go +++ b/server/internal/grpcapi/streams.go @@ -163,6 +163,8 @@ func (s *Service) ReportDiscovery(ctx context.Context, rep *npv1.DiscoveryReport "сусідів", st.NeighborsSeen, "зіставлено", st.NeighborsResolved, "лінків_створено", st.LinksCreated, "лінків_оновлено", st.LinksUpdated) + s.refreshInterfaceChecks(ctx, agent, rep) + return &npv1.DiscoveryAck{ Accepted: true, NeighborsResolved: uint32(st.NeighborsResolved), @@ -170,6 +172,70 @@ func (s *Service) ReportDiscovery(ctx context.Context, rep *npv1.DiscoveryReport }, nil } +// refreshInterfaceChecks створює або оновлює snmp.if-чеки за щойно +// виявленими інтерфейсами й одразу штовхає зміну живому зонду. +// +// Без цього кроку автовиявлення наповнює inv.interfaces, але ніхто їх +// не опитує: на мапі є лінки й немає трафіку. Чекати наступного +// перепідключення агента, щоб він забрав новий план, — це години +// порожніх графіків після кожного нового комутатора. +func (s *Service) refreshInterfaceChecks(ctx context.Context, agent *store.Agent, rep *npv1.DiscoveryReport) { + seen := make(map[string]bool) + for _, r := range rep.GetInterfaces() { + if id := r.GetDeviceId(); id != "" { + seen[id] = true + } + } + if len(seen) == 0 { + return + } + + var ( + tasks []*npv1.Task + deviceIDs []string + ) + for deviceID := range seen { + task, err := s.store.EnsureInterfaceChecks(ctx, agent, deviceID) + if err != nil { + // Обрізаний список портів — теж помилка, але чек при цьому + // створено, тому це попередження, а не привід зупинитись. + s.log.Warn("snmp.if-чек", "device", deviceID, "err", err) + } + if task != nil { + tasks = append(tasks, task) + deviceIDs = append(deviceIDs, deviceID) + } + } + if len(tasks) == 0 { + return + } + + // Новий хеш обов'язковий: інакше після реконекту агент доповість + // старий, сервер вирішить, що план застарів, і перезаллє все. + hash, err := s.store.PlanHash(ctx, agent) + if err != nil { + s.log.Warn("не вдалося перерахувати хеш плану", "agent", agent.ID, "err", err) + return + } + + devices, err := s.store.DeviceTargets(ctx, agent, deviceIDs) + if err != nil { + s.log.Warn("не вдалося зібрати описи пристроїв", "agent", agent.ID, "err", err) + return + } + + delivered := s.PushToAgent(agent.ID, &npv1.ControlDown{ + Payload: &npv1.ControlDown_TaskDelta{TaskDelta: &npv1.TaskDelta{ + PlanHash: hash, + Upsert: tasks, + UpsertDevices: devices, + }}, + }) + + s.log.Info("snmp.if-чеки оновлено", + "agent", agent.ID, "задач", len(tasks), "надіслано_наживо", delivered) +} + // --------------------------------------------------------------------- // Конфігурації // --------------------------------------------------------------------- diff --git a/server/internal/store/autochecks.go b/server/internal/store/autochecks.go new file mode 100644 index 0000000..a3787b2 --- /dev/null +++ b/server/internal/store/autochecks.go @@ -0,0 +1,230 @@ +package store + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "time" + + "github.com/jackc/pgx/v5" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" + "google.golang.org/protobuf/types/known/durationpb" +) + +// Скільки інтерфейсів максимум класти в один snmp.if-чек. +// +// Модуль б'є запити порціями по 24 змінні, а на кожен порт припадає 10 +// OID. 256 портів — це вже 107 PDU за один цикл опитування; далі чек +// перестає вкладатись у власний таймаут раніше, ніж у ліміти пристрою. +// Комутатори з більшою кількістю портів треба ділити на кілька чеків — +// поки що просто обрізаємо й пишемо про це в журнал. +const MaxInterfacesPerCheck = 256 + +// InterfaceCheckInterval — типовий інтервал опитування лічильників. +// +// 60 секунд — компроміс: частіше не має сенсу для 32-бітних лічильників +// на повільних каналах, рідше — анімація трафіку на мапі стає слайдшоу. +const InterfaceCheckInterval = 60 * time.Second + +// ifCheckParams — те, що лягає в core.checks.params для snmp.if. +type ifCheckParams struct { + UseHCCounters bool `json:"use_hc_counters"` + Interfaces []ifCheckTarget `json:"interfaces"` +} + +type ifCheckTarget struct { + IfIndex int64 `json:"if_index"` + InterfaceID string `json:"interface_id"` + SpeedBps uint64 `json:"speed_bps"` +} + +// EnsureInterfaceChecks створює або оновлює snmp.if-чек для пристрою за +// поточним вмістом inv.interfaces. +// +// Навіщо це на сервері, а не на агенті: агент не має права вирішувати, +// що опитувати — це впирається в ліміти тарифу й у те, які інтерфейси +// оператор позначив як непотрібні. Агент лише виконує список. +// +// Повертає задачу для TaskDelta, якщо щось змінилось. nil означає +// «нічого робити»: або немає SNMP-креденшела, або немає інтерфейсів, +// або список не змінився з минулого разу. +func (s *Store) EnsureInterfaceChecks(ctx context.Context, a *Agent, deviceID string) (*npv1.Task, error) { + var task *npv1.Task + + err := s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + // Без SNMP-креденшела чек лише щохвилини писав би помилку + // автентифікації — це шум, а не моніторинг. + var hasCred bool + if err := tx.QueryRow(ctx, ` + SELECT EXISTS ( + SELECT 1 FROM inv.device_credentials dc + JOIN inv.credentials c ON c.id = dc.credential_id + WHERE dc.device_id = $1 + AND c.tenant_id = $2 + AND c.proto IN ('snmp_v2c','snmp_v3') + ) + `, deviceID, a.TenantID).Scan(&hasCred); err != nil { + return err + } + if !hasCred { + return nil + } + + rows, err := tx.Query(ctx, ` + SELECT if_index, id::text, COALESCE(speed_bps, 0) + FROM inv.interfaces + WHERE device_id = $1 + AND tenant_id = $2 + AND monitored + AND if_index IS NOT NULL + -- Loopback і відсутні порти графіка не дають, а місце + -- в PDU займають. + AND COALESCE(type, '') <> 'softwareLoopback' + AND oper_status <> 'notPresent' + ORDER BY if_index + LIMIT $3 + `, deviceID, a.TenantID, MaxInterfacesPerCheck+1) + if err != nil { + return err + } + defer rows.Close() + + params := ifCheckParams{UseHCCounters: true} + for rows.Next() { + var t ifCheckTarget + if err := rows.Scan(&t.IfIndex, &t.InterfaceID, &t.SpeedBps); err != nil { + return err + } + params.Interfaces = append(params.Interfaces, t) + } + if err := rows.Err(); err != nil { + return err + } + + truncated := false + if len(params.Interfaces) > MaxInterfacesPerCheck { + params.Interfaces = params.Interfaces[:MaxInterfacesPerCheck] + truncated = true + } + if len(params.Interfaces) == 0 { + return nil + } + + payload, err := json.Marshal(params) + if err != nil { + return err + } + + // Шукаємо існуючий чек будь-яких параметрів: унікальний індекс + // core.checks включає md5(params), тому наївний upsert плодив би + // новий рядок на кожну зміну складу портів. + var ( + checkID string + oldJSON string + interval int32 + timeout int32 + retries int32 + ) + err = tx.QueryRow(ctx, ` + SELECT id::text, params::text, interval_sec, timeout_ms, retries + FROM core.checks + WHERE device_id = $1 AND tenant_id = $2 AND check_type = 'snmp.if' + ORDER BY created_at + LIMIT 1 + `, deviceID, a.TenantID).Scan(&checkID, &oldJSON, &interval, &timeout, &retries) + + switch { + case err == nil: + if sameInterfaceSet(oldJSON, payload) { + return nil + } + if _, err := tx.Exec(ctx, ` + UPDATE core.checks + SET params = $2::jsonb, enabled = true, updated_at = now() + WHERE id = $1 + `, checkID, string(payload)); err != nil { + return err + } + + case errors.Is(err, pgx.ErrNoRows): + interval = int32(InterfaceCheckInterval / time.Second) + timeout, retries = 15000, 1 + if err := tx.QueryRow(ctx, ` + INSERT INTO core.checks + (tenant_id, device_id, check_type, params, interval_sec, timeout_ms, retries) + VALUES ($1, $2, 'snmp.if', $3::jsonb, $4, $5, $6) + RETURNING id::text + `, a.TenantID, deviceID, string(payload), interval, timeout, retries).Scan(&checkID); err != nil { + return err + } + + default: + return err + } + + if truncated { + return fmt.Errorf("пристрій %s має понад %d інтерфейсів: чек обрізано", + deviceID, MaxInterfacesPerCheck) + } + + intervalDur := time.Duration(interval) * time.Second + task = &npv1.Task{ + CheckId: checkID, + DeviceId: deviceID, + CheckType: "snmp.if", + ParamsJson: payload, + Interval: durationpb.New(intervalDur), + Timeout: durationpb.New(time.Duration(timeout) * time.Millisecond), + Retries: uint32(retries), + Enabled: true, + ScheduleOffset: durationpb.New(ScheduleOffset(checkID, intervalDur)), + } + return nil + }) + + return task, err +} + +// sameInterfaceSet порівнює склад портів, ігноруючи порядок ключів у JSON. +// +// Пряме порівняння рядків давало б хибну зміну щоразу, коли Postgres +// інакше впорядкує ключі jsonb, і агент отримував би новий план на +// кожен обхід автовиявлення. +func sameInterfaceSet(oldJSON string, newJSON []byte) bool { + var a, b ifCheckParams + if err := json.Unmarshal([]byte(oldJSON), &a); err != nil { + return false + } + if err := json.Unmarshal(newJSON, &b); err != nil { + return false + } + if a.UseHCCounters != b.UseHCCounters || len(a.Interfaces) != len(b.Interfaces) { + return false + } + + seen := make(map[string]ifCheckTarget, len(a.Interfaces)) + for _, t := range a.Interfaces { + seen[t.InterfaceID] = t + } + for _, t := range b.Interfaces { + prev, ok := seen[t.InterfaceID] + if !ok || prev.IfIndex != t.IfIndex || prev.SpeedBps != t.SpeedBps { + return false + } + } + return true +} + +// PlanHash перераховує хеш плану без побудови самого плану. +// +// Потрібен після зміни чеків: агент має отримати новий хеш разом із +// дельтою, інакше після реконекту він доповість старий, сервер вирішить, +// що план застарів, і перезаллє все повністю. +func (s *Store) PlanHash(ctx context.Context, a *Agent) ([]byte, error) { + plan, err := s.BuildPlan(ctx, a) + if err != nil { + return nil, err + } + return plan.GetPlanHash(), nil +} diff --git a/server/internal/store/plan.go b/server/internal/store/plan.go index 01d833d..e53a4ba 100644 --- a/server/internal/store/plan.go +++ b/server/internal/store/plan.go @@ -122,6 +122,42 @@ func (s *Store) BuildPlan(ctx context.Context, a *Agent) (*npv1.TaskPlan, error) return plan, err } +// DeviceTargets віддає описи пристроїв для TaskDelta. +// +// Дельта може нести задачу на пристрій, якого агент ще не знає (його +// щойно завели або він уперше потрапив під цей зонд). Задача без опису +// пристрою — це задача без адреси, тому їх шлють разом. +func (s *Store) DeviceTargets(ctx context.Context, a *Agent, ids []string) ([]*npv1.DeviceTarget, error) { + if len(ids) == 0 { + return nil, nil + } + + var out []*npv1.DeviceTarget + err := s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + rows, err := tx.Query(ctx, ` + SELECT id::text, name, + COALESCE(host(address), COALESCE(fqdn, '')), + COALESCE(chassis_id, ''), COALESCE(system_name, '') + FROM inv.devices + WHERE tenant_id = $1 AND id = ANY($2::uuid[]) AND deleted_at IS NULL + `, a.TenantID, ids) + if err != nil { + return err + } + defer rows.Close() + + for rows.Next() { + d := &npv1.DeviceTarget{} + if err := rows.Scan(&d.DeviceId, &d.Name, &d.Address, &d.ChassisId, &d.SystemName); err != nil { + return err + } + out = append(out, d) + } + return rows.Err() + }) + return out, err +} + // ModulesForPlan визначає, які модулі треба активувати на зонді. // // Виводиться з типів чеків у плані плюс явно дозволених у diff --git a/test/contract/go.mod b/test/contract/go.mod index 45b75bb..0923744 100644 --- a/test/contract/go.mod +++ b/test/contract/go.mod @@ -1,12 +1,19 @@ module github.com/netpulse/netpulse/test/contract -go 1.24 +go 1.25.0 require ( github.com/netpulse/netpulse/gen/go v0.0.0 - google.golang.org/grpc v1.76.0 + google.golang.org/grpc v1.83.0 google.golang.org/protobuf v1.36.12 ) +require ( + golang.org/x/net v0.55.0 // indirect + golang.org/x/sys v0.45.0 // indirect + golang.org/x/text v0.37.0 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect +) + // Згенерований код лежить у репозиторії поруч, а не тягнеться з проксі. replace github.com/netpulse/netpulse/gen/go => ../../gen/go