Етап 2: модуль topology — LLDP/CDP/ARP/FDB та інвентар портів

Спільний SNMP-транспорт винесено в snmpx: ним користуються два модулі.
Звіти автовиявлення йдуть окремим RPC, не телеметричним стрімом — вони
рідкі, великі й не прив'язані до моменту часу так, як метрики.

Живий прогін проти справжнього snmpd+lldpd знайшов дві помилки:

1. Префікс типу чека не збігався з ключем модуля (topo.discover при
   плагіні topology). Агент маршрутизує задачі саме за префіксом, тож
   зонд відхиляв би їх. Інваріант закріплено обмеженням у БД.
2. Унікальний індекс topo.neighbors схлопував ARP-сусідів: ключ не
   включав MAC, а chassis_id/port_id в ARP немає взагалі.

Перевірено на живих даних: інвентар із ifTable, 2 ARP-сусіди, зіставлення
шлюзу за MAC (впевненість 90), зведений лінк із capacity 10 Гбіт/с.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
zotac 2026-08-14 14:42:01 +03:00
parent 8a92ef8a45
commit ad7672c32f
15 changed files with 1673 additions and 168 deletions

View file

@ -259,9 +259,77 @@ Postgres, а нав'язує тип обом уживанням параметр
- `EnrollmentService` — зонди заводяться вставкою в `core.agents`.
- Сервер не надсилає `TaskDelta` (лише повний план) і не ініціює `ConfigJob`.
---
## 2026-08-14 — Етап 2 (частина 4): модуль topology, автовиявлення наскрізь
### Створено
```
agent/internal/snmpx/ спільний SNMP-транспорт + розбір індексів OID
agent/internal/modules/topology/ LLDP, CDP, ARP, FDB + інвентар портів
server/cmd/netpulse-secret/ заведення шифрованих секретів із CLI
```
Плюс: `scheduler.OnDiscovery` і `TriggerNow`, доставка звітів у `session`,
обробка `DiscoveryRequest`, реєстрація модуля в `main`.
### Прийняті рішення
1. **SNMP-транспорт винесено в `snmpx`** — ним користуються два модулі, а два
незалежні набори однієї логіки це два незалежні набори багів.
2. **Звіти автовиявлення йдуть окремим RPC**, не телеметричним стрімом: вони
рідкі, великі й не прив'язані до моменту часу так, як метрики.
3. **Черга звітів коротка й витісняє найстаріший.** Знімок топології актуальний
рівно доти, доки описує поточний стан; накопичувати застарілі немає сенсу.
4. **`ifHighSpeed` має пріоритет над `ifSpeed`:** 32-бітне поле впирається в
4.29 Гбіт/с, і на 10G порт показував би неправильний знаменник для `util_pct`.
5. **`ifName` перекриває `ifDescr`:** саме `ifName` віддає LLDP як port-id, тому
саме за ним зійдеться лінк.
6. **`netpulse-secret` як окремий інструмент** — секрети не можна вставити
звичайним SQL, а UI ще немає.
### Знайдено живим прогоном (обидва — справжні помилки)
1. **Префікс типу чека не збігався з ключем модуля.** У сіді було
`topo.discover` при плагіні `topology` і `ssl.expiry` при плагіні `http`.
Агент маршрутизує задачі саме за префіксом, тож зонд відхилив би їх як
адресовані неіснуючому модулю. Перейменовано на `topology.discover` і
`http.ssl_expiry`, а інваріант закріплено обмеженням у БД
`check_types_prefix_matches_plugin` — щоб наступна така неузгодженість
не доїхала до поля.
2. **Унікальний індекс `topo.neighbors` схлопував ARP-сусідів.** Ключ складався
з `chassis_id` + `port_id`, яких в ARP і FDB немає взагалі: з двох сусідів на
одному порту зберігався один. Додано `remote_mac` до індексу.
### Перевірено проти справжнього SNMP-агента
На стенді піднято `snmpd` + `lldpd` (LLDP-MIB через AgentX). Живий прогін
агент → сервер → БД:
- інвентар портів зі справжнього `ifTable`: `lo` і `eth0`, MAC `bc:24:11:07:68:67`,
10 Гбіт/с саме через `ifHighSpeed`;
- 2 сусіди зі справжньої ARP-таблиці, обидва з локальним портом `eth0`;
- шлюз `192.168.1.1` зіставлено за MAC із впевненістю 90;
- лінк `snmp-host:eth0 → gateway` зведено, `capacity_bps` 10 Гбіт/с;
- три проходи поспіль — лінк один, дублікатів немає;
- `snmp.get` збирає `sys.uptime`, помилок чеків немає.
Усі три набори тестів (агент, сервер, контракт) проходять з `-race`.
### Не перевірено
- **LLDP і CDP на живих сусідах.** `lldpd` зареєстрував MIB, але сусідів немає —
поруч немає другого пристрою, що шле LLDP. Розбір `lldpRemTable` і
`cdpCacheTable` покритий лише тестами на індекси OID; ARP-гілка того самого
коду перевірена на живих даних.
- **`snmp.if`**: сервер поки не генерує для нього перелік інтерфейсів, тому чек
нікуди не призначається.
### Далі
- Модуль topology на агенті: LLDP/CDP/ARP/FDB → `topo.neighbors`.
- Сервер має сам створювати `snmp.if`-чеки з `inv.interfaces` після
автовиявлення — тоді запрацює анімація трафіку на мапі.
- `EnrollmentService` + видача сертифікатів.
- Планувальник NCM: бекап за cron і за Syslog-подією.
- REST/WebSocket API для UI поверх тієї ж БД.

View file

@ -23,9 +23,11 @@ internal/
buffer.go обмежений буфер із бюджетом пам'яті
scheduler/ min-heap за часом запуску, семафор паралельності
session/ gRPC-клієнт, рукостискання, реконект, ack
snmpx/ спільний SNMP-транспорт + розбір індексів OID
modules/
icmp/ ping, RTT, jitter, втрати
snmp/ лічильники інтерфейсів (v2c/v3), довільні OID
topology/ LLDP, CDP, ARP, FDB + інвентар портів
```
## Рішення, які варто розуміти
@ -108,6 +110,17 @@ internal/
"metric_key": "cpu.util", "unit": "pct", "scale": 1}]}
```
`topology.discover` — сусіди й інвентар портів беруться з одного SNMP-обходу,
тому окремий чек для інтерфейсів був би зайвим трафіком.
```json
{"protos": ["lldp", "cdp", "arp", "fdb"], "collect_interfaces": true}
```
Префікс типу чека **зобов'язаний** дорівнювати ключу модуля: саме за ним агент
обирає виконавця. Реєстр відхиляє модуль, який оголошує чужий тип, а в БД те саме
закріплено обмеженням `check_types_prefix_matches_plugin`.
## Стан перевірки
Стенд Debian 13, Go 1.25.13. `go vet` чисто, тести з `-race`:
@ -117,6 +130,7 @@ internal/
| `internal/scheduler` | ok |
| `internal/telemetry` | ok |
| `internal/session` | ok |
| `internal/snmpx` | ok |
| `internal/modules/icmp` | ok |
Заміряно на релізному бінарнику (`CGO_ENABLED=0`, `-trimpath -s -w`):
@ -138,14 +152,36 @@ internal/
| `TestBufferRequeuePreservesData` | розрив між Send і Ack не втрачає ні семпли, ні дескриптори |
| `TestSessionFullCycle` | Hello → план у зустрічному напрямку → виконання з креденшелами → телеметрія на сервері |
| `TestSessionReconnectsAfterDrop` | агент повертається сам після розриву |
| `TestSplitIndex` | розбір індексу OID, включно з пасткою «схожий префікс» (`.1.2.20` не є нащадком `.1.2`) |
| `TestFormatMAC` | сирі байти, `00:11:...`, `0011.2233.4455` і дефіси зводяться до одного вигляду |
| `TestPingLoopback` | реальний ICMP-сокет, не мок |
### Перевірка проти справжнього SNMP-агента
На стенді піднято `snmpd` + `lldpd` (LLDP-MIB через AgentX). Живий прогін
`topology.discover` і `snmp.get` кожні 10 с:
| Що перевірено | Результат |
|---------------|-----------|
| Інвентар портів зі справжнього `ifTable` | `lo` (softwareLoopback) і `eth0` (ethernetCsmacd, MAC `bc:24:11:07:68:67`) |
| `ifHighSpeed` замість 32-бітного `ifSpeed` | 10 Гбіт/с — 32-бітне поле впиралося б у 4.29 Гбіт/с |
| Сусіди зі справжньої ARP-таблиці | 2 записи, обидва з локальним портом `eth0` |
| Зіставлення за MAC | шлюз `192.168.1.1` розпізнано, впевненість 90 |
| Зведення лінка | `snmp-host:eth0 → gateway`, джерело `arp`, `capacity_bps` 10 Гбіт/с |
| Три проходи поспіль | лінк один, дублікатів немає |
| Метрики `snmp.get` | `sys.uptime` збирається, помилок чеків немає |
**Обмеження перевірки:** `TestPingLoopback` пропускається під звичайним
користувачем в unprivileged LXC — ядро не дає ані unprivileged-, ані raw-сокета,
і `sysctl net.ipv4.ping_group_range` там недоступний. Під root на тому ж стенді
тест проходить. У проді агенту треба `CAP_NET_RAW` або дозволений
`ping_group_range`.
Модуль `snmp` **не перевірявся проти живого пристрою** — на стенді немає SNMP-агента.
Він компілюється й проходить `vet`, але логіка перевороту лічильників і розрахунку
`util_pct` чекає на перевірку з реальним обладнанням.
**LLDP і CDP не перевірені на живих сусідах:** `lldpd` на стенді зареєстрував
LLDP-MIB, але сусідів у нього немає — поруч немає другого пристрою, який шле LLDP.
Розбір `lldpRemTable` і `cdpCacheTable` покритий лише тестами на індекси OID;
ARP-гілка того самого коду перевірена на живих даних. CDP-гілка чекає на Cisco.
Так само не перевірено `snmp.if`: сервер поки не генерує для нього перелік
інтерфейсів, тому чек нікуди не призначається. Логіка перевороту лічильників і
`util_pct` лишається без живої перевірки.

