Автовиявлення наповнювало 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 <noreply@anthropic.com>
230 lines
8.1 KiB
Go
230 lines
8.1 KiB
Go
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
|
||
}
|