Автовиявлення наповнювало 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>
191 lines
6 KiB
Go
191 lines
6 KiB
Go
package store
|
||
|
||
import (
|
||
"context"
|
||
"crypto/sha256"
|
||
"fmt"
|
||
"hash/fnv"
|
||
"time"
|
||
|
||
"github.com/jackc/pgx/v5"
|
||
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
|
||
"google.golang.org/protobuf/types/known/durationpb"
|
||
)
|
||
|
||
// ScheduleOffset розводить задачі всередині інтервалу.
|
||
//
|
||
// Рахує СЕРВЕР, і рахує детерміновано від check_id: інакше після
|
||
// кожного перезапуску агента 5000 чеків з інтервалом 60 с з'їжджали б
|
||
// у нову випадкову фазу, а раз на хвилину мережа отримувала б сплеск.
|
||
// Функція чиста — те саме значення на будь-якому вузлі сервера.
|
||
func ScheduleOffset(checkID string, interval time.Duration) time.Duration {
|
||
if interval <= 0 {
|
||
return 0
|
||
}
|
||
h := fnv.New64a()
|
||
_, _ = h.Write([]byte(checkID))
|
||
return time.Duration(h.Sum64() % uint64(interval))
|
||
}
|
||
|
||
// BuildPlan збирає повний план задач для зонда.
|
||
func (s *Store) BuildPlan(ctx context.Context, a *Agent) (*npv1.TaskPlan, error) {
|
||
var plan *npv1.TaskPlan
|
||
|
||
err := s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error {
|
||
rows, err := tx.Query(ctx, `
|
||
SELECT c.id::text,
|
||
c.device_id::text,
|
||
COALESCE(c.interface_id::text, ''),
|
||
c.check_type,
|
||
c.params::text,
|
||
c.interval_sec,
|
||
c.timeout_ms,
|
||
c.retries,
|
||
d.name,
|
||
COALESCE(host(d.address), COALESCE(d.fqdn, '')),
|
||
COALESCE(d.chassis_id, ''),
|
||
COALESCE(d.system_name, '')
|
||
FROM core.checks c
|
||
JOIN inv.devices d ON d.id = c.device_id
|
||
WHERE c.tenant_id = $1
|
||
AND d.agent_id = $2
|
||
AND c.enabled
|
||
AND d.enabled
|
||
AND d.deleted_at IS NULL
|
||
ORDER BY c.id
|
||
`, a.TenantID, a.ID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
defer rows.Close()
|
||
|
||
p := &npv1.TaskPlan{Final: true}
|
||
devices := make(map[string]*npv1.DeviceTarget)
|
||
hasher := sha256.New()
|
||
|
||
for rows.Next() {
|
||
var (
|
||
checkID, deviceID, ifaceID, checkType, params string
|
||
intervalSec, timeoutMs, retries int32
|
||
devName, devAddr, chassisID, sysName string
|
||
)
|
||
if err := rows.Scan(&checkID, &deviceID, &ifaceID, &checkType, ¶ms,
|
||
&intervalSec, &timeoutMs, &retries,
|
||
&devName, &devAddr, &chassisID, &sysName); err != nil {
|
||
return err
|
||
}
|
||
|
||
interval := time.Duration(intervalSec) * time.Second
|
||
task := &npv1.Task{
|
||
CheckId: checkID,
|
||
DeviceId: deviceID,
|
||
InterfaceId: ifaceID,
|
||
CheckType: checkType,
|
||
ParamsJson: []byte(params),
|
||
Interval: durationpb.New(interval),
|
||
Timeout: durationpb.New(time.Duration(timeoutMs) * time.Millisecond),
|
||
Retries: uint32(retries),
|
||
Enabled: true,
|
||
ScheduleOffset: durationpb.New(ScheduleOffset(checkID, interval)),
|
||
}
|
||
p.Tasks = append(p.Tasks, task)
|
||
|
||
// Хеш плану рахуємо з полів, які реально впливають на
|
||
// поведінку агента. Зміна опису пристрою не має змушувати
|
||
// переливати 50 000 задач.
|
||
fmt.Fprintf(hasher, "%s|%s|%s|%s|%d|%d|%d|%s\n",
|
||
checkID, deviceID, ifaceID, checkType,
|
||
intervalSec, timeoutMs, retries, params)
|
||
|
||
if _, ok := devices[deviceID]; !ok {
|
||
devices[deviceID] = &npv1.DeviceTarget{
|
||
DeviceId: deviceID,
|
||
Name: devName,
|
||
Address: devAddr,
|
||
ChassisId: chassisID,
|
||
SystemName: sysName,
|
||
}
|
||
}
|
||
}
|
||
if err := rows.Err(); err != nil {
|
||
return err
|
||
}
|
||
|
||
for _, d := range devices {
|
||
p.Devices = append(p.Devices, d)
|
||
}
|
||
p.PlanHash = hasher.Sum(nil)
|
||
plan = p
|
||
return nil
|
||
})
|
||
|
||
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 визначає, які модулі треба активувати на зонді.
|
||
//
|
||
// Виводиться з типів чеків у плані плюс явно дозволених у
|
||
// core.agents.enabled_modules: зонд не має вмикати нічого, що йому
|
||
// не знадобиться, — це і пам'ять, і зайва поверхня атаки.
|
||
func ModulesForPlan(plan *npv1.TaskPlan, allowed []string) *npv1.ModuleControl {
|
||
needed := make(map[string]bool)
|
||
for _, t := range plan.GetTasks() {
|
||
if key, _, ok := cutDot(t.GetCheckType()); ok {
|
||
needed[key] = true
|
||
}
|
||
}
|
||
for _, m := range allowed {
|
||
needed[m] = true
|
||
}
|
||
|
||
mc := &npv1.ModuleControl{Exclusive: true}
|
||
for key := range needed {
|
||
mc.Modules = append(mc.Modules, &npv1.ModuleSpec{Key: key, Enabled: true})
|
||
}
|
||
return mc
|
||
}
|
||
|
||
func cutDot(s string) (before, after string, found bool) {
|
||
for i := 0; i < len(s); i++ {
|
||
if s[i] == '.' {
|
||
return s[:i], s[i+1:], true
|
||
}
|
||
}
|
||
return s, "", false
|
||
}
|