View file

@ -95,7 +95,11 @@ type Result struct {
Metrics []Metric
Icmp *npv1.IcmpResult
Interfaces []*npv1.InterfaceCounters
Neighbors []*npv1.NeighborRecord
// Сусіди й інвентар портів ідуть на сервер не телеметричним
// стрімом, а окремим ReportDiscovery: вони рідкі, великі й не
// прив'язані до моменту часу так, як метрики.
Neighbors []*npv1.NeighborRecord
InterfaceRecords []*npv1.InterfaceRecord
// Непрозоре навантаження, яке не є метрикою (напр. http.status → тіло відповіді).
Payload []byte
}

View file

@ -16,12 +16,12 @@ import (
"encoding/json"
"fmt"
"math"
"math/big"
"sync"
"time"
"github.com/gosnmp/gosnmp"
"github.com/netpulse/netpulse/agent/internal/module"
"github.com/netpulse/netpulse/agent/internal/snmpx"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
@ -113,7 +113,7 @@ func (m *Module) Close() error {
}
func (m *Module) Run(ctx context.Context, task module.Task) (module.Result, error) {
client, err := dial(ctx, task)
client, err := snmpx.Dial(ctx, task.Target.Address, task.Credentials, task.Timeout)
if err != nil {
return module.Result{}, err
}
@ -348,160 +348,12 @@ func normalizeOID(oid string) string {
// ---------------------------------------------------------------------
// Транспорт
//
// Дзвінки, вибір креденшелів і розбиття запитів на порції живуть у
// пакеті snmpx: ними користується ще й модуль topology, а два
// незалежні набори однієї логіки — два незалежні набори багів.
// ---------------------------------------------------------------------
// getAll розбиває запит на порції: більшість агентів на пристроях
// відмовляють, якщо в одному PDU більше ~30 змінних.
func getAll(ctx context.Context, client *gosnmp.GoSNMP, oids []string) (map[string]uint64, error) {
const chunk = 24
out := make(map[string]uint64, len(oids))
for start := 0; start < len(oids); start += chunk {
if err := ctx.Err(); err != nil {
return nil, err
}
end := start + chunk
if end > len(oids) {
end = len(oids)
}
pkt, err := client.Get(oids[start:end])
if err != nil {
return nil, fmt.Errorf("snmp get: %w", err)
}
for _, pdu := range pkt.Variables {
switch pdu.Type {
case gosnmp.NoSuchObject, gosnmp.NoSuchInstance, gosnmp.EndOfMibView, gosnmp.Null:
continue
}
v := gosnmp.ToBigInt(pdu.Value)
if v == nil || v.Sign() < 0 || v.Cmp(maxUint64) > 0 {
continue
}
out[pdu.Name] = v.Uint64()
}
}
return out, nil
}
var maxUint64 = new(big.Int).SetUint64(math.MaxUint64)
func dial(ctx context.Context, task module.Task) (*gosnmp.GoSNMP, error) {
cred := pickCredential(task.Credentials)
if cred == nil {
return nil, fmt.Errorf("немає SNMP-креденшелів для пристрою %s", task.Target.Name)
}
port := uint16(161)
if cred.Port > 0 && cred.Port < 65536 {
port = uint16(cred.Port)
}
timeout := task.Timeout
if timeout <= 0 {
timeout = 5 * time.Second
}
client := &gosnmp.GoSNMP{
Target: task.Target.Address,
Port: port,
Timeout: timeout,
Retries: 1,
Context: ctx,
}
switch cred.Transport {
case npv1.Transport_TRANSPORT_SNMP_V2C:
client.Version = gosnmp.Version2c
client.Community = cred.GetCommunity()
case npv1.Transport_TRANSPORT_SNMP_V3:
v3 := cred.GetSnmpV3()
if v3 == nil {
return nil, fmt.Errorf("SNMPv3 без параметрів безпеки")
}
client.Version = gosnmp.Version3
client.SecurityModel = gosnmp.UserSecurityModel
client.MsgFlags = msgFlags(v3.Level)
client.ContextName = v3.ContextName
name := v3.SecurityName
if name == "" {
name = cred.Username
}
client.SecurityParameters = &gosnmp.UsmSecurityParameters{
UserName: name,
AuthenticationProtocol: authProto(v3.AuthProtocol),
AuthenticationPassphrase: v3.AuthPassword,
PrivacyProtocol: privProto(v3.PrivProtocol),
PrivacyPassphrase: v3.PrivPassword,
}
default:
return nil, fmt.Errorf("креденшел %s не для SNMP", cred.Transport)
}
if err := client.Connect(); err != nil {
return nil, fmt.Errorf("не вдалося відкрити SNMP-сесію до %s: %w", task.Target.Address, err)
}
return client, nil
}
func pickCredential(creds []*npv1.Credential) *npv1.Credential {
for _, c := range creds {
if c.Transport == npv1.Transport_TRANSPORT_SNMP_V3 {
return c
}
}
for _, c := range creds {
if c.Transport == npv1.Transport_TRANSPORT_SNMP_V2C {
return c
}
}
return nil
}
func msgFlags(level npv1.SnmpV3Options_SecurityLevel) gosnmp.SnmpV3MsgFlags {
switch level {
case npv1.SnmpV3Options_SECURITY_LEVEL_AUTH_PRIV:
return gosnmp.AuthPriv
case npv1.SnmpV3Options_SECURITY_LEVEL_AUTH_NO_PRIV:
return gosnmp.AuthNoPriv
default:
return gosnmp.NoAuthNoPriv
}
}
func authProto(name string) gosnmp.SnmpV3AuthProtocol {
switch name {
case "MD5":
return gosnmp.MD5
case "SHA":
return gosnmp.SHA
case "SHA224":
return gosnmp.SHA224
case "SHA256":
return gosnmp.SHA256
case "SHA384":
return gosnmp.SHA384
case "SHA512":
return gosnmp.SHA512
default:
return gosnmp.NoAuth
}
}
func privProto(name string) gosnmp.SnmpV3PrivProtocol {
switch name {
case "DES":
return gosnmp.DES
case "AES":
return gosnmp.AES
case "AES192":
return gosnmp.AES192
case "AES256":
return gosnmp.AES256
default:
return gosnmp.NoPriv
}
return snmpx.GetUints(ctx, client, oids)
}

View file

@ -0,0 +1,713 @@
// Package topology — автовиявлення сусідів і інвентаризація портів.
//
// Модуль НЕ вирішує, хто з ким з'єднаний. Він доповідає сире:
// «на порту X бачу chassis Y, port Z». Зведення в лінки — на сервері
// (topo.neighbors → topo.links), і саме тому правила зіставлення можна
// міняти без оновлення зондів у полі.
package topology
import (
"context"
"encoding/json"
"fmt"
"strings"
"github.com/gosnmp/gosnmp"
"github.com/netpulse/netpulse/agent/internal/module"
"github.com/netpulse/netpulse/agent/internal/snmpx"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/protobuf/types/known/timestamppb"
)
// LLDP-MIB (1.0.8802.1.1.2)
const (
oidLldpLocPortIDSubtype = ".1.0.8802.1.1.2.1.3.7.1.2"
oidLldpLocPortID = ".1.0.8802.1.1.2.1.3.7.1.3"
oidLldpRemChassisIDSubtype = ".1.0.8802.1.1.2.1.4.1.1.4"
oidLldpRemChassisID = ".1.0.8802.1.1.2.1.4.1.1.5"
oidLldpRemPortIDSubtype = ".1.0.8802.1.1.2.1.4.1.1.6"
oidLldpRemPortID = ".1.0.8802.1.1.2.1.4.1.1.7"
oidLldpRemPortDesc = ".1.0.8802.1.1.2.1.4.1.1.8"
oidLldpRemSysName = ".1.0.8802.1.1.2.1.4.1.1.9"
oidLldpRemSysCapEnabled = ".1.0.8802.1.1.2.1.4.1.1.12"
)
// CISCO-CDP-MIB (cdpCacheTable). Індекс: ifIndex.deviceIndex —
// перший компонент одразу дає локальний порт, на відміну від LLDP.
const (
oidCdpCacheAddress = ".1.3.6.1.4.1.9.9.23.1.2.1.1.4"
oidCdpCacheVersion = ".1.3.6.1.4.1.9.9.23.1.2.1.1.5"
oidCdpCacheDeviceID = ".1.3.6.1.4.1.9.9.23.1.2.1.1.6"
oidCdpCacheDevicePort = ".1.3.6.1.4.1.9.9.23.1.2.1.1.7"
oidCdpCachePlatform = ".1.3.6.1.4.1.9.9.23.1.2.1.1.8"
oidCdpCacheCapabilities = ".1.3.6.1.4.1.9.9.23.1.2.1.1.9"
)
// IP-MIB / BRIDGE-MIB
const (
oidIPNetToMediaPhysAddress = ".1.3.6.1.2.1.4.22.1.2"
oidDot1dTpFdbPort = ".1.3.6.1.2.1.17.4.3.1.2"
oidDot1dBasePortIfIndex = ".1.3.6.1.2.1.17.1.4.1.2"
)
// IF-MIB
const (
oidIfDescr = ".1.3.6.1.2.1.2.2.1.2"
oidIfType = ".1.3.6.1.2.1.2.2.1.3"
oidIfMtu = ".1.3.6.1.2.1.2.2.1.4"
oidIfSpeed = ".1.3.6.1.2.1.2.2.1.5"
oidIfPhysAddress = ".1.3.6.1.2.1.2.2.1.6"
oidIfAdminStatus = ".1.3.6.1.2.1.2.2.1.7"
oidIfOperStatus = ".1.3.6.1.2.1.2.2.1.8"
oidIfName = ".1.3.6.1.2.1.31.1.1.1.1"
oidIfHighSpeed = ".1.3.6.1.2.1.31.1.1.1.15"
oidIfAlias = ".1.3.6.1.2.1.31.1.1.1.18"
)
// Params — вміст params_json для topology.discover.
type Params struct {
// lldp | cdp | arp | fdb. Порожньо — lldp + cdp.
Protos []string `json:"protos"`
// Збирати інвентар портів разом із сусідами. Обидва беруться з
// одного SNMP-обходу, тому окремий чек був би зайвим трафіком.
CollectInterfaces *bool `json:"collect_interfaces"`
}
type Module struct{}
func New() *Module { return &Module{} }
func (m *Module) Key() string { return "topology" }
func (m *Module) CheckTypes() []string { return []string{"topology.discover"} }
func (m *Module) Close() error { return nil }
func (m *Module) Run(ctx context.Context, task module.Task) (module.Result, error) {
var p Params
if len(task.Params) > 0 {
if err := json.Unmarshal(task.Params, &p); err != nil {
return module.Result{}, fmt.Errorf("невалідні params для topology.discover: %w", err)
}
}
if len(p.Protos) == 0 {
p.Protos = []string{"lldp", "cdp"}
}
collectIfaces := p.CollectInterfaces == nil || *p.CollectInterfaces
client, err := snmpx.Dial(ctx, task.Target.Address, task.Credentials, task.Timeout)
if err != nil {
return module.Result{}, err
}
defer client.Conn.Close()
var (
res module.Result
errs []string
)
// Інвентар портів потрібен першим: LLDP оперує власною нумерацією
// портів, і без ifName/ifDescr її нема на що відобразити.
ifaces, err := collectInterfaces(ctx, client, task.DeviceID)
if err != nil {
errs = append(errs, "інтерфейси: "+err.Error())
} else if collectIfaces {
res.InterfaceRecords = ifaces
}
byIndex := make(map[int64]*npv1.InterfaceRecord, len(ifaces))
for _, r := range ifaces {
byIndex[r.IfIndex] = r
}
for _, proto := range p.Protos {
var (
found []*npv1.NeighborRecord
perr error
)
switch strings.ToLower(proto) {
case "lldp":
found, perr = collectLLDP(ctx, client, task.DeviceID, byIndex)
case "cdp":
found, perr = collectCDP(ctx, client, task.DeviceID, byIndex)
case "arp":
found, perr = collectARP(ctx, client, task.DeviceID, byIndex)
case "fdb":
found, perr = collectFDB(ctx, client, task.DeviceID, byIndex)
default:
continue
}
if perr != nil {
// Відсутність LLDP-MIB — норма для половини обладнання,
// а не привід завалити весь чек: CDP може працювати.
errs = append(errs, proto+": "+perr.Error())
continue
}
res.Neighbors = append(res.Neighbors, found...)
}
if len(res.Neighbors) == 0 && len(res.InterfaceRecords) == 0 && len(errs) > 0 {
return module.Result{}, fmt.Errorf("автовиявлення не дало результату: %s", strings.Join(errs, "; "))
}
if len(errs) > 0 {
res.Payload, _ = json.Marshal(map[string]any{"warnings": errs})
}
return res, nil
}
// ---------------------------------------------------------------------
// Інтерфейси
// ---------------------------------------------------------------------
func collectInterfaces(ctx context.Context, c *gosnmp.GoSNMP, deviceID string) ([]*npv1.InterfaceRecord, error) {
recs := make(map[uint64]*npv1.InterfaceRecord)
get := func(idx uint64) *npv1.InterfaceRecord {
r, ok := recs[idx]
if !ok {
r = &npv1.InterfaceRecord{DeviceId: deviceID, IfIndex: int64(idx)}
recs[idx] = r
}
return r
}
columns := []struct {
oid string
apply func(r *npv1.InterfaceRecord, pdu gosnmp.SnmpPDU)
}{
{oidIfDescr, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
if r.Name == "" {
r.Name = snmpx.AsString(p)
}
}},
// ifName точніший за ifDescr і має пріоритет: саме його
// віддає LLDP як port-id, тому саме за ним зійдеться лінк.
{oidIfName, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
if s := snmpx.AsString(p); s != "" {
r.Name = s
}
}},
{oidIfAlias, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
r.Alias = snmpx.AsString(p)
}},
{oidIfType, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
if v, ok := snmpx.AsUint(p); ok {
r.Type = ifTypeName(v)
}
}},
{oidIfMtu, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
if v, ok := snmpx.AsUint(p); ok {
r.Mtu = uint32(v)
}
}},
{oidIfSpeed, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
if v, ok := snmpx.AsUint(p); ok && r.SpeedBps == 0 {
r.SpeedBps = v
}
}},
// ifSpeed 32-бітний і на 10G+ упирається в 4294967295.
// ifHighSpeed (у Мбіт/с) — єдине джерело правди для швидких портів.
{oidIfHighSpeed, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
if v, ok := snmpx.AsUint(p); ok && v > 0 {
r.SpeedBps = v * 1_000_000
}
}},
{oidIfPhysAddress, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
if b, ok := snmpx.AsBytes(p); ok {
r.Mac = snmpx.FormatMAC(b)
}
}},
{oidIfAdminStatus, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
if v, ok := snmpx.AsUint(p); ok {
r.AdminStatus = adminStatusName(v)
}
}},
{oidIfOperStatus, func(r *npv1.InterfaceRecord, p gosnmp.SnmpPDU) {
if v, ok := snmpx.AsUint(p); ok {
r.OperStatus = operStatusName(v)
}
}},
}
for _, col := range columns {
oid := col.oid
apply := col.apply
err := snmpx.Walk(ctx, c, oid, func(pdu gosnmp.SnmpPDU) error {
idx, ok := snmpx.SplitIndex(oid, pdu.Name)
if !ok || len(idx) != 1 {
return nil
}
apply(get(idx[0]), pdu)
return nil
})
if err != nil && oid == oidIfDescr {
// Без ifDescr таблиці інтерфейсів фактично немає.
return nil, err
}
}
out := make([]*npv1.InterfaceRecord, 0, len(recs))
for _, r := range recs {
if r.Name == "" {
r.Name = fmt.Sprintf("if%d", r.IfIndex)
}
out = append(out, r)
}
return out, nil
}
// ---------------------------------------------------------------------
// LLDP
// ---------------------------------------------------------------------
// collectLLDP читає lldpRemTable.
//
// Індекс рядка: timeMark.lldpRemLocalPortNum.lldpRemIndex. Другий
// компонент — НЕ ifIndex, а власна нумерація LLDP. Більшість
// реалізацій робить її рівною ifIndex, але покладатись на це не можна,
// тому спершу пробуємо відобразити через lldpLocPortId (ім'я порту).
func collectLLDP(ctx context.Context, c *gosnmp.GoSNMP, deviceID string, byIndex map[int64]*npv1.InterfaceRecord) ([]*npv1.NeighborRecord, error) {
localPorts, err := lldpLocalPortMap(ctx, c)
if err != nil {
return nil, err
}
byRow := make(map[string]*npv1.NeighborRecord)
subtypes := make(map[string]uint64)
portSubtypes := make(map[string]uint64)
row := func(idx snmpx.Index) *npv1.NeighborRecord {
key := idx.Key()
r, ok := byRow[key]
if !ok {
r = &npv1.NeighborRecord{
DeviceId: deviceID,
Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_LLDP,
SeenAt: timestamppb.Now(),
}
if portNum, ok := idx.At(1); ok {
resolveLocalPort(r, portNum, localPorts, byIndex)
}
byRow[key] = r
}
return r
}
walkCol := func(oid string, fn func(r *npv1.NeighborRecord, idx snmpx.Index, pdu gosnmp.SnmpPDU)) error {
return snmpx.Walk(ctx, c, oid, func(pdu gosnmp.SnmpPDU) error {
idx, ok := snmpx.SplitIndex(oid, pdu.Name)
if !ok || len(idx) < 3 {
return nil
}
fn(row(idx), idx, pdu)
return nil
})
}
// Підтипи читаємо першими: від них залежить, як тлумачити
// chassis-id і port-id — як MAC чи як текст.
if err := walkCol(oidLldpRemChassisIDSubtype, func(_ *npv1.NeighborRecord, idx snmpx.Index, pdu gosnmp.SnmpPDU) {
if v, ok := snmpx.AsUint(pdu); ok {
subtypes[idx.Key()] = v
}
}); err != nil {
return nil, err
}
_ = walkCol(oidLldpRemPortIDSubtype, func(_ *npv1.NeighborRecord, idx snmpx.Index, pdu gosnmp.SnmpPDU) {
if v, ok := snmpx.AsUint(pdu); ok {
portSubtypes[idx.Key()] = v
}
})
if err := walkCol(oidLldpRemChassisID, func(r *npv1.NeighborRecord, idx snmpx.Index, pdu gosnmp.SnmpPDU) {
raw, _ := snmpx.AsBytes(pdu)
// subtype 4 = macAddress
if subtypes[idx.Key()] == 4 {
if mac := snmpx.FormatMAC(raw); mac != "" {
r.RemoteChassisId = mac
r.RemoteMac = mac
return
}
}
r.RemoteChassisId = snmpx.AsString(pdu)
}); err != nil {
return nil, err
}
_ = walkCol(oidLldpRemPortID, func(r *npv1.NeighborRecord, idx snmpx.Index, pdu gosnmp.SnmpPDU) {
raw, _ := snmpx.AsBytes(pdu)
// subtype 3 = macAddress
if portSubtypes[idx.Key()] == 3 {
if mac := snmpx.FormatMAC(raw); mac != "" {
r.RemotePortId = mac
return
}
}
r.RemotePortId = snmpx.AsString(pdu)
})
_ = walkCol(oidLldpRemPortDesc, func(r *npv1.NeighborRecord, _ snmpx.Index, pdu gosnmp.SnmpPDU) {
r.RemotePortDescr = snmpx.AsString(pdu)
})
_ = walkCol(oidLldpRemSysName, func(r *npv1.NeighborRecord, _ snmpx.Index, pdu gosnmp.SnmpPDU) {
r.RemoteSystemName = snmpx.AsString(pdu)
})
_ = walkCol(oidLldpRemSysCapEnabled, func(r *npv1.NeighborRecord, _ snmpx.Index, pdu gosnmp.SnmpPDU) {
if raw, ok := snmpx.AsBytes(pdu); ok {
r.RemoteCapabilities = decodeLldpCaps(raw)
}
})
out := make([]*npv1.NeighborRecord, 0, len(byRow))
for _, r := range byRow {
// Сусід без жодного ідентифікатора не зіставиться ні з чим —
// краще не засмічувати topo.neighbors.
if r.RemoteChassisId == "" && r.RemoteSystemName == "" && r.RemoteMac == "" {
continue
}
out = append(out, r)
}
return out, nil
}
// lldpLocalPortMap: lldpRemLocalPortNum → ім'я локального порту.
func lldpLocalPortMap(ctx context.Context, c *gosnmp.GoSNMP) (map[uint64]string, error) {
subtypes := make(map[uint64]uint64)
_ = snmpx.Walk(ctx, c, oidLldpLocPortIDSubtype, func(pdu gosnmp.SnmpPDU) error {
if idx, ok := snmpx.SplitIndex(oidLldpLocPortIDSubtype, pdu.Name); ok && len(idx) == 1 {
if v, ok := snmpx.AsUint(pdu); ok {
subtypes[idx[0]] = v
}
}
return nil
})
names := make(map[uint64]string)
err := snmpx.Walk(ctx, c, oidLldpLocPortID, func(pdu gosnmp.SnmpPDU) error {
idx, ok := snmpx.SplitIndex(oidLldpLocPortID, pdu.Name)
if !ok || len(idx) != 1 {
return nil
}
// subtype 5 = interfaceName, 7 = local — обидва текстові.
names[idx[0]] = snmpx.AsString(pdu)
return nil
})
if err != nil {
return names, err
}
return names, nil
}
// resolveLocalPort прив'язує запис до локального інтерфейсу.
func resolveLocalPort(r *npv1.NeighborRecord, portNum uint64, localPorts map[uint64]string, byIndex map[int64]*npv1.InterfaceRecord) {
name := localPorts[portNum]
if name != "" {
r.LocalPortName = name
for _, ifc := range byIndex {
if strings.EqualFold(ifc.Name, name) || strings.EqualFold(ifc.Alias, name) {
r.LocalIfIndex = ifc.IfIndex
return
}
}
}
// Запасний варіант: більшість реалізацій робить lldpLocPortNum
// рівним ifIndex. Приймаємо лише якщо такий інтерфейс справді є —
// інакше сервер отримає посилання в нікуди.
if ifc, ok := byIndex[int64(portNum)]; ok {
r.LocalIfIndex = ifc.IfIndex
if r.LocalPortName == "" {
r.LocalPortName = ifc.Name
}
}
}
// decodeLldpCaps розбирає бітову маску можливостей (LLDP-MIB).
func decodeLldpCaps(raw []byte) []string {
if len(raw) == 0 {
return nil
}
names := []string{"other", "repeater", "bridge", "wlan-ap", "router",
"telephone", "docsis", "station"}
var out []string
bits := raw[0]
for i, name := range names {
if i >= 8 {
break
}
if bits&(1<<(7-uint(i))) != 0 {
out = append(out, name)
}
}
return out
}
// ---------------------------------------------------------------------
// CDP
// ---------------------------------------------------------------------
func collectCDP(ctx context.Context, c *gosnmp.GoSNMP, deviceID string, byIndex map[int64]*npv1.InterfaceRecord) ([]*npv1.NeighborRecord, error) {
byRow := make(map[string]*npv1.NeighborRecord)
row := func(idx snmpx.Index) *npv1.NeighborRecord {
key := idx.Key()
r, ok := byRow[key]
if !ok {
r = &npv1.NeighborRecord{
DeviceId: deviceID,
Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_CDP,
SeenAt: timestamppb.Now(),
}
// У CDP перший компонент індексу — одразу ifIndex.
if ifIdx, ok := idx.At(0); ok {
if ifc, exists := byIndex[int64(ifIdx)]; exists {
r.LocalIfIndex = ifc.IfIndex
r.LocalPortName = ifc.Name
} else {
r.LocalIfIndex = int64(ifIdx)
}
}
byRow[key] = r
}
return r
}
walkCol := func(oid string, fn func(r *npv1.NeighborRecord, pdu gosnmp.SnmpPDU)) error {
return snmpx.Walk(ctx, c, oid, func(pdu gosnmp.SnmpPDU) error {
idx, ok := snmpx.SplitIndex(oid, pdu.Name)
if !ok || len(idx) < 2 {
return nil
}
fn(row(idx), pdu)
return nil
})
}
if err := walkCol(oidCdpCacheDeviceID, func(r *npv1.NeighborRecord, pdu gosnmp.SnmpPDU) {
s := snmpx.AsString(pdu)
r.RemoteSystemName = s
// CDP device-id часто є серійником або MAC, а не іменем.
// Якщо це схоже на MAC — кладемо і в chassis_id, бо саме за
// ним сервер зіставляє надійніше, ніж за іменем.
if raw, ok := snmpx.AsBytes(pdu); ok {
if mac := snmpx.FormatMAC(raw); mac != "" {
r.RemoteChassisId = mac
r.RemoteMac = mac
}
}
if r.RemoteChassisId == "" {
r.RemoteChassisId = s
}
}); err != nil {
return nil, err
}
_ = walkCol(oidCdpCacheDevicePort, func(r *npv1.NeighborRecord, pdu gosnmp.SnmpPDU) {
r.RemotePortId = snmpx.AsString(pdu)
})
_ = walkCol(oidCdpCachePlatform, func(r *npv1.NeighborRecord, pdu gosnmp.SnmpPDU) {
r.RemotePlatform = snmpx.AsString(pdu)
})
_ = walkCol(oidCdpCacheVersion, func(r *npv1.NeighborRecord, pdu gosnmp.SnmpPDU) {
if r.RemotePortDescr == "" {
r.RemotePortDescr = truncate(snmpx.AsString(pdu), 120)
}
})
_ = walkCol(oidCdpCacheAddress, func(r *npv1.NeighborRecord, pdu gosnmp.SnmpPDU) {
if raw, ok := snmpx.AsBytes(pdu); ok {
r.RemoteMgmtIp = snmpx.FormatIP(raw)
}
})
_ = walkCol(oidCdpCacheCapabilities, func(r *npv1.NeighborRecord, pdu gosnmp.SnmpPDU) {
if raw, ok := snmpx.AsBytes(pdu); ok {
r.RemoteCapabilities = decodeCdpCaps(raw)
}
})
out := make([]*npv1.NeighborRecord, 0, len(byRow))
for _, r := range byRow {
if r.RemoteChassisId == "" && r.RemoteSystemName == "" {
continue
}
out = append(out, r)
}
return out, nil
}
func decodeCdpCaps(raw []byte) []string {
if len(raw) < 4 {
return nil
}
bits := uint32(raw[0])<<24 | uint32(raw[1])<<16 | uint32(raw[2])<<8 | uint32(raw[3])
names := map[uint32]string{
0x01: "router", 0x02: "transparent-bridge", 0x04: "source-route-bridge",
0x08: "switch", 0x10: "host", 0x20: "igmp", 0x40: "repeater",
}
var out []string
for bit, name := range names {
if bits&bit != 0 {
out = append(out, name)
}
}
return out
}
// ---------------------------------------------------------------------
// ARP і FDB
//
// Це не «сусіди» в сенсі LLDP: ARP каже лише, що якийсь MAC живе за
// цим портом. Але для пристроїв без LLDP/CDP (принтери, IP-камери,
// дешеві комутатори) це єдиний шанс потрапити на мапу, тому сервер
// приймає їх з низькою впевненістю.
// ---------------------------------------------------------------------
func collectARP(ctx context.Context, c *gosnmp.GoSNMP, deviceID string, byIndex map[int64]*npv1.InterfaceRecord) ([]*npv1.NeighborRecord, error) {
var out []*npv1.NeighborRecord
err := snmpx.Walk(ctx, c, oidIPNetToMediaPhysAddress, func(pdu gosnmp.SnmpPDU) error {
idx, ok := snmpx.SplitIndex(oidIPNetToMediaPhysAddress, pdu.Name)
if !ok || len(idx) < 5 {
return nil
}
ip, ok := idx.IPv4At(1)
if !ok {
return nil
}
raw, _ := snmpx.AsBytes(pdu)
mac := snmpx.FormatMAC(raw)
if mac == "" {
return nil
}
r := &npv1.NeighborRecord{
DeviceId: deviceID,
Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_ARP,
RemoteMac: mac,
RemoteMgmtIp: ip,
SeenAt: timestamppb.Now(),
}
if ifc, exists := byIndex[int64(idx[0])]; exists {
r.LocalIfIndex = ifc.IfIndex
r.LocalPortName = ifc.Name
} else {
r.LocalIfIndex = int64(idx[0])
}
out = append(out, r)
return nil
})
return out, err
}
func collectFDB(ctx context.Context, c *gosnmp.GoSNMP, deviceID string, byIndex map[int64]*npv1.InterfaceRecord) ([]*npv1.NeighborRecord, error) {
// dot1dTpFdbPort віддає НЕ ifIndex, а номер мосту (bridge port).
// Без dot1dBasePortIfIndex ці числа ні з чим не зіставити.
bridgeToIf := make(map[uint64]int64)
_ = snmpx.Walk(ctx, c, oidDot1dBasePortIfIndex, func(pdu gosnmp.SnmpPDU) error {
idx, ok := snmpx.SplitIndex(oidDot1dBasePortIfIndex, pdu.Name)
if !ok || len(idx) != 1 {
return nil
}
if v, ok := snmpx.AsUint(pdu); ok {
bridgeToIf[idx[0]] = int64(v)
}
return nil
})
var out []*npv1.NeighborRecord
err := snmpx.Walk(ctx, c, oidDot1dTpFdbPort, func(pdu gosnmp.SnmpPDU) error {
idx, ok := snmpx.SplitIndex(oidDot1dTpFdbPort, pdu.Name)
if !ok || len(idx) < 6 {
return nil
}
mac, ok := idx.MACAt(len(idx) - 6)
if !ok {
return nil
}
port, ok := snmpx.AsUint(pdu)
if !ok || port == 0 {
return nil
}
r := &npv1.NeighborRecord{
DeviceId: deviceID,
Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_FDB,
RemoteMac: mac,
SeenAt: timestamppb.Now(),
}
if ifIdx, exists := bridgeToIf[port]; exists {
r.LocalIfIndex = ifIdx
if ifc, ok := byIndex[ifIdx]; ok {
r.LocalPortName = ifc.Name
}
}
out = append(out, r)
return nil
})
return out, err
}
// ---------------------------------------------------------------------
func adminStatusName(v uint64) string {
switch v {
case 1:
return "up"
case 2:
return "down"
case 3:
return "testing"
default:
return "unknown"
}
}
func operStatusName(v uint64) string {
switch v {
case 1:
return "up"
case 2:
return "down"
case 3:
return "testing"
case 5:
return "dormant"
case 6:
return "notPresent"
case 7:
return "lowerLayerDown"
default:
return "unknown"
}
}
func ifTypeName(v uint64) string {
switch v {
case 6:
return "ethernetCsmacd"
case 24:
return "softwareLoopback"
case 53:
return "propVirtual"
case 131:
return "tunnel"
case 135:
return "l2vlan"
case 136:
return "l3ipvlan"
case 161:
return "ieee8023adLag"
default:
return fmt.Sprintf("type%d", v)
}
}
func truncate(s string, n int) string {
if len(s) <= n {
return s
}
return s[:n]
}

