Netpulse_SasS/server/internal/store/plan.go
zotac aaf067a46b Етап 2: сервер сам заводить snmp.if-чеки з виявлених інтерфейсів
Автовиявлення наповнювало 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>
2026-08-14 15:10:02 +03:00

191 lines
6 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.

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, &params,
&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
}