View file

@ -12,6 +12,7 @@ package scheduler
import (
"container/heap"
"context"
"strings"
"sync"
"time"
@ -30,6 +31,13 @@ type Sink interface {
// StatusFunc доповідає серверу про життєвий цикл задачі.
type StatusFunc func(*npv1.TaskStatusUpdate)
// DiscoveryFunc отримує результати автовиявлення.
//
// Вони не йдуть телеметричним стрімом: сусіди й інвентар портів рідкі,
// великі й не прив'язані до моменту часу так, як метрики. Для них
// окремий RPC ReportDiscovery.
type DiscoveryFunc func(res module.Result)
// CredentialFunc віддає облікові дані пристрою на момент виконання.
//
// Навмисно не поле в Task: креденшели мають TTL і оновлюються окремим
@ -53,6 +61,7 @@ type Scheduler struct {
sink Sink
onStatus StatusFunc
creds CredentialFunc
onDisco DiscoveryFunc
mu sync.Mutex
entries map[string]*entry
@ -70,6 +79,7 @@ type Config struct {
Sink Sink
OnStatus StatusFunc
Credentials CredentialFunc
OnDiscovery DiscoveryFunc
MaxConcurrency int
// Підміна годинника в тестах.
Now func() time.Time
@ -88,11 +98,15 @@ func New(cfg Config) *Scheduler {
if cfg.Credentials == nil {
cfg.Credentials = func(string) []*npv1.Credential { return nil }
}
if cfg.OnDiscovery == nil {
cfg.OnDiscovery = func(module.Result) {}
}
return &Scheduler{
registry: cfg.Registry,
sink: cfg.Sink,
onStatus: cfg.OnStatus,
creds: cfg.Credentials,
onDisco: cfg.OnDiscovery,
entries: make(map[string]*entry),
wake: make(chan struct{}, 1),
sem: make(chan struct{}, cfg.MaxConcurrency),
@ -210,6 +224,45 @@ type TaskMeta struct {
Offset time.Duration
}
// TriggerNow ставить у чергу негайне виконання задач, що підпадають
// під фільтр. Потрібне для кнопки «Запустити виявлення» в UI: сервер
// надсилає DiscoveryRequest, а не чекає наступного циклу розкладу.
//
// Повертає, скільки задач зрушено — сервер має знати, що його запит
// не потрапив у порожнечу.
func (s *Scheduler) TriggerNow(deviceIDs []string, checkTypePrefix string) int {
s.mu.Lock()
defer s.mu.Unlock()
want := make(map[string]bool, len(deviceIDs))
for _, id := range deviceIDs {
want[id] = true
}
now := s.nowFn()
n := 0
for _, e := range s.entries {
if checkTypePrefix != "" && !strings.HasPrefix(e.task.CheckType, checkTypePrefix) {
continue
}
if len(want) > 0 && !want[e.task.DeviceID] {
continue
}
if e.running {
continue
}
e.nextRun = now
if e.index >= 0 && e.index < len(s.queue) {
heap.Fix(&s.queue, e.index)
}
n++
}
if n > 0 {
s.wakeup()
}
return n
}
// Len — скільки задач у розкладі.
func (s *Scheduler) Len() int {
s.mu.Lock()
@ -410,6 +463,9 @@ func (s *Scheduler) execute(ctx context.Context, e *entry) {
s.sink.Add(task.DeviceID, module.ModuleKey(task.CheckType), res)
s.sink.AddCheckResult(cr)
if len(res.Neighbors) > 0 || len(res.InterfaceRecords) > 0 {
s.onDisco(res)
}
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: task.CheckID,
State: npv1.TaskStatusUpdate_STATE_SUCCEEDED,

View file

@ -63,6 +63,10 @@ type Session struct {
// опитування через повільний контрольний канал.
statusCh chan *npv1.TaskStatusUpdate
// Звіти автовиявлення. Окремо від телеметрії: вони рідкі, великі
// й їдуть унарним ReportDiscovery, а не стрімом.
discoCh chan *npv1.DiscoveryReport
credMu sync.RWMutex
creds map[string][]*npv1.Credential
credExpiry time.Time
@ -94,6 +98,7 @@ func New(cfg Config) *Session {
cfg: cfg,
log: cfg.Logger,
statusCh: make(chan *npv1.TaskStatusUpdate, 512),
discoCh: make(chan *npv1.DiscoveryReport, 16),
creds: make(map[string][]*npv1.Credential),
devices: make(map[string]*npv1.DeviceTarget),
startedAt: time.Now(),
@ -120,6 +125,40 @@ func (s *Session) Credentials(deviceID string) []*npv1.Credential {
return s.creds[deviceID]
}
// ReportDiscovery — колбек для планувальника.
//
// Черга навмисно коротка: звіт автовиявлення актуальний рівно доти,
// доки описує поточний стан мережі. Накопичувати десяток застарілих
// знімків, поки немає зв'язку, немає сенсу — наступний обхід віддасть
// свіжіші дані.
func (s *Session) ReportDiscovery(res module.Result) {
if len(res.Neighbors) == 0 && len(res.InterfaceRecords) == 0 {
return
}
rep := &npv1.DiscoveryReport{
AgentId: s.cfg.AgentID,
Neighbors: res.Neighbors,
Interfaces: res.InterfaceRecords,
Final: true,
StartedAt: timestamppb.Now(),
FinishedAt: timestamppb.Now(),
}
select {
case s.discoCh <- rep:
default:
// Витісняємо найстаріший, а не відкидаємо новий.
select {
case <-s.discoCh:
default:
}
select {
case s.discoCh <- rep:
default:
}
}
}
// ReportStatus — колбек для планувальника.
func (s *Session) ReportStatus(u *npv1.TaskStatusUpdate) {
select {
@ -220,7 +259,9 @@ func (s *Session) runOnce(ctx context.Context) error {
out := make(chan *npv1.ControlUp, 256)
var wg sync.WaitGroup
errCh := make(chan error, 4)
// Місткість дорівнює числу горутин: інакше та, що впала останньою,
// заблокується на записі й ніколи не дочекається wg.Wait().
errCh := make(chan error, 8)
spawn := func(name string, fn func() error) {
wg.Add(1)
@ -237,6 +278,7 @@ func (s *Session) runOnce(ctx context.Context) error {
spawn("heartbeat", func() error { return s.heartbeatLoop(sctx, out, welcome) })
spawn("status", func() error { return s.statusLoop(sctx, out) })
spawn("telemetry", func() error { return s.telemetryLoop(sctx, client, welcome) })
spawn("discovery", func() error { return s.discoveryLoop(sctx, client) })
// Читання команд — у цій же горутині.
readErr := s.controlLoop(sctx, ctrl, out)
@ -404,6 +446,14 @@ func (s *Session) controlLoop(ctx context.Context, ctrl npv1.AgentService_Contro
s.skewNanos.Store(int64(time.Until(p.Ping.ServerTime.AsTime())))
}
case *npv1.ControlDown_DiscoveryRequest:
if s.cfg.Scheduler == nil {
continue
}
n := s.cfg.Scheduler.TriggerNow(p.DiscoveryRequest.GetDeviceIds(), "topo.")
s.log.Info("сервер попросив запустити автовиявлення",
"run_id", p.DiscoveryRequest.GetRunId(), адач_зрушено", n)
case *npv1.ControlDown_Directive:
if stop := s.applyDirective(p.Directive); stop {
return nil
@ -546,6 +596,35 @@ func (s *Session) buildTasks(in []*npv1.Task) ([]module.Task, []scheduler.TaskMe
return tasks, meta
}
// ---------------------------------------------------------------------
// Автовиявлення
// ---------------------------------------------------------------------
func (s *Session) discoveryLoop(ctx context.Context, client npv1.AgentServiceClient) error {
for {
select {
case <-ctx.Done():
return nil
case rep := <-s.discoCh:
ack, err := client.ReportDiscovery(ctx, rep)
if err != nil {
if ctx.Err() != nil {
return nil
}
// Звіт втрачено свідомо: наступний обхід за розкладом
// віддасть свіжіші дані, а ретрай застарілого знімка
// топології нічого не додає.
s.log.Warn("звіт автовиявлення не доставлено", "err", err)
continue
}
s.log.Info("автовиявлення прийнято сервером",
"сусідів", len(rep.Neighbors),
"зіставлено", ack.GetNeighborsResolved(),
"лінків", ack.GetLinksCreated())
}
}
}
// ---------------------------------------------------------------------
// Телеметрія
// ---------------------------------------------------------------------

View file

@ -0,0 +1,233 @@
// Package snmpx — спільний SNMP-транспорт для модулів агента.
//
// Винесено окремо, бо ним користуються два різні модулі: snmp (метрики)
// і topology (виявлення сусідів). Дублювати логіку вибору креденшелів
// і розбиття запитів на порції в обох — вірний спосіб отримати два
// різні набори багів.
package snmpx
import (
"context"
"fmt"
"math"
"math/big"
"strings"
"time"
"github.com/gosnmp/gosnmp"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
)
// MaxVarsPerPDU — скільки змінних класти в один запит.
//
// Більшість агентів на обладнанні відмовляють приблизно після 30;
// беремо із запасом, бо відмова виглядає як «пристрій не відповідає»
// і плутає діагностику.
const MaxVarsPerPDU = 24
// Dial відкриває SNMP-сесію. Викликач зобов'язаний закрити Conn.
func Dial(ctx context.Context, address string, creds []*npv1.Credential, timeout time.Duration) (*gosnmp.GoSNMP, error) {
cred := PickCredential(creds)
if cred == nil {
return nil, fmt.Errorf("немає SNMP-креденшелів для %s", address)
}
port := uint16(161)
if cred.Port > 0 && cred.Port < 65536 {
port = uint16(cred.Port)
}
if timeout <= 0 {
timeout = 5 * time.Second
}
client := &gosnmp.GoSNMP{
Target: hostOnly(address),
Port: port,
Timeout: timeout,
Retries: 1,
Context: ctx,
// Пристрої з великими таблицями інтерфейсів інакше віддають
// tooBig на кожен walk.
MaxOids: MaxVarsPerPDU,
}
switch cred.Transport {
case npv1.Transport_TRANSPORT_SNMP_V2C:
client.Version = gosnmp.Version2c
client.Community = cred.GetCommunity()
case npv1.Transport_TRANSPORT_SNMP_V3:
v3 := cred.GetSnmpV3()
if v3 == nil {
return nil, fmt.Errorf("SNMPv3 без параметрів безпеки")
}
client.Version = gosnmp.Version3
client.SecurityModel = gosnmp.UserSecurityModel
client.MsgFlags = msgFlags(v3.Level)
client.ContextName = v3.ContextName
name := v3.SecurityName
if name == "" {
name = cred.Username
}
client.SecurityParameters = &gosnmp.UsmSecurityParameters{
UserName: name,
AuthenticationProtocol: authProto(v3.AuthProtocol),
AuthenticationPassphrase: v3.AuthPassword,
PrivacyProtocol: privProto(v3.PrivProtocol),
PrivacyPassphrase: v3.PrivPassword,
}
default:
return nil, fmt.Errorf("креденшел %s не для SNMP", cred.Transport)
}
if err := client.Connect(); err != nil {
return nil, fmt.Errorf("SNMP-сесія до %s: %w", address, err)
}
return client, nil
}
func hostOnly(address string) string {
if i := strings.LastIndex(address, ":"); i > 0 && !strings.Contains(address, "]") {
if strings.Count(address, ":") == 1 {
return address[:i]
}
}
return address
}
// PickCredential віддає перевагу v3 над v2c: якщо адміністратор
// налаштував обидва, шифрований варіант очевидно бажаніший.
func PickCredential(creds []*npv1.Credential) *npv1.Credential {
for _, c := range creds {
if c.Transport == npv1.Transport_TRANSPORT_SNMP_V3 {
return c
}
}
for _, c := range creds {
if c.Transport == npv1.Transport_TRANSPORT_SNMP_V2C {
return c
}
}
return nil
}
// GetUints робить Get порціями й повертає лише числові значення.
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 {
if err := ctx.Err(); err != nil {
return nil, err
}
end := min(start+MaxVarsPerPDU, len(oids))
pkt, err := c.Get(oids[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
}
}
}
return out, nil
}
// AsUint дістає беззнакове ціле з PDU, відсіюючи службові типи.
func AsUint(pdu gosnmp.SnmpPDU) (uint64, bool) {
switch pdu.Type {
case gosnmp.NoSuchObject, gosnmp.NoSuchInstance, gosnmp.EndOfMibView, gosnmp.Null:
return 0, false
}
v := gosnmp.ToBigInt(pdu.Value)
if v == nil || v.Sign() < 0 || v.Cmp(maxUint64) > 0 {
return 0, false
}
return v.Uint64(), true
}
var maxUint64 = new(big.Int).SetUint64(math.MaxUint64)
// AsBytes дістає OCTET STRING.
func AsBytes(pdu gosnmp.SnmpPDU) ([]byte, bool) {
b, ok := pdu.Value.([]byte)
return b, ok
}
// AsString дістає текст, прибираючи нульові байти й пробіли по краях.
//
// Обладнання регулярно повертає ім'я порту з кінцевим \x00 або
// вирівняне пробілами; без нормалізації те саме ім'я не збігається
// саме з собою при зіставленні лінків.
func AsString(pdu gosnmp.SnmpPDU) string {
switch v := pdu.Value.(type) {
case []byte:
return strings.TrimSpace(strings.TrimRight(string(v), "\x00"))
case string:
return strings.TrimSpace(v)
default:
if u, ok := AsUint(pdu); ok {
return fmt.Sprint(u)
}
return ""
}
}
// Walk обходить піддерево, викликаючи fn на кожній змінній.
func Walk(ctx context.Context, c *gosnmp.GoSNMP, root string, fn func(gosnmp.SnmpPDU) error) error {
if err := ctx.Err(); err != nil {
return err
}
if c.Version == gosnmp.Version1 {
return c.Walk(root, fn)
}
return c.BulkWalk(root, fn)
}
func msgFlags(level npv1.SnmpV3Options_SecurityLevel) gosnmp.SnmpV3MsgFlags {
switch level {
case npv1.SnmpV3Options_SECURITY_LEVEL_AUTH_PRIV:
return gosnmp.AuthPriv
case npv1.SnmpV3Options_SECURITY_LEVEL_AUTH_NO_PRIV:
return gosnmp.AuthNoPriv
default:
return gosnmp.NoAuthNoPriv
}
}
func authProto(name string) gosnmp.SnmpV3AuthProtocol {
switch name {
case "MD5":
return gosnmp.MD5
case "SHA":
return gosnmp.SHA
case "SHA224":
return gosnmp.SHA224
case "SHA256":
return gosnmp.SHA256
case "SHA384":
return gosnmp.SHA384
case "SHA512":
return gosnmp.SHA512
default:
return gosnmp.NoAuth
}
}
func privProto(name string) gosnmp.SnmpV3PrivProtocol {
switch name {
case "DES":
return gosnmp.DES
case "AES":
return gosnmp.AES
case "AES192":
return gosnmp.AES192
case "AES256":
return gosnmp.AES256
default:
return gosnmp.NoPriv
}
}

138
agent/internal/snmpx/oid.go Normal file
View file

@ -0,0 +1,138 @@
package snmpx
import (
"fmt"
"net"
"strconv"
"strings"
)
// Index — суфікс OID після базового піддерева, розібраний на числа.
//
// Уся робота з SNMP-таблицями зводиться до одного: узяти OID змінної,
// відрізати базу й прочитати індекс. Для ifTable це ifIndex, для
// cdpCacheTable — (ifIndex, deviceIndex), для ARP — (ifIndex, IPv4),
// для LLDP — (timeMark, localPortNum, remIndex).
type Index []uint64
// SplitIndex відрізає базовий OID і повертає індекс.
// Обидва OID можуть бути з провідною крапкою або без неї.
func SplitIndex(base, full string) (Index, bool) {
b := strings.TrimPrefix(base, ".")
f := strings.TrimPrefix(full, ".")
if !strings.HasPrefix(f, b+".") {
return nil, false
}
rest := f[len(b)+1:]
if rest == "" {
return nil, false
}
parts := strings.Split(rest, ".")
idx := make(Index, 0, len(parts))
for _, p := range parts {
n, err := strconv.ParseUint(p, 10, 64)
if err != nil {
return nil, false
}
idx = append(idx, n)
}
return idx, true
}
// At повертає i-й елемент індексу.
func (ix Index) At(i int) (uint64, bool) {
if i < 0 || i >= len(ix) {
return 0, false
}
return ix[i], true
}
// Key склеює індекс назад у рядок — зручно як ключ мапи, коли треба
// зіставити кілька колонок однієї таблиці.
func (ix Index) Key() string {
parts := make([]string, len(ix))
for i, v := range ix {
parts[i] = strconv.FormatUint(v, 10)
}
return strings.Join(parts, ".")
}
// IPv4At читає чотири октети індексу як IPv4-адресу.
// Так влаштований індекс ipNetToMediaTable: ifIndex.a.b.c.d
func (ix Index) IPv4At(offset int) (string, bool) {
if offset < 0 || offset+4 > len(ix) {
return "", false
}
for i := offset; i < offset+4; i++ {
if ix[i] > 255 {
return "", false
}
}
return fmt.Sprintf("%d.%d.%d.%d", ix[offset], ix[offset+1], ix[offset+2], ix[offset+3]), true
}
// MACAt читає шість октетів індексу як MAC.
// Так влаштований індекс dot1dTpFdbTable.
func (ix Index) MACAt(offset int) (string, bool) {
if offset < 0 || offset+6 > len(ix) {
return "", false
}
b := make([]byte, 6)
for i := 0; i < 6; i++ {
if ix[offset+i] > 255 {
return "", false
}
b[i] = byte(ix[offset+i])
}
return net.HardwareAddr(b).String(), true
}
// FormatMAC перетворює сирі байти в канонічний вигляд.
//
// Обладнання віддає MAC по-різному: 6 сирих байтів (норма), текстом
// із дефісами або двокрапками, іноді з великими літерами. Зводимо все
// до одного вигляду, інакше зіставлення сусіда з інвентарем провалиться
// через косметичну різницю.
func FormatMAC(raw []byte) string {
if len(raw) == 6 {
return net.HardwareAddr(raw).String()
}
s := strings.TrimSpace(strings.TrimRight(string(raw), "\x00"))
if s == "" {
return ""
}
if hw, err := net.ParseMAC(s); err == nil {
return hw.String()
}
// Формат Cisco: 0011.2233.4455
cleaned := strings.NewReplacer(".", "", ":", "", "-", "", " ", "").Replace(s)
if len(cleaned) == 12 {
b := make([]byte, 6)
for i := 0; i < 6; i++ {
v, err := strconv.ParseUint(cleaned[i*2:i*2+2], 16, 8)
if err != nil {
return ""
}
b[i] = byte(v)
}
return net.HardwareAddr(b).String()
}
return ""
}
// FormatIP розбирає адресу з OCTET STRING (4 або 16 байтів) чи тексту.
func FormatIP(raw []byte) string {
switch len(raw) {
case 4, 16:
return net.IP(raw).String()
}
s := strings.TrimSpace(strings.TrimRight(string(raw), "\x00"))
if ip := net.ParseIP(s); ip != nil {
return ip.String()
}
return ""
}

View file

@ -0,0 +1,136 @@
package snmpx
import "testing"
func TestSplitIndex(t *testing.T) {
const base = ".1.3.6.1.2.1.2.2.1.2"
cases := []struct {
name string
full string
want Index
ok bool
}{
{"проста колонка", ".1.3.6.1.2.1.2.2.1.2.5", Index{5}, true},
{"без провідної крапки", "1.3.6.1.2.1.2.2.1.2.5", Index{5}, true},
{"складений індекс", ".1.3.6.1.2.1.2.2.1.2.10.3", Index{10, 3}, true},
{"чужа гілка", ".1.3.6.1.2.1.2.2.1.3.5", nil, false},
{"сам базовий OID без індексу", ".1.3.6.1.2.1.2.2.1.2", nil, false},
// Класична пастка: .1.2.20 не є нащадком .1.2, якщо порівнювати
// рядки без роздільника — сюди й ловляться чужі колонки.
{"схожий префікс", ".1.3.6.1.2.1.2.2.1.20.5", nil, false},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
got, ok := SplitIndex(base, c.full)
if ok != c.ok {
t.Fatalf("ok = %v, очікували %v", ok, c.ok)
}
if !ok {
return
}
if len(got) != len(c.want) {
t.Fatalf("індекс = %v, очікували %v", got, c.want)
}
for i := range got {
if got[i] != c.want[i] {
t.Fatalf("індекс = %v, очікували %v", got, c.want)
}
}
})
}
}
// ipNetToMediaTable індексується як ifIndex.a.b.c.d — саме звідси
// беруться ARP-сусіди.
func TestIndexIPv4At(t *testing.T) {
idx, ok := SplitIndex(".1.3.6.1.2.1.4.22.1.2", ".1.3.6.1.2.1.4.22.1.2.3.10.0.0.7")
if !ok {
t.Fatal("індекс не розібрався")
}
if idx[0] != 3 {
t.Fatalf("ifIndex = %d", idx[0])
}
ip, ok := idx.IPv4At(1)
if !ok || ip != "10.0.0.7" {
t.Fatalf("IP = %q ok=%v", ip, ok)
}
if _, ok := idx.IPv4At(3); ok {
t.Fatal("прочитано IP за межами індексу")
}
}
// dot1dTpFdbTable індексується шістьма октетами MAC.
func TestIndexMACAt(t *testing.T) {
idx, ok := SplitIndex(".1.3.6.1.2.1.17.4.3.1.2", ".1.3.6.1.2.1.17.4.3.1.2.0.17.34.51.68.85")
if !ok {
t.Fatal("індекс не розібрався")
}
mac, ok := idx.MACAt(0)
if !ok || mac != "00:11:22:33:44:55" {
t.Fatalf("MAC = %q ok=%v", mac, ok)
}
}
func TestIndexOutOfRange(t *testing.T) {
idx := Index{1, 2}
if _, ok := idx.At(5); ok {
t.Fatal("At повернув значення за межами")
}
if _, ok := idx.MACAt(0); ok {
t.Fatal("MACAt зібрав MAC із двох чисел")
}
}
// Обладнання віддає MAC у трьох різних виглядах — усі мають зійтися
// до одного, інакше зіставлення сусіда провалиться на косметиці.
func TestFormatMAC(t *testing.T) {
want := "00:11:22:33:44:55"
cases := []struct {
name string
in []byte
want string
}{
{"сирі 6 байтів", []byte{0x00, 0x11, 0x22, 0x33, 0x44, 0x55}, want},
{"текст із двокрапками", []byte("00:11:22:33:44:55"), want},
{"формат Cisco", []byte("0011.2233.4455"), want},
{"великі літери з дефісами", []byte("00-11-22-33-44-55"), want},
{"порожньо", []byte{}, ""},
{"сміття", []byte("не мак"), ""},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
if got := FormatMAC(c.in); got != c.want {
t.Fatalf("FormatMAC = %q, очікували %q", got, c.want)
}
})
}
}
func TestFormatIP(t *testing.T) {
if got := FormatIP([]byte{10, 0, 0, 1}); got != "10.0.0.1" {
t.Fatalf("IPv4 = %q", got)
}
if got := FormatIP([]byte("192.168.1.1")); got != "192.168.1.1" {
t.Fatalf("текстовий IP = %q", got)
}
if got := FormatIP([]byte("не адреса")); got != "" {
t.Fatalf("сміття перетворилось на %q", got)
}
}
func TestHostOnly(t *testing.T) {
if got := hostOnly("10.0.0.1:161"); got != "10.0.0.1" {
t.Fatalf("hostOnly = %q", got)
}
if got := hostOnly("10.0.0.1"); got != "10.0.0.1" {
t.Fatalf("hostOnly = %q", got)
}
// IPv6 не має бути порізаний по двокрапках.
if got := hostOnly("fe80::1"); got != "fe80::1" {
t.Fatalf("hostOnly = %q", got)
}
}

View file

@ -88,7 +88,14 @@ CREATE TABLE core.check_types (
name text NOT NULL,
-- JSON Schema параметрів + перелік метрик, які повертає чек
params_schema jsonb NOT NULL DEFAULT '{}'::jsonb,
metrics jsonb NOT NULL DEFAULT '[]'::jsonb
metrics jsonb NOT NULL DEFAULT '[]'::jsonb,
-- Префікс типу чека ЗОБОВ'ЯЗАНИЙ дорівнювати ключу плагіна: саме за
-- ним агент обирає модуль-виконавця ("snmp.if" -> модуль snmp).
-- Без цього обмеження неузгодженість помічається аж у полі, коли
-- зонд відхиляє задачу як адресовану неіснуючому модулю.
CONSTRAINT check_types_prefix_matches_plugin
CHECK (key LIKE plugin_key || '.%')
);
CREATE TABLE core.checks (

View file

@ -38,9 +38,15 @@ CREATE TABLE topo.neighbors (
last_seen_at timestamptz NOT NULL DEFAULT now(),
raw jsonb NOT NULL DEFAULT '{}'::jsonb
);
-- Ключ ідентичності сусіда залежить від протоколу, тому в індексі всі
-- три ознаки. LLDP/CDP розрізняють сусідів парою chassis+port; ARP і FDB
-- не мають жодної з них — там єдиний розрізняльник це MAC. Без нього всі
-- ARP-записи одного порту схлопуються в один рядок (перевірено на живій
-- ARP-таблиці: з двох сусідів зберігався один).
CREATE UNIQUE INDEX neighbors_uniq
ON topo.neighbors (device_id, COALESCE(interface_id, '00000000-0000-0000-0000-000000000000'::uuid),
proto, COALESCE(remote_chassis_id,''), COALESCE(remote_port_id,''));
proto, COALESCE(remote_chassis_id,''), COALESCE(remote_port_id,''),
COALESCE(remote_mac::text,''));
CREATE INDEX neighbors_tenant_idx ON topo.neighbors (tenant_id, last_seen_at DESC);
CREATE INDEX neighbors_resolved_idx ON topo.neighbors (resolved_device_id);

View file

@ -85,9 +85,9 @@ INSERT INTO core.plugins (key, name, version, scope, description, manifest, min_
('snmp', 'SNMP v2c/v3', '1.0.0', 'agent', 'Опитування OID, інтерфейси, CPU/RAM/сенсори',
'{"metrics":["if.*","cpu.util","mem.used","sensor.*"],"checks":["snmp.get","snmp.walk","snmp.if"]}', 'pro', true),
('http', 'HTTP/HTTPS/SSL', '1.0.0', 'agent', 'Код відповіді, час, строк дії сертифіката',
'{"metrics":["http.status","http.latency","ssl.days_left"],"checks":["http.status","ssl.expiry"]}', 'free', true),
'{"metrics":["http.status","http.latency","ssl.days_left"],"checks":["http.status","http.ssl_expiry"]}', 'free', true),
('topology', 'Topology Discovery', '1.0.0', 'both', 'LLDP/CDP/ARP/FDB автовиявлення зв''язків',
'{"protocols":["lldp","cdp","arp","fdb"],"checks":["topo.discover"]}', 'pro', true),
'{"protocols":["lldp","cdp","arp","fdb"],"checks":["topology.discover"]}', 'pro', true),
('ncm', 'Config Backup (NCM)', '1.0.0', 'both', 'SSH/Telnet бекап конфігів + Git diff',
'{"transports":["ssh","telnet","api"],"checks":["ncm.backup"]}', 'enterprise', true),
('modbus', 'Modbus-TCP', '1.0.0', 'agent', 'Інвертори, UPS, BMS',
@ -117,10 +117,10 @@ INSERT INTO core.check_types (key, plugin_key, name, params_schema, metrics) VAL
('http.status', 'http', 'HTTP Status',
'{"type":"object","required":["url"],"properties":{"url":{"type":"string"},"expect_status":{"type":"integer","default":200},"keyword":{"type":"string"}}}',
'["http.status","http.latency_ms"]'),
('ssl.expiry', 'http', 'SSL Certificate Expiry',
('http.ssl_expiry','http','SSL Certificate Expiry',
'{"type":"object","required":["host"],"properties":{"host":{"type":"string"},"port":{"type":"integer","default":443}}}',
'["ssl.days_left"]'),
('topo.discover','topology','Neighbor Discovery',
('topology.discover','topology','Neighbor Discovery',
'{"type":"object","properties":{"protos":{"type":"array","default":["lldp","cdp","arp"]}}}',
'[]'),
('ncm.backup', 'ncm', 'Config Backup',

View file

@ -0,0 +1,176 @@
// Команда netpulse-secret — заведення зашифрованих секретів у core.secrets.
//
// Потрібна, бо секрети не можна вставити звичайним SQL: у БД лягає лише
// шифротекст, а ключ живе поза нею. Доки немає UI, це єдиний коректний
// спосіб завести пароль SSH чи SNMP-community — і він же лишається
// робочим інструментом для скриптів масового заведення.
//
// netpulse-secret -dsn ... -dek k1=<hex> -tenant <uuid> -kind snmp_v3 \
// -value public -credential snmp-ro -proto snmp_v2c \
// -attach-device <uuid>
//
// Значення читається з -value або зі stdin: у скриптах воно не має
// потрапляти в історію команд і в список процесів.
package main
import (
"bufio"
"context"
"encoding/base64"
"encoding/hex"
"errors"
"flag"
"fmt"
"os"
"strings"
"github.com/jackc/pgx/v5"
"github.com/netpulse/netpulse/server/internal/crypto"
)
func main() {
if err := run(); err != nil {
fmt.Fprintln(os.Stderr, "netpulse-secret:", err)
os.Exit(1)
}
}
func run() error {
var (
dsn = flag.String("dsn", os.Getenv("NETPULSE_DSN"), "DSN PostgreSQL")
dek = flag.String("dek", os.Getenv("NETPULSE_DEK"), "ключ: id=<hex|base64>")
tenant = flag.String("tenant", "", "uuid тенанта")
kind = flag.String("kind", "generic", "core.secret_kind: ssh_password|ssh_key|snmp_v3|api_token|...")
value = flag.String("value", "", "значення секрету; порожнє — читати зі stdin")
credName = flag.String("credential", "", "створити inv.credentials з цим іменем")
proto = flag.String("proto", "", "inv.credential_proto для -credential")
username = flag.String("username", "", "логін для -credential")
port = flag.Int("port", 0, "порт для -credential")
options = flag.String("options", "{}", "JSON-опції креденшела (SNMPv3 тощо)")
attach = flag.String("attach-device", "", "uuid пристрою, до якого прив'язати креденшел")
)
flag.Parse()
if *dsn == "" || *tenant == "" {
return errors.New("потрібні -dsn і -tenant")
}
ring, err := buildKeyring(*dek)
if err != nil {
return err
}
plaintext := *value
if plaintext == "" {
data, err := readStdin()
if err != nil {
return err
}
plaintext = data
}
if plaintext == "" {
return errors.New("порожнє значення секрету")
}
ctx := context.Background()
conn, err := pgx.Connect(ctx, *dsn)
if err != nil {
return fmt.Errorf("підключення: %w", err)
}
defer conn.Close(ctx)
// AAD прив'язує шифротекст до тенанта й типу: переставити рядок
// на іншого тенанта прямим UPDATE не вийде — розшифрування впаде.
aad := *tenant + "|" + *kind
sec, err := ring.Encrypt([]byte(plaintext), aad)
if err != nil {
return err
}
tx, err := conn.Begin(ctx)
if err != nil {
return err
}
defer func() { _ = tx.Rollback(ctx) }()
var secretID string
if err := tx.QueryRow(ctx, `
INSERT INTO core.secrets (tenant_id, kind, key_id, nonce, ciphertext, auth_tag, aad)
VALUES ($1, $2::core.secret_kind, $3, $4, $5, $6, $7)
RETURNING id::text
`, *tenant, *kind, sec.KeyID, sec.Nonce, sec.Ciphertext, sec.AuthTag, aad).Scan(&secretID); err != nil {
return fmt.Errorf("запис секрету: %w", err)
}
fmt.Println("secret_id:", secretID)
if *credName != "" {
if *proto == "" {
return errors.New("для -credential потрібен -proto")
}
var credID string
if err := tx.QueryRow(ctx, `
INSERT INTO inv.credentials (tenant_id, name, proto, username, port, secret_id, options)
VALUES ($1, $2, $3::inv.credential_proto, NULLIF($4,''), NULLIF($5,0)::int, $6, $7::jsonb)
ON CONFLICT (tenant_id, name) DO UPDATE
SET secret_id = EXCLUDED.secret_id,
username = EXCLUDED.username,
port = EXCLUDED.port,
options = EXCLUDED.options,
updated_at = now()
RETURNING id::text
`, *tenant, *credName, *proto, *username, *port, secretID, *options).Scan(&credID); err != nil {
return fmt.Errorf("запис креденшела: %w", err)
}
fmt.Println("credential_id:", credID)
if *attach != "" {
if _, err := tx.Exec(ctx, `
INSERT INTO inv.device_credentials (device_id, credential_id, priority)
VALUES ($1, $2, 10)
ON CONFLICT (device_id, credential_id) DO NOTHING
`, *attach, credID); err != nil {
return fmt.Errorf("прив'язка до пристрою: %w", err)
}
fmt.Println("прив'язано до пристрою:", *attach)
}
}
return tx.Commit(ctx)
}
func readStdin() (string, error) {
r := bufio.NewReader(os.Stdin)
line, err := r.ReadString('\n')
if err != nil && line == "" {
return "", err
}
return strings.TrimRight(line, "\r\n"), nil
}
func buildKeyring(spec string) (*crypto.Keyring, error) {
ring := crypto.NewKeyring()
if strings.TrimSpace(spec) == "" {
return nil, errors.New("не вказано -dek")
}
for _, part := range strings.Split(spec, ",") {
part = strings.TrimSpace(part)
if part == "" {
continue
}
id, raw, ok := strings.Cut(part, "=")
if !ok {
return nil, fmt.Errorf("некоректний запис ключа: %q", part)
}
key, err := hex.DecodeString(raw)
if err != nil {
key, err = base64.StdEncoding.DecodeString(raw)
if err != nil {
return nil, fmt.Errorf("ключ %q: не hex і не base64", id)
}
}
if err := ring.Add(id, key); err != nil {
return nil, err
}
}
return ring, nil
}

View file

@ -130,7 +130,8 @@ func (s *Store) applyNeighbor(ctx context.Context, a *Agent, n *npv1.NeighborRec
VALUES ($1,$2,$3,$4::topo.discovery_proto,$5,$6,$7,$8,
$9::inet,$10::macaddr,$11,$12,$13,$14,$15,now(),$16::jsonb)
ON CONFLICT (device_id, COALESCE(interface_id, '00000000-0000-0000-0000-000000000000'::uuid),
proto, COALESCE(remote_chassis_id,''), COALESCE(remote_port_id,''))
proto, COALESCE(remote_chassis_id,''), COALESCE(remote_port_id,''),
COALESCE(remote_mac::text,''))
DO UPDATE SET
remote_system_name = EXCLUDED.remote_system_name,
remote_port_descr = EXCLUDED.remote_port_descr,