Етап 2: Go-агент — планувальник, буфер, сесія, модулі ICMP і SNMP

Модулі вкомпільовані (без .so): один бінарник на Alpine, Windows і роутер.
Креденшели беруться на момент виконання, бо мають TTL. Розклад вирівняний
по сітці інтервалу, тому переживає рестарт. Буфер обмежений і за кількістю,
і за пам'яттю; при переповненні викидає найстаріше, зміни статусу — останніми.

Перевірено на Debian 13 / Go 1.25: go vet чисто, go test -race усі пакети ok.
Релізний бінарник 12 МБ, базовий RSS 11.6 МБ при бюджеті 30 МБ.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
zotac 2026-08-14 03:46:46 +03:00
parent 82036a4592
commit c8075ee2a6
16 changed files with 4364 additions and 5 deletions

View file

@ -127,9 +127,72 @@ Debian 13 (192.168.1.203): protoc 3.21.12, Go 1.24.4, buf 1.72.0.
`RPC_REQUEST_STANDARD_NAME` та сусідніх правил — вони припускають пари
запит-відповідь, а `ControlUp`/`ControlDown` це незалежні потоки подій.
### Далі (Етап 2, частина 2)
---
Сам Go-агент: скелет із gRPC-клієнтом і реконектом, планувальник задач
зі `schedule_offset`, модулі ICMP і SNMP, збір сусідів LLDP/CDP/ARP,
буфер телеметрії з обмеженням пам'яті (бюджет RSS < 30 МБ).
Паралельно — серверна сторона `AgentService` поверх схеми з Етапу 1.
## 2026-08-14 — Етап 2 (частина 2): Go-агент
### Створено
```
netpulse/
├── .gitignore, .gitattributes (репозиторій: LF усюди, крім .ps1)
└── agent/
├── README.md будова, рішення, параметри чеків
├── go.mod replace → ../gen/go
├── cmd/netpulse-agent/ точка входу, GOMEMLIMIT, keepalive
└── internal/
├── config/ прапорці + NETPULSE_*, mTLS
├── module/ контракт модуля, реєстр, маршрутизація
├── telemetry/ interner.go (series_ref) + buffer.go
├── scheduler/ min-heap, семафор, schedule_offset
├── session/ gRPC-клієнт, реконект, ack, план задач
└── modules/icmp, /snmp
```
### Прийняті рішення
1. **Модулі вкомпільовані, без динамічного завантаження.** Один бінарник має
працювати на Alpine, Windows і роутері з musl. «Активація» = дозвіл сервера.
2. **Креденшели беруться на момент виконання**, не з плану: у них TTL, і
прострочені не віддаються взагалі — інакше агент заблокує обліковий запис
на половині комутаторів клієнта.
3. **Розклад вирівняний по сітці інтервалу**, тому після рестарту задача
повертається у свій слот, а не з'їжджає.
4. **Буфер викидає найстаріше.** Після відновлення зв'язку цінніший поточний
стан. Зміни статусу викидаються останніми й пролазять у батч першими.
5. **Швидкості інтерфейсів рахує агент** (знає фактичний інтервал); при
перевороті лічильника — `counter_reset` замість стрибка на терабіт.
6. **Один писар у контрольний стрім** — gRPC не допускає паралельних Send.
7. **`GOMEMLIMIT` 48 МБ у коді:** хай GC працює агресивніше, ніж OOM killer
осліпить моніторинг саме тоді, коли він потрібен.
### Перевірено на стенді
Debian 13, Go 1.25.13. `go vet` чисто, `go test ./... -race`усі пакети ok.
Релізний бінарник (`CGO_ENABLED=0 -trimpath -s -w`): **12 МБ**, базовий
**RSS 11.6 МБ** у циклі реконекту (бюджет 30 МБ).
### Не перевірено
- `TestPingLoopback` під звичайним користувачем в unprivileged LXC пропускається:
ядро не дає ані unprivileged-, ані raw-сокета, `sysctl ping_group_range`
недоступний. **Під root на тому ж стенді тест проходить** — ICMP-модуль
перевірений проти реального сокета. У проді потрібен `CAP_NET_RAW`.
- Модуль **snmp не перевірявся проти живого пристрою** — на стенді немає
SNMP-агента. Компілюється й проходить vet; логіка перевороту лічильників
і `util_pct` чекає на реальне обладнання.
### Репозиторій
`https://git.zotac.keenetic.link/zotac/Netpulse_SasS.git` (Forgejo).
Читання анонімне, **push вимагає токена** — Forgejo не пускає навіть у
публічний репозиторій без автентифікації (`Credentials are incorrect`).
Коміти лежать локально в `main` і чекають на токен.
### Далі (Етап 2, частина 3)
- Серверна сторона `AgentService` поверх схеми Етапу 1: запис телеметрії в
hypertables, резолвер `series_ref``ts.series.id`, планувальник, що
роздає `TaskPlan` із `schedule_offset`.
- Модуль topology на агенті: LLDP/CDP/ARP/FDB → `topo.neighbors`.
- `EnrollmentService`: видача сертифікатів зондам.

151
agent/README.md Normal file
View file

@ -0,0 +1,151 @@
# NetPulse Agent
Легкий зонд збору телеметрії. Єдиний статичний бінарник, усі з'єднання вихідні —
у мережі клієнта не треба відкривати жодного порту.
```bash
go build -trimpath -ldflags "-s -w" -o netpulse-agent ./cmd/netpulse-agent
```
```bash
./netpulse-agent -server monitor.example.com:443 -agent-id <uuid> -cert agent.pem -key agent.key
```
## Будова
```
cmd/netpulse-agent/ точка входу, збирання компонентів
internal/
config/ прапорці + NETPULSE_*, mTLS
module/ контракт модуля, реєстр, маршрутизація check_type → модуль
telemetry/
interner.go series_ref: серія реєструється раз за сесію
buffer.go обмежений буфер із бюджетом пам'яті
scheduler/ min-heap за часом запуску, семафор паралельності
session/ gRPC-клієнт, рукостискання, реконект, ack
modules/
icmp/ ping, RTT, jitter, втрати
snmp/ лічильники інтерфейсів (v2c/v3), довільні OID
```
## Рішення, які варто розуміти
**Модулі вкомпільовані, а не завантажуються динамічно.** Плагінність на агенті
навмисно простіша, ніж на сервері: жодних `.so`. Причина — той самий бінарник має
працювати на Alpine, Windows і роутері з musl, а бюджет RSS не дозволяє тягнути
рантайм плагінів. «Активація» означає, що сервер дозволив модуль цьому зонду
(`ControlDown.ModuleControl`).
**Креденшели беруться на момент виконання, а не з плану.** У них TTL. Якби вони
копіювались у задачу, після ротації паролів агент довбав би пристрої простроченими
даними, доки не приїде новий план — і заблокував би обліковий запис на половині
комутаторів. Прострочені не віддаються взагалі: явна помилка «немає креденшелів»
краща за блокування.
**Розклад вирівняний по сітці інтервалу, а не «зараз + інтервал».** Після рестарту
агента задача повертається у свій слот. `schedule_offset` рахує сервер
детерміновано від `check_id`, тому 5000 чеків з інтервалом 60 с розведені й
переживають перезапуск без збою фази.
**Буфер викидає найстаріші дані, а не найновіші.** Коли зв'язок відновиться,
оператору потрібен передусім поточний стан мережі. Зміни статусу пристрою
викидаються останніми: без них мапа показуватиме пристрій живим, поки він лежить.
Кількість викинутого їде в `AgentHealth.dropped_samples`, щоб діра була видимою.
**Швидкості інтерфейсів рахує агент.** Лише він знає фактичний інтервал між двома
опитуваннями; серверний розрахунок за номінальним інтервалом помиляється рівно на
мережеву затримку й джитер планувальника. При перевороті лічильника виставляється
`counter_reset`, і швидкості не заповнюються — краще діра в графіку, ніж стрибок
на терабіт і хибний алерт.
**Один писар у контрольний стрім.** gRPC не допускає паралельних `Send`, а слати
треба і heartbeat, і статуси задач. Тому окрема горутина-писар і канал. Читання
підтверджень телеметрії — на іншій горутині: `Send` + `Recv` на різних горутинах
це єдине, що gRPC дозволяє робити паралельно на одному стрімі.
**`GOMEMLIMIT` виставляється в коді (48 МБ).** Зонд часто живе на роутері або в
контейнері з 64 МБ. Хай збирач працює агресивніше, ніж OOM killer вб'є процес і
осліпить моніторинг саме тоді, коли він потрібен.
## Параметри
| Прапорець | Змінна | Призначення |
|-----------|--------|-------------|
| `-server` | `NETPULSE_SERVER` | адреса сервера `host:port` |
| `-agent-id` | `NETPULSE_AGENT_ID` | ідентифікатор зонда з `Enroll` |
| `-cert` / `-key` / `-ca` | `NETPULSE_CERT` / `_KEY` / `_CA` | mTLS |
| `-insecure` | `NETPULSE_INSECURE=1` | без TLS, лише локальний стенд |
| `-modules` | `NETPULSE_MODULES` | активні до першої відповіді сервера |
| `-max-concurrency` | `NETPULSE_MAX_CONCURRENCY` | скільки чеків паралельно |
| `-buffer-items` / `-buffer-bytes` | `NETPULSE_BUFFER_*` | ліміти буфера |
Перевірка сертифіката сервера не вимикається прапорцем: зонд ходить через
інтернет, і «тимчасово без перевірки» — це той самий канал, яким їдуть паролі від
усіх комутаторів клієнта.
## Параметри чеків
`icmp.ping`
```json
{"count": 3, "packet_size": 56, "interval_ms": 100}
```
`snmp.if` — перелік інтерфейсів приходить від сервера, бо він уже є в
`inv.interfaces`. Агент не ходить по `ifTable`, щоб з'ясувати, що існує:
виявлення нових інтерфейсів — робота модуля topology, і вона окрема саме тому,
що впирається в ліміт тарифу.
```json
{"use_hc_counters": true,
"interfaces": [{"if_index": 1, "interface_id": "<uuid>", "speed_bps": 1000000000}]}
```
`snmp.get`
```json
{"oids": [{"oid": "1.3.6.1.4.1.9.9.109.1.1.1.1.8.1",
"metric_key": "cpu.util", "unit": "pct", "scale": 1}]}
```
## Стан перевірки
Стенд Debian 13, Go 1.25.13. `go vet` чисто, тести з `-race`:
| Пакет | Результат |
|-------|-----------|
| `internal/scheduler` | ok |
| `internal/telemetry` | ok |
| `internal/session` | ok |
| `internal/modules/icmp` | ok |
Заміряно на релізному бінарнику (`CGO_ENABLED=0`, `-trimpath -s -w`):
- розмір — **12 МБ**
- базовий RSS у циклі реконекту — **11.6 МБ** (бюджет 30 МБ)
Що саме доводять тести:
| Тест | Що перевіряє |
|------|--------------|
| `TestNextRunSurvivesRestart` | слот задачі не з'їжджає після перезапуску агента |
| `TestOffsetsSpreadLoad` | 60 різних offset дають 60 різних моментів запуску |
| `TestSchedulerSkipsWhenPreviousStillRunning` | повільний чек не множить сам себе; пропуск явний |
| `TestSchedulerRejectsUnknownModule` | задача на неактивний модуль відхиляється зі слідом |
| `TestInternerLabelOrderDoesNotMatter` | недетермінований порядок map не плодить серій |
| `TestBufferDropsOldestOnOverflow` | при переповненні лишається найновіше, викинуте пораховане |
| `TestBufferPrioritisesStatusChanges` | зміна стану пролазить у батч поперед метрик |
| `TestBufferRequeuePreservesData` | розрив між Send і Ack не втрачає ні семпли, ні дескриптори |
| `TestSessionFullCycle` | Hello → план у зустрічному напрямку → виконання з креденшелами → телеметрія на сервері |
| `TestSessionReconnectsAfterDrop` | агент повертається сам після розриву |
| `TestPingLoopback` | реальний ICMP-сокет, не мок |
**Обмеження перевірки:** `TestPingLoopback` пропускається під звичайним
користувачем в unprivileged LXC — ядро не дає ані unprivileged-, ані raw-сокета,
і `sysctl net.ipv4.ping_group_range` там недоступний. Під root на тому ж стенді
тест проходить. У проді агенту треба `CAP_NET_RAW` або дозволений
`ping_group_range`.
Модуль `snmp` **не перевірявся проти живого пристрою** — на стенді немає SNMP-агента.
Він компілюється й проходить `vet`, але логіка перевороту лічильників і розрахунку
`util_pct` чекає на перевірку з реальним обладнанням.

14
agent/go.mod Normal file
View file

@ -0,0 +1,14 @@
module github.com/netpulse/netpulse/agent
go 1.24
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
google.golang.org/protobuf v1.36.12
)
// Згенерований із proto код лежить у репозиторії поруч.
replace github.com/netpulse/netpulse/gen/go => ../gen/go

48
agent/go.sum Normal file
View file

@ -0,0 +1,48 @@
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI=
github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY=
github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag=
github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE=
github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek=
github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps=
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/gosnmp/gosnmp v1.42.1 h1:MEJxhpC5v1coL3tFRix08PYmky9nyb1TLRRgJAmXm8A=
github.com/gosnmp/gosnmp v1.42.1/go.mod h1:CxVS6bXqmWZlafUj9pZUnQX5e4fAltqPcijxWpCitDo=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA=
github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64=
go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y=
go.opentelemetry.io/otel v1.44.0 h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU=
go.opentelemetry.io/otel v1.44.0/go.mod h1:BMgjTHL9WPRlRjL2oZCBTL4whCGtXch2H4BhOPIAyYc=
go.opentelemetry.io/otel/metric v1.44.0 h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSYCdDdmc=
go.opentelemetry.io/otel/metric v1.44.0/go.mod h1:8O7hanEPBNgEMmybD3s2VBKcgWOCsA6tzHBPODAiquo=
go.opentelemetry.io/otel/sdk v1.44.0 h1:nHYwb9lK+fJPU/dnT6s7W7Z8itMWyqrnVfbheVYrZ58=
go.opentelemetry.io/otel/sdk v1.44.0/go.mod h1:Osuydd3Se74nqjAKxid74N5eC+jfEqfTegHRnq58oK0=
go.opentelemetry.io/otel/sdk/metric v1.44.0 h1:3LlKgI+VjbVsjNRFZJZAJ30WjXC5VkNRks6si09iEfI=
go.opentelemetry.io/otel/sdk/metric v1.44.0/go.mod h1:5B5pMARnXxKhltooO4xUuCBorl65a4EpnTalObqOigA=
go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk=
go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE=
golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8=
golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww=
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY=
golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc=
golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38=
gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4=
gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E=
google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa h1:mZHHdPZl0dbGHCflZgAq/Q468DWVFcU2whhB2KAo8fk=
google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8=
google.golang.org/grpc v1.83.0 h1:JeNZEKJFbQxArAMl+hiytHauacDNqJUllNfmIMmpqnQ=
google.golang.org/grpc v1.83.0/go.mod h1:kDyl6SKsiHKt0uylY5gtn5cEjkrIOhQOGDgIc4JGwzQ=
google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc=
google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=

View file

@ -0,0 +1,157 @@
// Package config — налаштування агента з прапорців і змінних оточення.
//
// Конфіг-файлу навмисно немає: у зонда рівно жменя параметрів, а
// відсутність YAML економить залежність і кілька сотень кілобайт
// бінарника. Усе інше агент отримує від сервера в Welcome і ModuleControl.
package config
import (
"crypto/tls"
"crypto/x509"
"errors"
"flag"
"fmt"
"os"
"strconv"
"strings"
"time"
)
type Config struct {
// Куди підключатись: host:port.
Endpoint string
AgentID string
Hostname string
// mTLS. Порожні шляхи + Insecure = локальний стенд без шифрування.
CertFile string
KeyFile string
CAFile string
Insecure bool
MaxConcurrency int
BufferMaxItems int
BufferMaxBytes int
MinBackoff time.Duration
MaxBackoff time.Duration
LogLevel string
LogJSON bool
// Модулі, увімкнені до першої відповіді сервера.
DefaultModules []string
}
func envOr(key, def string) string {
if v := os.Getenv(key); v != "" {
return v
}
return def
}
func envInt(key string, def int) int {
if v := os.Getenv(key); v != "" {
if n, err := strconv.Atoi(v); err == nil {
return n
}
}
return def
}
// Parse читає прапорці поверх змінних оточення (NETPULSE_*).
func Parse(args []string) (*Config, error) {
host, _ := os.Hostname()
c := &Config{}
fs := flag.NewFlagSet("netpulse-agent", flag.ContinueOnError)
fs.StringVar(&c.Endpoint, "server", envOr("NETPULSE_SERVER", ""), "адреса сервера, host:port")
fs.StringVar(&c.AgentID, "agent-id", envOr("NETPULSE_AGENT_ID", ""), "ідентифікатор зонда")
fs.StringVar(&c.Hostname, "hostname", envOr("NETPULSE_HOSTNAME", host), "ім'я хоста для журналу сервера")
fs.StringVar(&c.CertFile, "cert", envOr("NETPULSE_CERT", ""), "клієнтський сертифікат (mTLS)")
fs.StringVar(&c.KeyFile, "key", envOr("NETPULSE_KEY", ""), "приватний ключ")
fs.StringVar(&c.CAFile, "ca", envOr("NETPULSE_CA", ""), "CA сервера")
fs.BoolVar(&c.Insecure, "insecure", os.Getenv("NETPULSE_INSECURE") == "1",
"без TLS — лише для локального стенду")
fs.IntVar(&c.MaxConcurrency, "max-concurrency", envInt("NETPULSE_MAX_CONCURRENCY", 64),
"скільки чеків виконувати паралельно")
fs.IntVar(&c.BufferMaxItems, "buffer-items", envInt("NETPULSE_BUFFER_ITEMS", 50_000),
"ліміт записів у буфері телеметрії")
fs.IntVar(&c.BufferMaxBytes, "buffer-bytes", envInt("NETPULSE_BUFFER_BYTES", 8<<20),
"оціночний ліміт пам'яті буфера в байтах")
fs.DurationVar(&c.MinBackoff, "min-backoff", time.Second, "мінімальна пауза перед реконектом")
fs.DurationVar(&c.MaxBackoff, "max-backoff", 2*time.Minute, "максимальна пауза перед реконектом")
fs.StringVar(&c.LogLevel, "log-level", envOr("NETPULSE_LOG_LEVEL", "info"), "debug|info|warn|error")
fs.BoolVar(&c.LogJSON, "log-json", os.Getenv("NETPULSE_LOG_JSON") == "1", "журнал у JSON")
mods := fs.String("modules", envOr("NETPULSE_MODULES", "icmp"),
"модулі, увімкнені до відповіді сервера (через кому)")
if err := fs.Parse(args); err != nil {
return nil, err
}
for _, m := range strings.Split(*mods, ",") {
if m = strings.TrimSpace(m); m != "" {
c.DefaultModules = append(c.DefaultModules, m)
}
}
return c, c.validate()
}
func (c *Config) validate() error {
if c.Endpoint == "" {
return errors.New("не вказано адресу сервера (-server або NETPULSE_SERVER)")
}
if c.AgentID == "" {
return errors.New("не вказано -agent-id: зонд має бути зареєстрований через Enroll")
}
if !c.Insecure && (c.CertFile == "" || c.KeyFile == "") {
return errors.New("потрібні -cert і -key для mTLS (або -insecure для локального стенду)")
}
if c.MaxConcurrency <= 0 {
return errors.New("-max-concurrency має бути додатним")
}
return nil
}
// TLSConfig збирає mTLS-конфігурацію.
//
// Перевірка сертифіката сервера обов'язкова й не вимикається прапорцем:
// зонд ходить через інтернет, і «тимчасово без перевірки» — це той
// самий канал, яким їдуть паролі від усіх комутаторів клієнта.
func (c *Config) TLSConfig() (*tls.Config, error) {
if c.Insecure {
return nil, nil
}
cert, err := tls.LoadX509KeyPair(c.CertFile, c.KeyFile)
if err != nil {
return nil, fmt.Errorf("клієнтський сертифікат: %w", err)
}
tc := &tls.Config{
Certificates: []tls.Certificate{cert},
MinVersion: tls.VersionTLS13,
}
if c.CAFile != "" {
pem, err := os.ReadFile(c.CAFile)
if err != nil {
return nil, fmt.Errorf("CA: %w", err)
}
pool := x509.NewCertPool()
if !pool.AppendCertsFromPEM(pem) {
return nil, errors.New("CA: не вдалося розібрати PEM")
}
tc.RootCAs = pool
}
return tc, nil
}

View file

@ -0,0 +1,261 @@
// Package module визначає контракт модуля опитування на боці агента.
//
// Плагінність на агенті влаштована простіше, ніж на сервері: динамічного
// завантаження .so немає навмисно. Причина — єдиний статичний бінарник,
// який має працювати на Alpine, Windows і роутері з musl, і бюджет
// RSS < 30 МБ. Тому модулі вкомпільовані, а "активація" означає, що
// сервер дозволив використання модуля цим зондом.
package module
import (
"context"
"encoding/json"
"fmt"
"sort"
"strings"
"sync"
"time"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
)
// Target — пристрій, який треба опитати.
type Target struct {
DeviceID string
Name string
Address string
}
// Task — одиниця роботи. Дзеркалить npv1.Task, але вже з розібраними
// полями: модулю не потрібно знати про protobuf.
type Task struct {
CheckID string
DeviceID string
InterfaceID string
// "<plugin>.<check>", напр. icmp.ping
CheckType string
// Непрозорий JSON: кожен модуль розбирає свою частину сам.
Params json.RawMessage
Target Target
Credentials []*npv1.Credential
Timeout time.Duration
}
// CheckName повертає частину після крапки: "icmp.ping" → "ping".
func (t Task) CheckName() string {
_, name, _ := strings.Cut(t.CheckType, ".")
return name
}
// ModuleKey повертає частину до крапки: "icmp.ping" → "icmp".
func ModuleKey(checkType string) string {
key, _, _ := strings.Cut(checkType, ".")
return key
}
// Metric — узагальнений вимір. Модуль не знає про series_ref: інтернування
// живе в межах сесії, а модуль може пережити кілька сесій.
type Metric struct {
MetricKey string
Unit string
Labels map[string]string
InterfaceID string
Value float64
Ts time.Time
}
// SeriesKey — канонічний ключ серії. Порядок labels зафіксовано сортуванням,
// інакше та сама серія отримувала б різні ref між запусками.
func (m Metric) SeriesKey(deviceID string) string {
var sb strings.Builder
sb.WriteString(deviceID)
sb.WriteByte(0)
sb.WriteString(m.InterfaceID)
sb.WriteByte(0)
sb.WriteString(m.MetricKey)
if len(m.Labels) > 0 {
keys := make([]string, 0, len(m.Labels))
for k := range m.Labels {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
sb.WriteByte(0)
sb.WriteString(k)
sb.WriteByte('=')
sb.WriteString(m.Labels[k])
}
}
return sb.String()
}
// Result — те, що модуль повертає після виконання задачі.
type Result struct {
Metrics []Metric
Icmp *npv1.IcmpResult
Interfaces []*npv1.InterfaceCounters
Neighbors []*npv1.NeighborRecord
// Непрозоре навантаження, яке не є метрикою (напр. http.status → тіло відповіді).
Payload []byte
}
// Module — виконавець одного або кількох типів чеків.
type Module interface {
// Key — ідентифікатор модуля: icmp, snmp, topology, ncm.
Key() string
// CheckTypes — які типи чеків модуль уміє виконувати.
CheckTypes() []string
// Run виконує задачу. Дедлайн передається через ctx; модуль
// зобов'язаний його поважати, інакше планувальник заб'ється.
Run(ctx context.Context, task Task) (Result, error)
// Close звільняє ресурси (пули з'єднань, сокети).
Close() error
}
// ---------------------------------------------------------------------
// Реєстр
// ---------------------------------------------------------------------
// Registry — набір вкомпільованих модулів і їхній стан активації.
//
// Активацію диктує сервер (ControlDown.ModuleControl). Незареєстрований
// або неактивний модуль означає, що задача буде відхилена з
// STATE_REJECTED — це краще, ніж мовчки її пропустити: у сервера
// лишається слід у core.checks.last_error.
type Registry struct {
mu sync.RWMutex
modules map[string]Module
active map[string]bool
// checkType → module key
routes map[string]string
}
func NewRegistry() *Registry {
return &Registry{
modules: make(map[string]Module),
active: make(map[string]bool),
routes: make(map[string]string),
}
}
// Register додає вкомпільований модуль. Модуль неактивний, доки сервер
// його не увімкне — крім випадку, коли активацію ще не отримано взагалі
// (див. SetActive/EnsureDefaults).
func (r *Registry) Register(m Module) error {
r.mu.Lock()
defer r.mu.Unlock()
key := m.Key()
if _, dup := r.modules[key]; dup {
return fmt.Errorf("модуль %q вже зареєстровано", key)
}
for _, ct := range m.CheckTypes() {
if ModuleKey(ct) != key {
return fmt.Errorf("модуль %q оголошує чужий тип чека %q", key, ct)
}
if owner, dup := r.routes[ct]; dup {
return fmt.Errorf("тип чека %q вже належить модулю %q", ct, owner)
}
}
r.modules[key] = m
for _, ct := range m.CheckTypes() {
r.routes[ct] = key
}
return nil
}
// Compiled повертає ключі всіх вкомпільованих модулів — рівно те, що
// агент оголошує в Hello.build.compiled_modules.
func (r *Registry) Compiled() []string {
r.mu.RLock()
defer r.mu.RUnlock()
keys := make([]string, 0, len(r.modules))
for k := range r.modules {
keys = append(keys, k)
}
sort.Strings(keys)
return keys
}
// SetActive застосовує рішення сервера. exclusive = вимкнути все,
// чого немає в списку.
func (r *Registry) SetActive(enabled map[string]bool, exclusive bool) {
r.mu.Lock()
defer r.mu.Unlock()
if exclusive {
for k := range r.active {
r.active[k] = false
}
}
for k, v := range enabled {
if _, known := r.modules[k]; known {
r.active[k] = v
}
}
}
// EnsureDefaults вмикає перелічені модулі, якщо сервер ще нічого не
// сказав. Потрібно для першого запуску до отримання ModuleControl.
func (r *Registry) EnsureDefaults(keys ...string) {
r.mu.Lock()
defer r.mu.Unlock()
if len(r.active) > 0 {
return
}
for _, k := range keys {
if _, known := r.modules[k]; known {
r.active[k] = true
}
}
}
func (r *Registry) IsActive(key string) bool {
r.mu.RLock()
defer r.mu.RUnlock()
return r.active[key]
}
// ErrNoModule — задача адресована модулю, якого немає або він не активний.
type ErrNoModule struct {
CheckType string
Reason string
}
func (e *ErrNoModule) Error() string {
return fmt.Sprintf("чек %q: %s", e.CheckType, e.Reason)
}
// Resolve знаходить модуль для типу чека.
func (r *Registry) Resolve(checkType string) (Module, error) {
r.mu.RLock()
defer r.mu.RUnlock()
key, ok := r.routes[checkType]
if !ok {
return nil, &ErrNoModule{CheckType: checkType, Reason: "невідомий тип чека"}
}
if !r.active[key] {
return nil, &ErrNoModule{CheckType: checkType, Reason: "модуль " + key + " не активований сервером"}
}
return r.modules[key], nil
}
// CloseAll закриває всі модулі. Помилки збираються, а не глушаться:
// зависле з'єднання при завершенні — теж діагностика.
func (r *Registry) CloseAll() []error {
r.mu.Lock()
defer r.mu.Unlock()
var errs []error
for key, m := range r.modules {
if err := m.Close(); err != nil {
errs = append(errs, fmt.Errorf("модуль %s: %w", key, err))
}
}
return errs
}

View file

@ -0,0 +1,365 @@
// Package icmp — модуль доступності: ping, RTT, jitter, втрати.
//
// Це єдиний модуль, доступний на тарифі Free, і найгарячіший шлях у
// системі: саме його результат фарбує вузол на мапі.
package icmp
import (
"context"
"crypto/rand"
"encoding/binary"
"errors"
"fmt"
"math"
"net"
"os"
"sync"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"golang.org/x/net/icmp"
"golang.org/x/net/ipv4"
"golang.org/x/net/ipv6"
"google.golang.org/protobuf/types/known/timestamppb"
"encoding/json"
)
// Params — вміст Task.params_json для icmp.ping.
type Params struct {
Count int `json:"count"`
PacketSize int `json:"packet_size"`
// Пауза між пакетами всередині одного заміру.
IntervalMs int `json:"interval_ms"`
}
func (p *Params) applyDefaults() {
if p.Count <= 0 {
p.Count = 3
}
if p.Count > 32 {
p.Count = 32
}
if p.PacketSize < 16 {
p.PacketSize = 56
}
if p.PacketSize > 1400 {
p.PacketSize = 1400
}
if p.IntervalMs <= 0 {
p.IntervalMs = 100
}
}
// Module — реалізація module.Module.
type Module struct {
mu sync.Mutex
seq uint32
closed bool
}
func New() *Module { return &Module{} }
func (m *Module) Key() string { return "icmp" }
func (m *Module) CheckTypes() []string { return []string{"icmp.ping"} }
func (m *Module) Close() error {
m.mu.Lock()
defer m.mu.Unlock()
m.closed = true
return nil
}
func (m *Module) nextSeq() int {
m.mu.Lock()
defer m.mu.Unlock()
m.seq++
return int(m.seq & 0xffff)
}
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 для icmp.ping: %w", err)
}
}
p.applyDefaults()
addr, err := resolve(ctx, task.Target.Address)
if err != nil {
return module.Result{}, err
}
rtts, sent, err := m.probe(ctx, addr, p)
// Помилка сокета — це збій чека, а не "пристрій лежить": ці випадки
// не можна плутати, інакше проблема з правами на зонді виглядатиме
// як аварія в мережі клієнта.
if err != nil {
return module.Result{}, err
}
now := time.Now()
res := computeStats(rtts, sent)
icmpResult := &npv1.IcmpResult{
DeviceId: task.DeviceID,
CheckId: task.CheckID,
Ts: timestamppb.New(now),
RttAvgMs: res.avg,
RttMinMs: res.min,
RttMaxMs: res.max,
JitterMs: res.jitter,
LossPct: res.loss,
PacketsSent: uint32(sent),
PacketsRecv: uint32(len(rtts)),
Reachable: len(rtts) > 0,
}
out := module.Result{Icmp: icmpResult}
// Дублюємо в узагальнені метрики: широка таблиця живить мапу,
// а series/samples — довільні графіки й правила алертів.
if len(rtts) > 0 {
out.Metrics = append(out.Metrics,
module.Metric{MetricKey: "icmp.rtt_avg", Unit: "ms", Value: float64(res.avg), Ts: now},
module.Metric{MetricKey: "icmp.jitter", Unit: "ms", Value: float64(res.jitter), Ts: now},
)
}
out.Metrics = append(out.Metrics,
module.Metric{MetricKey: "icmp.loss_pct", Unit: "pct", Value: float64(res.loss), Ts: now},
)
return out, nil
}
type stats struct {
avg, min, max, jitter, loss float32
}
func computeStats(rtts []time.Duration, sent int) stats {
var s stats
if sent <= 0 {
return s
}
if len(rtts) == 0 {
s.loss = 100
return s
}
minRTT, maxRTT := rtts[0], rtts[0]
var total time.Duration
for _, d := range rtts {
total += d
if d < minRTT {
minRTT = d
}
if d > maxRTT {
maxRTT = d
}
}
avg := total / time.Duration(len(rtts))
// Jitter — середнє абсолютне відхилення між сусідніми замірами
// (RFC 3550 у спрощеному вигляді). Саме воно, а не розкид min/max,
// показує нестабільність каналу.
var jitter float64
if len(rtts) > 1 {
var acc float64
for i := 1; i < len(rtts); i++ {
acc += math.Abs(float64(rtts[i]-rtts[i-1])) / float64(time.Millisecond)
}
jitter = acc / float64(len(rtts)-1)
}
s.avg = float32(float64(avg) / float64(time.Millisecond))
s.min = float32(float64(minRTT) / float64(time.Millisecond))
s.max = float32(float64(maxRTT) / float64(time.Millisecond))
s.jitter = float32(jitter)
s.loss = float32(100 * float64(sent-len(rtts)) / float64(sent))
return s
}
func resolve(ctx context.Context, address string) (*net.IPAddr, error) {
if address == "" {
return nil, errors.New("порожня адреса опитування")
}
// Адреса може прийти з маскою або портом — беремо лише IP-частину.
if host, _, err := net.SplitHostPort(address); err == nil {
address = host
}
if ip := net.ParseIP(address); ip != nil {
return &net.IPAddr{IP: ip}, nil
}
var r net.Resolver
ips, err := r.LookupIPAddr(ctx, address)
if err != nil {
return nil, fmt.Errorf("не резолвиться %q: %w", address, err)
}
if len(ips) == 0 {
return nil, fmt.Errorf("не резолвиться %q: порожня відповідь", address)
}
return &ips[0], nil
}
// probe шле p.Count пакетів і збирає RTT тих, що повернулись.
func (m *Module) probe(ctx context.Context, addr *net.IPAddr, p Params) ([]time.Duration, int, error) {
v6 := addr.IP.To4() == nil
conn, err := listen(v6)
if err != nil {
return nil, 0, err
}
defer conn.Close()
var (
echoType icmp.Type = ipv4.ICMPTypeEcho
proto = 1
)
if v6 {
echoType = ipv6.ICMPTypeEchoRequest
proto = 58
}
// Токен у корисному навантаженні: відсіює чужі відповіді, коли
// сокет неprivileged і ядро мультиплексує ID між процесами.
token := make([]byte, 8)
if _, err := rand.Read(token); err != nil {
return nil, 0, err
}
id := os.Getpid() & 0xffff
rtts := make([]time.Duration, 0, p.Count)
sent := 0
// Ціль для udp-сокета вимагає *net.UDPAddr, для raw — *net.IPAddr.
var dst net.Addr = addr
if _, isUDP := conn.LocalAddr().(*net.UDPAddr); isUDP {
dst = &net.UDPAddr{IP: addr.IP, Zone: addr.Zone}
}
interval := time.Duration(p.IntervalMs) * time.Millisecond
buf := make([]byte, 1500)
for i := 0; i < p.Count; i++ {
if err := ctx.Err(); err != nil {
break
}
seq := m.nextSeq()
payload := make([]byte, p.PacketSize)
copy(payload, token)
binary.BigEndian.PutUint32(payload[8:], uint32(seq))
msg := icmp.Message{
Type: echoType,
Code: 0,
Body: &icmp.Echo{ID: id, Seq: seq, Data: payload},
}
wire, err := msg.Marshal(nil)
if err != nil {
return nil, sent, err
}
start := time.Now()
if _, err := conn.WriteTo(wire, dst); err != nil {
// Мережа недоступна з самого зонда — це збій чека.
return nil, sent, fmt.Errorf("не вдалося відправити ICMP: %w", err)
}
sent++
deadline := start.Add(timeoutFor(ctx, interval))
if got, ok := awaitReply(conn, buf, proto, token, seq, deadline); ok {
rtts = append(rtts, got.Sub(start))
}
if i < p.Count-1 {
remaining := interval - time.Since(start)
if remaining > 0 {
select {
case <-ctx.Done():
return rtts, sent, nil
case <-time.After(remaining):
}
}
}
}
return rtts, sent, nil
}
func timeoutFor(ctx context.Context, fallback time.Duration) time.Duration {
if dl, ok := ctx.Deadline(); ok {
if d := time.Until(dl); d > 0 && d < fallback*4 {
return d
}
}
if fallback < 200*time.Millisecond {
return 200 * time.Millisecond
}
return fallback * 4
}
// awaitReply читає, доки не побачить свою відповідь або не спливе дедлайн.
func awaitReply(conn *icmp.PacketConn, buf []byte, proto int, token []byte, seq int, deadline time.Time) (time.Time, bool) {
for {
if time.Now().After(deadline) {
return time.Time{}, false
}
if err := conn.SetReadDeadline(deadline); err != nil {
return time.Time{}, false
}
n, _, err := conn.ReadFrom(buf)
if err != nil {
return time.Time{}, false
}
received := time.Now()
msg, err := icmp.ParseMessage(proto, buf[:n])
if err != nil {
continue
}
echo, ok := msg.Body.(*icmp.Echo)
if !ok {
continue
}
if echo.Seq != seq || len(echo.Data) < 8 {
continue
}
if string(echo.Data[:8]) != string(token) {
continue
}
return received, true
}
}
// listen відкриває ICMP-сокет.
//
// Спершу пробуємо unprivileged datagram-сокет (Linux:
// net.ipv4.ping_group_range, macOS: доступний завжди) — щоб агент не
// вимагав root. Якщо ядро не дозволяє, падаємо на raw-сокет, який
// потребує CAP_NET_RAW.
func listen(v6 bool) (*icmp.PacketConn, error) {
udpNet, rawNet, addr := "udp4", "ip4:icmp", "0.0.0.0"
if v6 {
udpNet, rawNet, addr = "udp6", "ip6:ipv6-icmp", "::"
}
conn, udpErr := icmp.ListenPacket(udpNet, addr)
if udpErr == nil {
return conn, nil
}
conn, rawErr := icmp.ListenPacket(rawNet, addr)
if rawErr == nil {
return conn, nil
}
return nil, fmt.Errorf(
"немає доступу до ICMP-сокета (unprivileged: %v; raw: %v); "+
"дайте бінарнику CAP_NET_RAW або дозвольте sysctl net.ipv4.ping_group_range",
udpErr, rawErr)
}

View file

@ -0,0 +1,141 @@
package icmp
import (
"context"
"encoding/json"
"strings"
"testing"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
)
func TestComputeStats(t *testing.T) {
rtts := []time.Duration{
10 * time.Millisecond,
12 * time.Millisecond,
11 * time.Millisecond,
}
s := computeStats(rtts, 4)
if s.min != 10 || s.max != 12 {
t.Fatalf("min/max = %v/%v", s.min, s.max)
}
if s.avg != 11 {
t.Fatalf("avg = %v", s.avg)
}
// Jitter — середнє відхилення між СУСІДНІМИ замірами: |12-10| і |11-12|.
if s.jitter != 1.5 {
t.Fatalf("jitter = %v, очікували 1.5", s.jitter)
}
if s.loss != 25 {
t.Fatalf("loss = %v, очікували 25%%", s.loss)
}
}
func TestComputeStatsAllLost(t *testing.T) {
s := computeStats(nil, 3)
if s.loss != 100 {
t.Fatalf("loss = %v, очікували 100", s.loss)
}
if s.avg != 0 {
t.Fatalf("avg при повній втраті = %v", s.avg)
}
}
func TestParamsDefaults(t *testing.T) {
var p Params
p.applyDefaults()
if p.Count != 3 || p.PacketSize != 56 || p.IntervalMs != 100 {
t.Fatalf("значення за замовчуванням: %+v", p)
}
// Захист від параметрів, які перетворили б зонд на генератор трафіку.
p = Params{Count: 10_000, PacketSize: 65_000}
p.applyDefaults()
if p.Count > 32 || p.PacketSize > 1400 {
t.Fatalf("ліміти не застосовані: %+v", p)
}
}
func TestResolveRejectsEmpty(t *testing.T) {
if _, err := resolve(context.Background(), ""); err == nil {
t.Fatal("порожня адреса прийнята")
}
}
func TestResolveStripsPort(t *testing.T) {
addr, err := resolve(context.Background(), "127.0.0.1:161")
if err != nil {
t.Fatalf("resolve: %v", err)
}
if addr.IP.String() != "127.0.0.1" {
t.Fatalf("IP = %v", addr.IP)
}
}
// Реальний ping до localhost. Пропускається там, де ядро не дає
// ані unprivileged-, ані raw-сокета (типова CI-пісочниця).
func TestPingLoopback(t *testing.T) {
m := New()
defer m.Close()
params, _ := json.Marshal(Params{Count: 2, PacketSize: 32, IntervalMs: 20})
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
res, err := m.Run(ctx, module.Task{
CheckID: "chk-1",
DeviceID: "dev-1",
CheckType: "icmp.ping",
Params: params,
Target: module.Target{DeviceID: "dev-1", Address: "127.0.0.1"},
Timeout: 3 * time.Second,
})
if err != nil {
if strings.Contains(err.Error(), "немає доступу до ICMP-сокета") {
t.Skip("середовище не дозволяє ICMP:", err)
}
t.Fatalf("ping: %v", err)
}
if res.Icmp == nil {
t.Fatal("немає IcmpResult")
}
if !res.Icmp.Reachable {
t.Fatalf("127.0.0.1 недосяжний: loss=%v", res.Icmp.LossPct)
}
if res.Icmp.PacketsSent != 2 {
t.Fatalf("відправлено %d пакетів", res.Icmp.PacketsSent)
}
if res.Icmp.PacketsRecv == 0 {
t.Fatal("жодної відповіді від localhost")
}
// Результат має продублюватись в узагальнені метрики: широка
// таблиця живить мапу, series/samples — графіки й правила алертів.
var haveLoss bool
for _, mt := range res.Metrics {
if mt.MetricKey == "icmp.loss_pct" {
haveLoss = true
}
}
if !haveLoss {
t.Fatalf("немає метрики icmp.loss_pct: %+v", res.Metrics)
}
}
func TestModuleContract(t *testing.T) {
m := New()
defer m.Close()
if m.Key() != "icmp" {
t.Fatalf("Key = %q", m.Key())
}
for _, ct := range m.CheckTypes() {
if module.ModuleKey(ct) != m.Key() {
t.Fatalf("тип чека %q не належить модулю %q", ct, m.Key())
}
}
}

View file

@ -0,0 +1,507 @@
// Package snmp — модуль опитування по SNMP v2c/v3.
//
// Два типи чеків:
// snmp.if — лічильники інтерфейсів (те, що живить анімацію трафіку)
// snmp.get — довільні OID → узагальнені метрики
//
// Свідоме рішення: агент НЕ ходить по ifTable, щоб з'ясувати, які
// інтерфейси існують. Перелік (ifIndex → interface_id → speed_bps)
// приходить у параметрах задачі, бо він уже є в inv.interfaces на
// сервері. Виявлення нових інтерфейсів — робота модуля topology, і
// вона окрема саме тому, що впирається в ліміт тарифу.
package snmp
import (
"context"
"encoding/json"
"fmt"
"math"
"math/big"
"sync"
"time"
"github.com/gosnmp/gosnmp"
"github.com/netpulse/netpulse/agent/internal/module"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
)
// Базові 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"
)
// ---------------------------------------------------------------------
// Параметри задач
// ---------------------------------------------------------------------
// IfParams — params_json для snmp.if.
type IfParams struct {
UseHCCounters *bool `json:"use_hc_counters"`
Interfaces []IfaceTarget `json:"interfaces"`
}
type IfaceTarget struct {
IfIndex int `json:"if_index"`
InterfaceID string `json:"interface_id"`
// Номінальна швидкість порту — знаменник для util_pct.
SpeedBps uint64 `json:"speed_bps"`
}
// GetParams — params_json для snmp.get.
type GetParams struct {
OIDs []OIDSpec `json:"oids"`
}
type OIDSpec struct {
OID string `json:"oid"`
MetricKey string `json:"metric_key"`
Unit string `json:"unit"`
Labels map[string]string `json:"labels"`
// Множник: сенсори часто віддають десяті градуса цілим числом.
Scale float64 `json:"scale"`
}
// ---------------------------------------------------------------------
// Стан лічильників
// ---------------------------------------------------------------------
// counterState — попередній замір, потрібний для розрахунку швидкості.
//
// Швидкості рахує агент, бо лише він знає фактичний інтервал між двома
// опитуваннями. Серверний розрахунок за номінальним інтервалом дає
// похибку рівно на мережеву затримку й джитер планувальника.
type counterState struct {
ts time.Time
inOctets uint64
outOctets uint64
inUcast uint64
outUcast uint64
valid bool
}
type Module struct {
mu sync.Mutex
state map[string]*counterState // deviceID|ifIndex
}
func New() *Module {
return &Module{state: make(map[string]*counterState)}
}
func (m *Module) Key() string { return "snmp" }
func (m *Module) CheckTypes() []string { return []string{"snmp.if", "snmp.get"} }
func (m *Module) Close() error {
m.mu.Lock()
defer m.mu.Unlock()
m.state = make(map[string]*counterState)
return nil
}
func (m *Module) Run(ctx context.Context, task module.Task) (module.Result, error) {
client, err := dial(ctx, task)
if err != nil {
return module.Result{}, err
}
defer client.Conn.Close()
switch task.CheckName() {
case "if":
return m.runInterfaces(ctx, client, task)
case "get":
return m.runGet(ctx, client, task)
default:
return module.Result{}, fmt.Errorf("snmp: невідомий чек %q", task.CheckType)
}
}
// ---------------------------------------------------------------------
// snmp.if
// ---------------------------------------------------------------------
func (m *Module) runInterfaces(ctx context.Context, client *gosnmp.GoSNMP, task module.Task) (module.Result, error) {
var p IfParams
if len(task.Params) > 0 {
if err := json.Unmarshal(task.Params, &p); err != nil {
return module.Result{}, fmt.Errorf("невалідні params для snmp.if: %w", err)
}
}
if len(p.Interfaces) == 0 {
return module.Result{}, fmt.Errorf("snmp.if без переліку інтерфейсів: сервер має передати ifIndex → interface_id")
}
useHC := true
if p.UseHCCounters != nil {
useHC = *p.UseHCCounters
}
inOID, outOID := oidIfHCInOctets, oidIfHCOutOctets
if !useHC {
inOID, outOID = oidIfInOctets, oidIfOutOctets
}
// Один Get на всі інтерфейси замість walk: менше пакетів і жодного
// зайвого трафіку по портах, які нас не цікавлять.
var oids []string
for _, ifc := range p.Interfaces {
oids = append(oids,
fmt.Sprintf("%s.%d", inOID, ifc.IfIndex),
fmt.Sprintf("%s.%d", outOID, ifc.IfIndex),
fmt.Sprintf("%s.%d", oidIfInErrors, ifc.IfIndex),
fmt.Sprintf("%s.%d", oidIfOutErrors, ifc.IfIndex),
fmt.Sprintf("%s.%d", oidIfInDiscards, ifc.IfIndex),
fmt.Sprintf("%s.%d", oidIfOutDiscards, ifc.IfIndex),
fmt.Sprintf("%s.%d", oidIfOperStatus, ifc.IfIndex),
fmt.Sprintf("%s.%d", oidIfAdminStatus, ifc.IfIndex),
)
if useHC {
oids = append(oids,
fmt.Sprintf("%s.%d", oidIfHCInUcast, ifc.IfIndex),
fmt.Sprintf("%s.%d", oidIfHCOutUcast, ifc.IfIndex),
)
}
}
values, err := getAll(ctx, client, oids)
if err != nil {
return module.Result{}, err
}
now := time.Now()
out := module.Result{}
for _, ifc := range p.Interfaces {
suffix := fmt.Sprintf(".%d", ifc.IfIndex)
inOct, okIn := values[inOID+suffix]
outOct, okOut := values[outOID+suffix]
if !okIn || !okOut {
// Порт зник (модуль вийнято зі стека) — це не помилка чека.
continue
}
counters := &npv1.InterfaceCounters{
DeviceId: task.DeviceID,
InterfaceId: ifc.InterfaceID,
Ts: timestamppb.New(now),
InOctets: inOct,
OutOctets: outOct,
InErrors: values[oidIfInErrors+suffix],
OutErrors: values[oidIfOutErrors+suffix],
InDiscards: values[oidIfInDiscards+suffix],
OutDiscards: values[oidIfOutDiscards+suffix],
OperUp: values[oidIfOperStatus+suffix] == 1,
AdminUp: values[oidIfAdminStatus+suffix] == 1,
}
if useHC {
counters.InUcastPkts = values[oidIfHCInUcast+suffix]
counters.OutUcastPkts = values[oidIfHCOutUcast+suffix]
}
key := task.DeviceID + "|" + fmt.Sprint(ifc.IfIndex)
m.mu.Lock()
prev := m.state[key]
if prev == nil {
prev = &counterState{}
m.state[key] = prev
}
prevCopy := *prev
prev.ts = now
prev.inOctets = inOct
prev.outOctets = outOct
prev.inUcast = counters.InUcastPkts
prev.outUcast = counters.OutUcastPkts
prev.valid = true
m.mu.Unlock()
if prevCopy.valid {
elapsed := now.Sub(prevCopy.ts)
counters.Interval = durationpb.New(elapsed)
if inOct < prevCopy.inOctets || outOct < prevCopy.outOctets {
// Пристрій перезавантажився або лічильник перевернувся.
// Швидкості цього разу недійсні — краще діра в графіку,
// ніж стрибок на терабіт, який підніме хибний алерт.
counters.CounterReset = true
} else if elapsed > 0 {
secs := elapsed.Seconds()
counters.InBps = float64(inOct-prevCopy.inOctets) * 8 / secs
counters.OutBps = float64(outOct-prevCopy.outOctets) * 8 / secs
if useHC {
counters.InPps = float64(counters.InUcastPkts-prevCopy.inUcast) / secs
counters.OutPps = float64(counters.OutUcastPkts-prevCopy.outUcast) / secs
}
if ifc.SpeedBps > 0 {
counters.UtilInPct = clampPct(counters.InBps / float64(ifc.SpeedBps) * 100)
counters.UtilOutPct = clampPct(counters.OutBps / float64(ifc.SpeedBps) * 100)
}
}
}
out.Interfaces = append(out.Interfaces, counters)
if counters.Interval != nil && !counters.CounterReset {
labels := map[string]string{"if_index": fmt.Sprint(ifc.IfIndex)}
out.Metrics = append(out.Metrics,
module.Metric{MetricKey: "if.in_bps", Unit: "bps", Value: counters.InBps,
InterfaceID: ifc.InterfaceID, Labels: labels, Ts: now},
module.Metric{MetricKey: "if.out_bps", Unit: "bps", Value: counters.OutBps,
InterfaceID: ifc.InterfaceID, Labels: labels, Ts: now},
)
}
}
if len(out.Interfaces) == 0 {
return module.Result{}, fmt.Errorf("жоден із %d інтерфейсів не відповів", len(p.Interfaces))
}
return out, nil
}
func clampPct(v float64) float32 {
if math.IsNaN(v) || math.IsInf(v, 0) || v < 0 {
return 0
}
// Трохи понад 100% буває при розсинхроні інтервалу — не ріжемо
// одразу, але захищаємось від абсурду.
if v > 1000 {
return 1000
}
return float32(v)
}
// ---------------------------------------------------------------------
// snmp.get
// ---------------------------------------------------------------------
func (m *Module) runGet(ctx context.Context, client *gosnmp.GoSNMP, task module.Task) (module.Result, error) {
var p GetParams
if len(task.Params) > 0 {
if err := json.Unmarshal(task.Params, &p); err != nil {
return module.Result{}, fmt.Errorf("невалідні params для snmp.get: %w", err)
}
}
if len(p.OIDs) == 0 {
return module.Result{}, fmt.Errorf("snmp.get без жодного OID")
}
oids := make([]string, 0, len(p.OIDs))
for _, spec := range p.OIDs {
oids = append(oids, normalizeOID(spec.OID))
}
values, err := getAll(ctx, client, oids)
if err != nil {
return module.Result{}, err
}
now := time.Now()
out := module.Result{}
for _, spec := range p.OIDs {
raw, ok := values[normalizeOID(spec.OID)]
if !ok {
continue
}
scale := spec.Scale
if scale == 0 {
scale = 1
}
key := spec.MetricKey
if key == "" {
key = "snmp" + normalizeOID(spec.OID)
}
out.Metrics = append(out.Metrics, module.Metric{
MetricKey: key,
Unit: spec.Unit,
Labels: spec.Labels,
InterfaceID: task.InterfaceID,
Value: float64(raw) * scale,
Ts: now,
})
}
if len(out.Metrics) == 0 {
return module.Result{}, fmt.Errorf("жоден з %d OID не повернув значення", len(p.OIDs))
}
return out, nil
}
func normalizeOID(oid string) string {
if oid == "" || oid[0] == '.' {
return oid
}
return "." + oid
}
// ---------------------------------------------------------------------
// Транспорт
// ---------------------------------------------------------------------
// 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
}
}

View file

@ -0,0 +1,445 @@
// Package scheduler виконує план задач, отриманий від сервера.
//
// Дві вимоги визначають будову:
// 1. 5000 чеків з інтервалом 60 с не мають стартувати одночасно —
// інакше зонд раз на хвилину викидає сплеск на всю мережу.
// Розведення дає schedule_offset, який рахує СЕРВЕР (детерміновано
// від check_id), щоб воно зберігалось між перезапусками агента.
// 2. Повільний пристрій не має гальмувати решту: кожна задача має
// власний дедлайн, а паралельність обмежена семафором.
package scheduler
import (
"container/heap"
"context"
"sync"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
)
// Sink приймає результати. Реалізується telemetry.Buffer.
type Sink interface {
Add(deviceID, pluginKey string, res module.Result)
AddCheckResult(cr *npv1.CheckResult)
}
// StatusFunc доповідає серверу про життєвий цикл задачі.
type StatusFunc func(*npv1.TaskStatusUpdate)
// CredentialFunc віддає облікові дані пристрою на момент виконання.
//
// Навмисно не поле в Task: креденшели мають TTL і оновлюються окремим
// повідомленням від сервера. Якби вони копіювались у план, після
// ротації паролів агент довбав би пристрої простроченими даними,
// доки не приїде новий план.
type CredentialFunc func(deviceID string) []*npv1.Credential
type entry struct {
task module.Task
interval time.Duration
offset time.Duration
nextRun time.Time
running bool
index int
}
// Scheduler — планувальник із min-heap за часом наступного запуску.
type Scheduler struct {
registry *module.Registry
sink Sink
onStatus StatusFunc
creds CredentialFunc
mu sync.Mutex
entries map[string]*entry
queue taskHeap
wake chan struct{}
paused bool
sem chan struct{}
nowFn func() time.Time
}
type Config struct {
Registry *module.Registry
Sink Sink
OnStatus StatusFunc
Credentials CredentialFunc
MaxConcurrency int
// Підміна годинника в тестах.
Now func() time.Time
}
func New(cfg Config) *Scheduler {
if cfg.MaxConcurrency <= 0 {
cfg.MaxConcurrency = 64
}
if cfg.Now == nil {
cfg.Now = time.Now
}
if cfg.OnStatus == nil {
cfg.OnStatus = func(*npv1.TaskStatusUpdate) {}
}
if cfg.Credentials == nil {
cfg.Credentials = func(string) []*npv1.Credential { return nil }
}
return &Scheduler{
registry: cfg.Registry,
sink: cfg.Sink,
onStatus: cfg.OnStatus,
creds: cfg.Credentials,
entries: make(map[string]*entry),
wake: make(chan struct{}, 1),
sem: make(chan struct{}, cfg.MaxConcurrency),
nowFn: cfg.Now,
}
}
// SetPaused зупиняє видачу задач, не втрачаючи розкладу.
// Сервер вмикає це на grace-періоді після несплати: канал живий,
// опитування стоїть.
func (s *Scheduler) SetPaused(v bool) {
s.mu.Lock()
s.paused = v
s.mu.Unlock()
s.wakeup()
}
// NextRun рахує момент наступного запуску так, щоб він був вирівняний
// по сітці інтервалу й зсунутий на offset.
//
// Вирівнювання по сітці (а не "зараз + інтервал") робить розклад
// відтворюваним: після рестарту агента задача повертається у свій
// слот, а не з'їжджає на випадковий момент.
func NextRun(now time.Time, interval, offset time.Duration) time.Time {
if interval <= 0 {
interval = time.Minute
}
if offset < 0 || offset >= interval {
offset = time.Duration(0)
}
next := now.Truncate(interval).Add(offset)
for !next.After(now) {
next = next.Add(interval)
}
return next
}
// ---------------------------------------------------------------------
// Застосування плану
// ---------------------------------------------------------------------
// Apply замінює весь план (ControlDown.TaskPlan).
func (s *Scheduler) Apply(tasks []module.Task, meta []TaskMeta) {
s.mu.Lock()
defer s.mu.Unlock()
// Зберігаємо стан running: задача, що зараз виконується, не має
// зникнути з-під ніг у горутини, яка її виконує.
old := s.entries
s.entries = make(map[string]*entry, len(tasks))
s.queue = nil
now := s.nowFn()
for i, t := range tasks {
m := meta[i]
e := &entry{task: t, interval: m.Interval, offset: m.Offset}
if prev, ok := old[t.CheckID]; ok {
e.running = prev.running
}
e.nextRun = NextRun(now, e.interval, e.offset)
s.entries[t.CheckID] = e
heap.Push(&s.queue, e)
}
s.wakeup()
}
// Upsert застосовує інкрементальну зміну (ControlDown.TaskDelta).
func (s *Scheduler) Upsert(tasks []module.Task, meta []TaskMeta) {
s.mu.Lock()
defer s.mu.Unlock()
now := s.nowFn()
for i, t := range tasks {
m := meta[i]
if e, ok := s.entries[t.CheckID]; ok {
e.task = t
// Розклад перераховуємо лише якщо він реально змінився,
// інакше редагування опису пристрою збивало б фазу.
if e.interval != m.Interval || e.offset != m.Offset {
e.interval, e.offset = m.Interval, m.Offset
e.nextRun = NextRun(now, e.interval, e.offset)
heap.Fix(&s.queue, e.index)
}
continue
}
e := &entry{task: t, interval: m.Interval, offset: m.Offset}
e.nextRun = NextRun(now, e.interval, e.offset)
s.entries[t.CheckID] = e
heap.Push(&s.queue, e)
}
s.wakeup()
}
// Remove знімає задачі з розкладу.
func (s *Scheduler) Remove(checkIDs []string) {
s.mu.Lock()
defer s.mu.Unlock()
for _, id := range checkIDs {
e, ok := s.entries[id]
if !ok {
continue
}
if e.index >= 0 && e.index < len(s.queue) {
heap.Remove(&s.queue, e.index)
}
delete(s.entries, id)
}
s.wakeup()
}
// TaskMeta — розкладова частина задачі.
type TaskMeta struct {
Interval time.Duration
Offset time.Duration
}
// Len — скільки задач у розкладі.
func (s *Scheduler) Len() int {
s.mu.Lock()
defer s.mu.Unlock()
return len(s.entries)
}
// Running — скільки задач виконується просто зараз.
func (s *Scheduler) Running() int {
s.mu.Lock()
defer s.mu.Unlock()
n := 0
for _, e := range s.entries {
if e.running {
n++
}
}
return n
}
func (s *Scheduler) wakeup() {
select {
case s.wake <- struct{}{}:
default:
}
}
// ---------------------------------------------------------------------
// Головний цикл
// ---------------------------------------------------------------------
// Run крутиться до скасування контексту.
func (s *Scheduler) Run(ctx context.Context) {
timer := time.NewTimer(time.Hour)
defer timer.Stop()
var wg sync.WaitGroup
defer wg.Wait()
for {
now := s.nowFn()
due, wait := s.popDue(now)
for _, e := range due {
wg.Add(1)
go func(e *entry) {
defer wg.Done()
s.execute(ctx, e)
}(e)
}
if !timer.Stop() {
select {
case <-timer.C:
default:
}
}
timer.Reset(wait)
select {
case <-ctx.Done():
return
case <-timer.C:
case <-s.wake:
}
}
}
// popDue дістає задачі, яким час виконуватись, і повертає час до
// наступної. Задача, попередній запуск якої ще триває, пропускається
// з явним STATE_SKIPPED — тиха втрата циклу гірша за видиму.
func (s *Scheduler) popDue(now time.Time) ([]*entry, time.Duration) {
s.mu.Lock()
defer s.mu.Unlock()
if s.paused {
return nil, time.Second
}
var due []*entry
for s.queue.Len() > 0 {
top := s.queue[0]
if top.nextRun.After(now) {
break
}
heap.Pop(&s.queue)
top.nextRun = NextRun(now, top.interval, top.offset)
heap.Push(&s.queue, top)
if top.running {
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: top.task.CheckID,
State: npv1.TaskStatusUpdate_STATE_SKIPPED,
Ts: timestamppb.New(now),
Error: &npv1.Error{
Code: "still_running",
Message: "попередній запуск не завершився до наступного інтервалу",
Retryable: true,
},
})
continue
}
top.running = true
due = append(due, top)
}
wait := time.Hour
if s.queue.Len() > 0 {
if d := s.queue[0].nextRun.Sub(now); d < wait {
wait = d
}
}
if wait < time.Millisecond {
wait = time.Millisecond
}
return due, wait
}
func (s *Scheduler) execute(ctx context.Context, e *entry) {
defer func() {
s.mu.Lock()
e.running = false
s.mu.Unlock()
}()
task := e.task
// Креденшели беремо на момент виконання, а не з плану: у них TTL.
task.Credentials = s.creds(task.DeviceID)
mod, err := s.registry.Resolve(task.CheckType)
if err != nil {
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: task.CheckID,
State: npv1.TaskStatusUpdate_STATE_REJECTED,
Ts: timestamppb.Now(),
Error: &npv1.Error{Code: "no_module", Message: err.Error()},
})
return
}
// Семафор обмежує паралельність. Чекаємо на слот, але не довше,
// ніж дозволяє власний таймаут задачі.
select {
case s.sem <- struct{}{}:
defer func() { <-s.sem }()
case <-ctx.Done():
return
case <-time.After(task.Timeout):
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: task.CheckID,
State: npv1.TaskStatusUpdate_STATE_SKIPPED,
Ts: timestamppb.Now(),
Error: &npv1.Error{
Code: "concurrency_limit",
Message: "не дочекались вільного слота виконання",
Retryable: true,
},
})
return
}
timeout := task.Timeout
if timeout <= 0 {
timeout = 10 * time.Second
}
runCtx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
started := time.Now()
res, runErr := mod.Run(runCtx, task)
elapsed := time.Since(started)
cr := &npv1.CheckResult{
CheckId: task.CheckID,
DeviceId: task.DeviceID,
CheckType: task.CheckType,
Ts: timestamppb.New(started),
Duration: durationpb.New(elapsed),
Success: runErr == nil,
PayloadJson: res.Payload,
}
if runErr != nil {
cr.Error = &npv1.Error{
Code: classify(runErr, runCtx),
Message: runErr.Error(),
Retryable: true,
}
s.sink.AddCheckResult(cr)
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: task.CheckID,
State: npv1.TaskStatusUpdate_STATE_FAILED,
Ts: timestamppb.Now(),
Error: cr.Error,
})
return
}
s.sink.Add(task.DeviceID, module.ModuleKey(task.CheckType), res)
s.sink.AddCheckResult(cr)
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: task.CheckID,
State: npv1.TaskStatusUpdate_STATE_SUCCEEDED,
Ts: timestamppb.Now(),
})
}
func classify(err error, ctx context.Context) string {
if ctx.Err() == context.DeadlineExceeded {
return "timeout"
}
return "check_failed"
}
// ---------------------------------------------------------------------
// heap
// ---------------------------------------------------------------------
type taskHeap []*entry
func (h taskHeap) Len() int { return len(h) }
func (h taskHeap) Less(i, j int) bool { return h[i].nextRun.Before(h[j].nextRun) }
func (h taskHeap) Swap(i, j int) { h[i], h[j] = h[j], h[i]; h[i].index = i; h[j].index = j }
func (h *taskHeap) Push(x any) { e := x.(*entry); e.index = len(*h); *h = append(*h, e) }
func (h *taskHeap) Pop() any {
old := *h
n := len(old)
e := old[n-1]
old[n-1] = nil
e.index = -1
*h = old[:n-1]
return e
}

View file

@ -0,0 +1,361 @@
package scheduler
import (
"context"
"errors"
"sync"
"sync/atomic"
"testing"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
)
// ---------------------------------------------------------------------
// NextRun — розклад має бути відтворюваним
// ---------------------------------------------------------------------
func TestNextRunAlignsToGrid(t *testing.T) {
base := time.Date(2026, 8, 14, 12, 0, 3, 0, time.UTC)
got := NextRun(base, time.Minute, 7*time.Second)
want := time.Date(2026, 8, 14, 12, 0, 7, 0, time.UTC)
if !got.Equal(want) {
t.Fatalf("NextRun = %v, очікували %v", got, want)
}
// Момент уже минув у цій хвилині — переходимо в наступну.
got = NextRun(base.Add(10*time.Second), time.Minute, 7*time.Second)
want = time.Date(2026, 8, 14, 12, 1, 7, 0, time.UTC)
if !got.Equal(want) {
t.Fatalf("NextRun = %v, очікували %v", got, want)
}
}
// Головне, заради чого існує schedule_offset: після рестарту агента
// задача має повернутись у свій слот, а не з'їхати на випадковий момент.
func TestNextRunSurvivesRestart(t *testing.T) {
interval, offset := time.Minute, 23*time.Second
beforeRestart := NextRun(time.Date(2026, 8, 14, 12, 0, 5, 0, time.UTC), interval, offset)
afterRestart := NextRun(time.Date(2026, 8, 14, 12, 0, 19, 0, time.UTC), interval, offset)
if !beforeRestart.Equal(afterRestart) {
t.Fatalf("слот з'їхав після рестарту: %v != %v", beforeRestart, afterRestart)
}
}
// 5000 чеків з однаковим інтервалом не повинні стартувати одночасно.
func TestOffsetsSpreadLoad(t *testing.T) {
now := time.Date(2026, 8, 14, 12, 0, 0, 0, time.UTC)
interval := time.Minute
seen := make(map[time.Time]int)
for i := 0; i < 60; i++ {
offset := time.Duration(i) * time.Second
seen[NextRun(now, interval, offset)]++
}
if len(seen) != 60 {
t.Fatalf("60 різних offset дали лише %d різних моментів запуску", len(seen))
}
}
func TestNextRunRejectsOutOfRangeOffset(t *testing.T) {
now := time.Date(2026, 8, 14, 12, 0, 0, 0, time.UTC)
// Offset більший за інтервал — помилка сервера; не має ламати розклад.
got := NextRun(now, time.Minute, 90*time.Second)
if got.Sub(now) != time.Minute {
t.Fatalf("некоректний offset зіпсував розклад: %v", got.Sub(now))
}
}
// ---------------------------------------------------------------------
// Виконання
// ---------------------------------------------------------------------
type stubModule struct {
key string
runs atomic.Int64
delay time.Duration
failErr error
seen chan module.Task
}
func newStub(key string) *stubModule {
return &stubModule{key: key, seen: make(chan module.Task, 64)}
}
func (m *stubModule) Key() string { return m.key }
func (m *stubModule) CheckTypes() []string { return []string{m.key + ".probe"} }
func (m *stubModule) Close() error { return nil }
func (m *stubModule) Run(ctx context.Context, task module.Task) (module.Result, error) {
m.runs.Add(1)
select {
case m.seen <- task:
default:
}
if m.delay > 0 {
select {
case <-ctx.Done():
return module.Result{}, ctx.Err()
case <-time.After(m.delay):
}
}
if m.failErr != nil {
return module.Result{}, m.failErr
}
return module.Result{
Metrics: []module.Metric{{MetricKey: "stub.value", Unit: "pct", Value: 42}},
}, nil
}
type recordingSink struct {
mu sync.Mutex
results []module.Result
checks []*npv1.CheckResult
}
func (s *recordingSink) Add(deviceID, pluginKey string, res module.Result) {
s.mu.Lock()
defer s.mu.Unlock()
s.results = append(s.results, res)
}
func (s *recordingSink) AddCheckResult(cr *npv1.CheckResult) {
s.mu.Lock()
defer s.mu.Unlock()
s.checks = append(s.checks, cr)
}
func (s *recordingSink) counts() (int, int) {
s.mu.Lock()
defer s.mu.Unlock()
return len(s.results), len(s.checks)
}
func newHarness(t *testing.T, mod module.Module) (*Scheduler, *recordingSink, chan *npv1.TaskStatusUpdate) {
t.Helper()
reg := module.NewRegistry()
if err := reg.Register(mod); err != nil {
t.Fatalf("Register: %v", err)
}
reg.EnsureDefaults(mod.Key())
sink := &recordingSink{}
statuses := make(chan *npv1.TaskStatusUpdate, 128)
s := New(Config{
Registry: reg,
Sink: sink,
OnStatus: func(u *npv1.TaskStatusUpdate) {
select {
case statuses <- u:
default:
}
},
MaxConcurrency: 4,
Credentials: func(deviceID string) []*npv1.Credential {
return []*npv1.Credential{{CredentialId: "cred-" + deviceID}}
},
})
return s, sink, statuses
}
func TestSchedulerRunsTaskAndFillsCredentials(t *testing.T) {
mod := newStub("stub")
s, sink, statuses := newHarness(t, mod)
s.Apply(
[]module.Task{{
CheckID: "chk-1",
DeviceID: "dev-1",
CheckType: "stub.probe",
Target: module.Target{DeviceID: "dev-1", Address: "10.0.0.1"},
Timeout: time.Second,
}},
[]TaskMeta{{Interval: 50 * time.Millisecond}},
)
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
go s.Run(ctx)
select {
case task := <-mod.seen:
// Креденшели підставляються на момент виконання, а не з плану:
// у них TTL, і в плані вони б протухли.
if len(task.Credentials) != 1 || task.Credentials[0].CredentialId != "cred-dev-1" {
t.Fatalf("креденшели не підставились: %+v", task.Credentials)
}
case <-ctx.Done():
t.Fatal("задача жодного разу не виконалась")
}
// Дочекаємось, поки результат осяде в sink.
deadline := time.Now().Add(2 * time.Second)
for time.Now().Before(deadline) {
if r, c := sink.counts(); r > 0 && c > 0 {
break
}
time.Sleep(10 * time.Millisecond)
}
if r, c := sink.counts(); r == 0 || c == 0 {
t.Fatalf("результати не дійшли: results=%d checks=%d", r, c)
}
var gotSucceeded bool
for len(statuses) > 0 {
if (<-statuses).State == npv1.TaskStatusUpdate_STATE_SUCCEEDED {
gotSucceeded = true
}
}
if !gotSucceeded {
t.Fatal("сервер не отримав STATE_SUCCEEDED")
}
}
// Повільна задача не має накопичувати паралельні запуски: наступний цикл
// пропускається з явним STATE_SKIPPED.
func TestSchedulerSkipsWhenPreviousStillRunning(t *testing.T) {
mod := newStub("slow")
mod.delay = 400 * time.Millisecond
s, _, statuses := newHarness(t, mod)
s.Apply(
[]module.Task{{
CheckID: "chk-slow",
DeviceID: "dev-1",
CheckType: "slow.probe",
Timeout: 2 * time.Second,
}},
[]TaskMeta{{Interval: 50 * time.Millisecond}},
)
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
go s.Run(ctx)
deadline := time.Now().Add(1500 * time.Millisecond)
var skipped bool
for time.Now().Before(deadline) && !skipped {
select {
case u := <-statuses:
if u.State == npv1.TaskStatusUpdate_STATE_SKIPPED {
if u.Error == nil || u.Error.Code != "still_running" {
t.Fatalf("пропуск без зрозумілої причини: %+v", u.Error)
}
skipped = true
}
case <-time.After(20 * time.Millisecond):
}
}
if !skipped {
t.Fatal("повільна задача запускалась паралельно сама з собою")
}
// І при цьому вона таки виконувалась, а не зависла назавжди.
if mod.runs.Load() == 0 {
t.Fatal("задача не виконалась жодного разу")
}
}
// Задача на невідомий/неактивний модуль має бути відхилена явно —
// у сервера мусить лишитись слід у core.checks.last_error.
func TestSchedulerRejectsUnknownModule(t *testing.T) {
mod := newStub("stub")
s, _, statuses := newHarness(t, mod)
s.Apply(
[]module.Task{{
CheckID: "chk-x",
DeviceID: "dev-1",
CheckType: "netflow.probe",
Timeout: time.Second,
}},
[]TaskMeta{{Interval: 30 * time.Millisecond}},
)
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
go s.Run(ctx)
select {
case u := <-statuses:
if u.State != npv1.TaskStatusUpdate_STATE_REJECTED {
t.Fatalf("стан = %v, очікували REJECTED", u.State)
}
if u.Error == nil || u.Error.Code != "no_module" {
t.Fatalf("немає діагностики: %+v", u.Error)
}
case <-ctx.Done():
t.Fatal("невідомий тип чека проковтнуто мовчки")
}
}
func TestSchedulerReportsFailure(t *testing.T) {
mod := newStub("bad")
mod.failErr = errors.New("пристрій не відповідає")
s, sink, statuses := newHarness(t, mod)
s.Apply(
[]module.Task{{CheckID: "chk-f", DeviceID: "dev-1", CheckType: "bad.probe", Timeout: time.Second}},
[]TaskMeta{{Interval: 30 * time.Millisecond}},
)
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
go s.Run(ctx)
select {
case u := <-statuses:
if u.State != npv1.TaskStatusUpdate_STATE_FAILED {
t.Fatalf("стан = %v, очікували FAILED", u.State)
}
case <-ctx.Done():
t.Fatal("помилка виконання не доповіли")
}
// Невдалий чек теж має лишити слід у телеметрії — інакше
// «нічого не приходить» не відрізниш від «усе гаразд».
deadline := time.Now().Add(time.Second)
for time.Now().Before(deadline) {
if _, c := sink.counts(); c > 0 {
return
}
time.Sleep(10 * time.Millisecond)
}
t.Fatal("CheckResult про помилку не потрапив у буфер")
}
func TestSchedulerRemoveAndPause(t *testing.T) {
mod := newStub("stub")
s, _, _ := newHarness(t, mod)
s.Apply(
[]module.Task{{CheckID: "chk-1", DeviceID: "dev-1", CheckType: "stub.probe", Timeout: time.Second}},
[]TaskMeta{{Interval: 20 * time.Millisecond}},
)
if s.Len() != 1 {
t.Fatalf("у розкладі %d задач", s.Len())
}
s.SetPaused(true)
ctx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond)
defer cancel()
go s.Run(ctx)
time.Sleep(200 * time.Millisecond)
if n := mod.runs.Load(); n != 0 {
t.Fatalf("на паузі виконалось %d задач", n)
}
s.Remove([]string{"chk-1"})
if s.Len() != 0 {
t.Fatalf("після Remove лишилось %d задач", s.Len())
}
}

View file

@ -0,0 +1,657 @@
// Package session тримає з'єднання з сервером і все, що з нього випливає:
// рукостискання, застосування плану задач, відправку телеметрії,
// підтвердження, реконект.
//
// Усе з'єднання ініціює агент — у мережі клієнта немає ані відкритих
// портів, ані прокидання NAT. Тому «команда з сервера» технічно є
// повідомленням у зустрічному напрямку вже відкритого стріму Control.
package session
import (
"context"
"errors"
"fmt"
"log/slog"
"math/rand"
"runtime"
"sync"
"sync/atomic"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
"github.com/netpulse/netpulse/agent/internal/scheduler"
"github.com/netpulse/netpulse/agent/internal/telemetry"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/grpc"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
)
// Conn — те, що вміє *grpc.ClientConn. Винесено в інтерфейс, щоб тести
// підставляли bufconn без справжньої мережі.
type Conn interface {
grpc.ClientConnInterface
Close() error
}
// DialFunc відкриває з'єднання до сервера.
type DialFunc func(ctx context.Context) (Conn, error)
type Config struct {
AgentID string
Hostname string
Build *npv1.AgentBuild
Registry *module.Registry
Buffer *telemetry.Buffer
Scheduler *scheduler.Scheduler
Dial DialFunc
Logger *slog.Logger
MinBackoff time.Duration
MaxBackoff time.Duration
// Модулі, увімкнені до першої відповіді сервера.
DefaultModules []string
}
type Session struct {
cfg Config
log *slog.Logger
// Статуси задач — діагностика. Канал з буфером і скиданням при
// переповненні: краще втратити рядок журналу, ніж загальмувати
// опитування через повільний контрольний канал.
statusCh chan *npv1.TaskStatusUpdate
credMu sync.RWMutex
creds map[string][]*npv1.Credential
credExpiry time.Time
devMu sync.RWMutex
devices map[string]*npv1.DeviceTarget
planHash atomic.Pointer[[]byte]
nextBatch atomic.Uint64
lastAcked atomic.Uint64
skewNanos atomic.Int64
connected atomic.Bool
seq atomic.Uint64
startedAt time.Time
}
func New(cfg Config) *Session {
if cfg.Logger == nil {
cfg.Logger = slog.Default()
}
if cfg.MinBackoff <= 0 {
cfg.MinBackoff = time.Second
}
if cfg.MaxBackoff <= 0 {
cfg.MaxBackoff = 2 * time.Minute
}
return &Session{
cfg: cfg,
log: cfg.Logger,
statusCh: make(chan *npv1.TaskStatusUpdate, 512),
creds: make(map[string][]*npv1.Credential),
devices: make(map[string]*npv1.DeviceTarget),
startedAt: time.Now(),
}
}
// SetScheduler замикає взаємну залежність: планувальник потребує
// колбеків сесії (статуси, креденшели), а сесія — планувальника, щоб
// застосовувати план. Конструктором це не виражається, тому окремий сетер.
// Викликати до Run.
func (s *Session) SetScheduler(sc *scheduler.Scheduler) { s.cfg.Scheduler = sc }
// Credentials — провайдер для планувальника.
func (s *Session) Credentials(deviceID string) []*npv1.Credential {
s.credMu.RLock()
defer s.credMu.RUnlock()
// Прострочені креденшели не віддаємо: краще явна помилка
// "немає креденшелів", ніж заблокований обліковий запис на
// половині комутаторів клієнта.
if !s.credExpiry.IsZero() && time.Now().After(s.credExpiry) {
return nil
}
return s.creds[deviceID]
}
// ReportStatus — колбек для планувальника.
func (s *Session) ReportStatus(u *npv1.TaskStatusUpdate) {
select {
case s.statusCh <- u:
default:
}
}
// Connected — чи є жива сесія (для /healthz і самометрик).
func (s *Session) Connected() bool { return s.connected.Load() }
// ClockSkew — поправка годинника відносно сервера.
func (s *Session) ClockSkew() time.Duration {
return time.Duration(s.skewNanos.Load())
}
// ---------------------------------------------------------------------
// Цикл підключення
// ---------------------------------------------------------------------
// Run тримає з'єднання, доки не скасують контекст.
func (s *Session) Run(ctx context.Context) error {
backoff := s.cfg.MinBackoff
for {
if ctx.Err() != nil {
return ctx.Err()
}
err := s.runOnce(ctx)
s.connected.Store(false)
if ctx.Err() != nil {
return ctx.Err()
}
if err != nil {
s.log.Warn("сесію розірвано", "err", err, "retry_in", backoff)
}
// Джитер — щоб сотня зондів не ломанулась перепідключатись
// одночасно після рестарту сервера.
jitter := time.Duration(rand.Int63n(int64(backoff/2 + 1)))
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(backoff + jitter):
}
backoff *= 2
if backoff > s.cfg.MaxBackoff {
backoff = s.cfg.MaxBackoff
}
}
}
func (s *Session) runOnce(ctx context.Context) error {
conn, err := s.cfg.Dial(ctx)
if err != nil {
return fmt.Errorf("dial: %w", err)
}
defer conn.Close()
client := npv1.NewAgentServiceClient(conn)
sctx, cancel := context.WithCancel(ctx)
defer cancel()
ctrl, err := client.Control(sctx)
if err != nil {
return fmt.Errorf("відкрити Control: %w", err)
}
if err := ctrl.Send(s.hello()); err != nil {
return fmt.Errorf("Hello: %w", err)
}
first, err := ctrl.Recv()
if err != nil {
return fmt.Errorf("очікування Welcome: %w", err)
}
welcome := first.GetWelcome()
if welcome == nil {
return fmt.Errorf("замість Welcome прийшло %T", first.Payload)
}
s.applyWelcome(welcome)
s.connected.Store(true)
s.log.Info("сесію встановлено",
"session_id", welcome.SessionId,
"clock_skew", s.ClockSkew())
// Єдиний писар у контрольний стрім: gRPC не допускає паралельних
// Send, а слати треба і heartbeat, і статуси задач.
out := make(chan *npv1.ControlUp, 256)
var wg sync.WaitGroup
errCh := make(chan error, 4)
spawn := func(name string, fn func() error) {
wg.Add(1)
go func() {
defer wg.Done()
if err := fn(); err != nil && !errors.Is(err, context.Canceled) {
errCh <- fmt.Errorf("%s: %w", name, err)
}
cancel()
}()
}
spawn("writer", func() error { return s.writerLoop(sctx, ctrl, out) })
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) })
// Читання команд — у цій же горутині.
readErr := s.controlLoop(sctx, ctrl, out)
cancel()
wg.Wait()
close(errCh)
if readErr != nil && !errors.Is(readErr, context.Canceled) {
return readErr
}
for e := range errCh {
if e != nil {
return e
}
}
return nil
}
func (s *Session) hello() *npv1.ControlUp {
var hash []byte
if p := s.planHash.Load(); p != nil {
hash = *p
}
return &npv1.ControlUp{
Seq: s.seq.Add(1),
Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{
AgentId: s.cfg.AgentID,
Build: s.cfg.Build,
Hostname: s.cfg.Hostname,
// Хеш плану дозволяє серверу не перезаливати 50 000 задач
// після кожного обриву зв'язку.
TaskPlanHash: hash,
// Продовжуємо з місця розриву, а не з нуля.
LastAckedBatchId: s.lastAcked.Load(),
StartedAt: timestamppb.New(s.startedAt),
}},
}
}
func (s *Session) applyWelcome(w *npv1.Welcome) {
if w.ServerTime != nil {
s.skewNanos.Store(int64(time.Until(w.ServerTime.AsTime())))
}
s.cfg.Registry.EnsureDefaults(s.cfg.DefaultModules...)
}
// ---------------------------------------------------------------------
// Цикли
// ---------------------------------------------------------------------
func (s *Session) writerLoop(ctx context.Context, ctrl npv1.AgentService_ControlClient, out <-chan *npv1.ControlUp) error {
for {
select {
case <-ctx.Done():
return nil
case msg := <-out:
if err := ctrl.Send(msg); err != nil {
return err
}
}
}
}
func (s *Session) enqueue(ctx context.Context, out chan<- *npv1.ControlUp, msg *npv1.ControlUp) {
msg.Seq = s.seq.Add(1)
select {
case out <- msg:
case <-ctx.Done():
default:
// Черга забита — контрольний канал не встигає. Втратити
// heartbeat не страшно: сервер помітить це за таймаутом.
}
}
func (s *Session) heartbeatLoop(ctx context.Context, out chan<- *npv1.ControlUp, w *npv1.Welcome) error {
interval := 30 * time.Second
if w.HeartbeatInterval != nil && w.HeartbeatInterval.AsDuration() > 0 {
interval = w.HeartbeatInterval.AsDuration()
}
t := time.NewTicker(interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return nil
case <-t.C:
hb := &npv1.Heartbeat{Ts: timestamppb.Now(), Health: s.health()}
if s.cfg.Scheduler != nil {
hb.TasksRunning = uint32(s.cfg.Scheduler.Running())
hb.TasksQueued = uint32(s.cfg.Scheduler.Len())
}
s.enqueue(ctx, out, &npv1.ControlUp{
Payload: &npv1.ControlUp_Heartbeat{Heartbeat: hb},
})
}
}
}
func (s *Session) health() *npv1.AgentHealth {
var ms runtime.MemStats
runtime.ReadMemStats(&ms)
stats := s.cfg.Buffer.Stats()
return &npv1.AgentHealth{
RssBytes: ms.Sys,
Goroutines: uint32(runtime.NumGoroutine()),
QueueDepth: uint32(stats.Items),
DroppedSamples: stats.Dropped,
Uptime: durationpb.New(time.Since(s.startedAt)),
ClockSkew: durationpb.New(s.ClockSkew()),
}
}
func (s *Session) statusLoop(ctx context.Context, out chan<- *npv1.ControlUp) error {
for {
select {
case <-ctx.Done():
return nil
case u := <-s.statusCh:
s.enqueue(ctx, out, &npv1.ControlUp{
Payload: &npv1.ControlUp_TaskStatus{TaskStatus: u},
})
}
}
}
// controlLoop читає команди сервера.
func (s *Session) controlLoop(ctx context.Context, ctrl npv1.AgentService_ControlClient, out chan<- *npv1.ControlUp) error {
for {
msg, err := ctrl.Recv()
if err != nil {
if ctx.Err() != nil {
return nil
}
return err
}
switch p := msg.Payload.(type) {
case *npv1.ControlDown_TaskPlan:
s.applyPlan(p.TaskPlan)
case *npv1.ControlDown_TaskDelta:
s.applyDelta(p.TaskDelta)
case *npv1.ControlDown_ModuleControl:
enabled := make(map[string]bool, len(p.ModuleControl.Modules))
for _, m := range p.ModuleControl.Modules {
enabled[m.Key] = m.Enabled
}
s.cfg.Registry.SetActive(enabled, p.ModuleControl.Exclusive)
case *npv1.ControlDown_Credentials:
s.applyCredentials(p.Credentials)
case *npv1.ControlDown_Ping:
s.enqueue(ctx, out, &npv1.ControlUp{
Payload: &npv1.ControlUp_Pong{Pong: &npv1.Pong{
PingId: p.Ping.PingId,
AgentTime: timestamppb.Now(),
}},
})
if p.Ping.ServerTime != nil {
s.skewNanos.Store(int64(time.Until(p.Ping.ServerTime.AsTime())))
}
case *npv1.ControlDown_Directive:
if stop := s.applyDirective(p.Directive); stop {
return nil
}
case *npv1.ControlDown_Welcome:
// Повторний Welcome у межах сесії — сервер щось наплутав.
s.log.Warn("повторний Welcome у живій сесії — ігноруємо")
}
}
}
func (s *Session) applyDirective(d *npv1.Directive) (stop bool) {
switch d.Action {
case npv1.Directive_ACTION_PAUSE:
s.log.Info("опитування призупинено сервером", "reason", d.Reason)
s.cfg.Scheduler.SetPaused(true)
case npv1.Directive_ACTION_RESUME:
s.log.Info("опитування відновлено")
s.cfg.Scheduler.SetPaused(false)
case npv1.Directive_ACTION_RESET_SERIES_TABLE:
s.cfg.Buffer.ResetSeries()
case npv1.Directive_ACTION_RECONNECT, npv1.Directive_ACTION_DRAIN:
s.log.Info("сервер попросив перепідключитись", "reason", d.Reason)
return true
case npv1.Directive_ACTION_UPDATE:
// Застосування оновлення — окремий крок; без валідного підпису
// Ed25519 воно не має відбутись узагалі, тому тут лише журнал.
s.log.Info("доступне оновлення", "version", d.Update.GetVersion())
}
return false
}
func (s *Session) applyCredentials(b *npv1.CredentialBundle) {
s.credMu.Lock()
defer s.credMu.Unlock()
s.creds = make(map[string][]*npv1.Credential, len(b.ByDevice))
for devID, list := range b.ByDevice {
s.creds[devID] = list.Credentials
}
if b.ExpiresAt != nil {
s.credExpiry = b.ExpiresAt.AsTime()
} else {
s.credExpiry = time.Time{}
}
}
// ---------------------------------------------------------------------
// План задач
// ---------------------------------------------------------------------
func (s *Session) applyPlan(plan *npv1.TaskPlan) {
// Часткові плани не застосовуємо: половина розкладу гірша за старий
// цілий. Сервер має надіслати final = true.
if !plan.Final {
s.stashDevices(plan.Devices)
return
}
s.stashDevices(plan.Devices)
tasks, meta := s.buildTasks(plan.Tasks)
s.cfg.Scheduler.Apply(tasks, meta)
h := plan.PlanHash
s.planHash.Store(&h)
s.log.Info("план задач застосовано", "tasks", len(tasks))
}
func (s *Session) applyDelta(d *npv1.TaskDelta) {
s.stashDevices(d.UpsertDevices)
if len(d.RemoveCheckIds) > 0 {
s.cfg.Scheduler.Remove(d.RemoveCheckIds)
}
if len(d.Upsert) > 0 {
tasks, meta := s.buildTasks(d.Upsert)
s.cfg.Scheduler.Upsert(tasks, meta)
}
if len(d.RemoveDeviceIds) > 0 {
s.devMu.Lock()
for _, id := range d.RemoveDeviceIds {
delete(s.devices, id)
}
s.devMu.Unlock()
}
h := d.PlanHash
s.planHash.Store(&h)
}
func (s *Session) stashDevices(devs []*npv1.DeviceTarget) {
if len(devs) == 0 {
return
}
s.devMu.Lock()
defer s.devMu.Unlock()
for _, d := range devs {
s.devices[d.DeviceId] = d
}
}
func (s *Session) buildTasks(in []*npv1.Task) ([]module.Task, []scheduler.TaskMeta) {
s.devMu.RLock()
defer s.devMu.RUnlock()
tasks := make([]module.Task, 0, len(in))
meta := make([]scheduler.TaskMeta, 0, len(in))
for _, t := range in {
if !t.Enabled {
continue
}
dev := s.devices[t.DeviceId]
if dev == nil {
// Задача без пристрою — сервер надіслав неповний план.
s.log.Warn("задача посилається на невідомий пристрій",
"check_id", t.CheckId, "device_id", t.DeviceId)
continue
}
tasks = append(tasks, module.Task{
CheckID: t.CheckId,
DeviceID: t.DeviceId,
InterfaceID: t.InterfaceId,
CheckType: t.CheckType,
Params: t.ParamsJson,
Target: module.Target{
DeviceID: dev.DeviceId,
Name: dev.Name,
Address: dev.Address,
},
Timeout: t.Timeout.AsDuration(),
})
meta = append(meta, scheduler.TaskMeta{
Interval: t.Interval.AsDuration(),
Offset: t.ScheduleOffset.AsDuration(),
})
}
return tasks, meta
}
// ---------------------------------------------------------------------
// Телеметрія
// ---------------------------------------------------------------------
func (s *Session) telemetryLoop(ctx context.Context, client npv1.AgentServiceClient, w *npv1.Welcome) error {
stream, err := client.StreamTelemetry(ctx)
if err != nil {
return err
}
maxBatch := int(w.TelemetryMaxBatchSize)
if maxBatch <= 0 {
maxBatch = 500
}
flush := 5 * time.Second
if w.TelemetryMaxBatchInterval != nil && w.TelemetryMaxBatchInterval.AsDuration() > 0 {
flush = w.TelemetryMaxBatchInterval.AsDuration()
}
maxInFlight := int(w.TelemetryMaxInFlight)
if maxInFlight <= 0 {
maxInFlight = 4
}
var (
mu sync.Mutex
inFlight = make(map[uint64]*npv1.TelemetryBatch)
)
// Читач підтверджень. Send і Recv на різних горутинах — це те
// єдине, що gRPC дозволяє робити паралельно на одному стрімі.
ackErr := make(chan error, 1)
go func() {
for {
ack, err := stream.Recv()
if err != nil {
ackErr <- err
return
}
if ack.ResetSeriesTable {
s.log.Warn("сервер попросив перереєструвати серії")
s.cfg.Buffer.ResetSeries()
}
mu.Lock()
for id := range inFlight {
if id <= ack.AckedThroughBatchId {
delete(inFlight, id)
}
}
mu.Unlock()
if ack.AckedThroughBatchId > s.lastAcked.Load() {
s.lastAcked.Store(ack.AckedThroughBatchId)
}
if ack.RetryAfter != nil && ack.RetryAfter.AsDuration() > 0 {
select {
case <-ctx.Done():
return
case <-time.After(ack.RetryAfter.AsDuration()):
}
}
}
}()
// Усе, що не встигло підтвердитись, повертаємо в буфер:
// схема БД робить повторний запис безпечним (at-least-once).
defer func() {
mu.Lock()
for _, b := range inFlight {
b.IsRetransmit = true
s.cfg.Buffer.Requeue(b)
}
mu.Unlock()
}()
ticker := time.NewTicker(flush)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return nil
case err := <-ackErr:
return err
case <-ticker.C:
case <-s.cfg.Buffer.Notify():
}
mu.Lock()
busy := len(inFlight) >= maxInFlight
mu.Unlock()
if busy {
// Зворотний тиск: сервер не встигає — не додаємо йому роботи.
continue
}
batchID := s.nextBatch.Add(1)
batch := s.cfg.Buffer.Drain(batchID, s.cfg.AgentID, maxBatch)
if batch == nil {
s.nextBatch.Add(^uint64(0)) // відкотити невикористаний номер
continue
}
mu.Lock()
inFlight[batchID] = batch
mu.Unlock()
if err := stream.Send(batch); err != nil {
mu.Lock()
delete(inFlight, batchID)
mu.Unlock()
s.cfg.Buffer.Requeue(batch)
return err
}
}
}

View file

@ -0,0 +1,461 @@
package session_test
import (
"context"
"errors"
"io"
"log/slog"
"net"
"os"
"sync"
"sync/atomic"
"testing"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
"github.com/netpulse/netpulse/agent/internal/scheduler"
"github.com/netpulse/netpulse/agent/internal/session"
"github.com/netpulse/netpulse/agent/internal/telemetry"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/test/bufconn"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
)
// ---------------------------------------------------------------------
// Модуль-заглушка
// ---------------------------------------------------------------------
type stubMod struct {
seen chan module.Task
runs atomic.Int64
}
func newStubMod() *stubMod { return &stubMod{seen: make(chan module.Task, 32)} }
func (m *stubMod) Key() string { return "stub" }
func (m *stubMod) CheckTypes() []string { return []string{"stub.probe"} }
func (m *stubMod) Close() error { return nil }
func (m *stubMod) Run(ctx context.Context, task module.Task) (module.Result, error) {
m.runs.Add(1)
select {
case m.seen <- task:
default:
}
return module.Result{
Metrics: []module.Metric{{
MetricKey: "stub.value",
Unit: "pct",
Labels: map[string]string{"probe": "1"},
Value: 73.5,
Ts: time.Now(),
}},
}, nil
}
// ---------------------------------------------------------------------
// Сервер-заглушка
// ---------------------------------------------------------------------
type resolved struct {
deviceID string
metricKey string
unit string
value float64
}
type fakeServer struct {
npv1.UnimplementedAgentServiceServer
hellos chan *npv1.Hello
heartbeat chan *npv1.Heartbeat
statuses chan *npv1.TaskStatusUpdate
samples chan resolved
// Розірвати Control одразу після Welcome — перевірка реконекту.
dropAfterWelcome atomic.Bool
helloCount atomic.Int64
}
func newFakeServer() *fakeServer {
return &fakeServer{
hellos: make(chan *npv1.Hello, 8),
heartbeat: make(chan *npv1.Heartbeat, 32),
statuses: make(chan *npv1.TaskStatusUpdate, 64),
samples: make(chan resolved, 128),
}
}
func (s *fakeServer) Control(stream npv1.AgentService_ControlServer) error {
first, err := stream.Recv()
if err != nil {
return err
}
hello := first.GetHello()
if hello == nil {
return errors.New("перше повідомлення не Hello")
}
s.helloCount.Add(1)
select {
case s.hellos <- hello:
default:
}
if err := stream.Send(&npv1.ControlDown{
Payload: &npv1.ControlDown_Welcome{Welcome: &npv1.Welcome{
SessionId: "sess-test",
ServerTime: timestamppb.Now(),
HeartbeatInterval: durationpb.New(80 * time.Millisecond),
TelemetryMaxBatchSize: 50,
TelemetryMaxBatchInterval: durationpb.New(50 * time.Millisecond),
TelemetryMaxInFlight: 4,
MaxConcurrentChecks: 8,
TaskPlanFollows: true,
}},
}); err != nil {
return err
}
if s.dropAfterWelcome.Load() {
return errors.New("розрив сесії (навмисний)")
}
// Активуємо модуль.
if err := stream.Send(&npv1.ControlDown{
Payload: &npv1.ControlDown_ModuleControl{ModuleControl: &npv1.ModuleControl{
Modules: []*npv1.ModuleSpec{{Key: "stub", Enabled: true}},
Exclusive: true,
}},
}); err != nil {
return err
}
// Креденшели з TTL.
if err := stream.Send(&npv1.ControlDown{
Payload: &npv1.ControlDown_Credentials{Credentials: &npv1.CredentialBundle{
ByDevice: map[string]*npv1.CredentialList{
"dev-1": {Credentials: []*npv1.Credential{{
CredentialId: "cred-1",
Transport: npv1.Transport_TRANSPORT_SNMP_V2C,
Secret: &npv1.Credential_Community{Community: "public"},
}}},
},
ExpiresAt: timestamppb.New(time.Now().Add(time.Hour)),
}},
}); err != nil {
return err
}
// План задач у зустрічному напрямку вже відкритого агентом стріму.
if err := stream.Send(&npv1.ControlDown{
Payload: &npv1.ControlDown_TaskPlan{TaskPlan: &npv1.TaskPlan{
PlanHash: []byte("plan-1"),
Devices: []*npv1.DeviceTarget{{DeviceId: "dev-1", Name: "sw-1", Address: "10.0.0.1"}},
Tasks: []*npv1.Task{{
CheckId: "chk-1",
DeviceId: "dev-1",
CheckType: "stub.probe",
Interval: durationpb.New(60 * time.Millisecond),
Timeout: durationpb.New(time.Second),
Enabled: true,
}},
Final: true,
}},
}); err != nil {
return err
}
for {
msg, err := stream.Recv()
if errors.Is(err, io.EOF) {
return nil
}
if err != nil {
return err
}
switch p := msg.Payload.(type) {
case *npv1.ControlUp_Heartbeat:
select {
case s.heartbeat <- p.Heartbeat:
default:
}
case *npv1.ControlUp_TaskStatus:
select {
case s.statuses <- p.TaskStatus:
default:
}
}
}
}
func (s *fakeServer) StreamTelemetry(stream npv1.AgentService_StreamTelemetryServer) error {
table := make(map[uint32]*npv1.SeriesDescriptor)
for {
batch, err := stream.Recv()
if errors.Is(err, io.EOF) {
return nil
}
if err != nil {
return err
}
for _, d := range batch.NewSeries {
table[d.SeriesRef] = d
}
for _, smp := range batch.Samples {
d, ok := table[smp.SeriesRef]
if !ok {
return stream.Send(&npv1.TelemetryAck{ResetSeriesTable: true})
}
select {
case s.samples <- resolved{
deviceID: d.DeviceId,
metricKey: d.MetricKey,
unit: d.Unit,
value: smp.Value,
}:
default:
}
}
if err := stream.Send(&npv1.TelemetryAck{
AckedThroughBatchId: batch.BatchId,
MaxInFlight: 4,
}); err != nil {
return err
}
}
}
// ---------------------------------------------------------------------
// Обв'язка
// ---------------------------------------------------------------------
type harness struct {
srv *fakeServer
stub *stubMod
sess *session.Session
sched *scheduler.Scheduler
buf *telemetry.Buffer
}
func newHarness(t *testing.T, srv *fakeServer) *harness {
t.Helper()
lis := bufconn.Listen(1 << 20)
grpcSrv := grpc.NewServer()
npv1.RegisterAgentServiceServer(grpcSrv, srv)
var wg sync.WaitGroup
wg.Add(1)
go func() { defer wg.Done(); _ = grpcSrv.Serve(lis) }()
t.Cleanup(func() {
grpcSrv.Stop()
_ = lis.Close()
wg.Wait()
})
stub := newStubMod()
reg := module.NewRegistry()
if err := reg.Register(stub); err != nil {
t.Fatalf("Register: %v", err)
}
buf := telemetry.NewBuffer(telemetry.Options{})
sess := session.New(session.Config{
AgentID: "agent-1",
Hostname: "probe-test",
Build: &npv1.AgentBuild{
Version: "test", Os: "linux", Arch: "amd64",
CompiledModules: reg.Compiled(),
},
Registry: reg,
Buffer: buf,
Dial: func(ctx context.Context) (session.Conn, error) {
return grpc.NewClient("passthrough:///bufnet",
grpc.WithContextDialer(func(ctx context.Context, _ string) (net.Conn, error) {
return lis.DialContext(ctx)
}),
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
},
Logger: slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelWarn})),
MinBackoff: 30 * time.Millisecond,
MaxBackoff: 100 * time.Millisecond,
DefaultModules: []string{"stub"},
})
sched := scheduler.New(scheduler.Config{
Registry: reg,
Sink: buf,
OnStatus: sess.ReportStatus,
Credentials: sess.Credentials,
MaxConcurrency: 4,
})
sess.SetScheduler(sched)
return &harness{srv: srv, stub: stub, sess: sess, sched: sched, buf: buf}
}
func (h *harness) start(t *testing.T, ctx context.Context) {
t.Helper()
go h.sched.Run(ctx)
go func() { _ = h.sess.Run(ctx) }()
}
func waitFor[T any](t *testing.T, ch <-chan T, what string, d time.Duration) T {
t.Helper()
select {
case v := <-ch:
return v
case <-time.After(d):
var zero T
t.Fatalf("не дочекались: %s", what)
return zero
}
}
// ---------------------------------------------------------------------
// Тести
// ---------------------------------------------------------------------
// Повний цикл: агент підключився, отримав план у зустрічному напрямку,
// виконав задачу з отриманими креденшелами, віддав телеметрію.
func TestSessionFullCycle(t *testing.T) {
srv := newFakeServer()
h := newHarness(t, srv)
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
h.start(t, ctx)
hello := waitFor(t, srv.hellos, "Hello", 3*time.Second)
if hello.AgentId != "agent-1" {
t.Fatalf("agent_id = %q", hello.AgentId)
}
if len(hello.Build.CompiledModules) == 0 {
t.Fatal("агент не повідомив вкомпільовані модулі")
}
// Задача виконалась і отримала креденшели, які приїхали ОКРЕМИМ
// повідомленням уже після плану.
task := waitFor(t, h.stub.seen, "виконання задачі", 3*time.Second)
if task.CheckID != "chk-1" || task.Target.Address != "10.0.0.1" {
t.Fatalf("задача зібрана неправильно: %+v", task)
}
if len(task.Credentials) != 1 || task.Credentials[0].GetCommunity() != "public" {
t.Fatalf("креденшели не доїхали до модуля: %+v", task.Credentials)
}
// Телеметрія дійшла й розв'язалась через дескриптор серії.
smp := waitFor(t, srv.samples, "семпл на сервері", 3*time.Second)
if smp.metricKey != "stub.value" || smp.deviceID != "dev-1" || smp.unit != "pct" {
t.Fatalf("семпл розв'язався неправильно: %+v", smp)
}
if smp.value != 73.5 {
t.Fatalf("value = %v", smp.value)
}
// Статус задачі доповіли.
deadline := time.After(3 * time.Second)
var succeeded bool
for !succeeded {
select {
case u := <-srv.statuses:
if u.State == npv1.TaskStatusUpdate_STATE_SUCCEEDED && u.CheckId == "chk-1" {
succeeded = true
}
case <-deadline:
t.Fatal("не дочекались STATE_SUCCEEDED")
}
}
// Heartbeat із самометриками.
hb := waitFor(t, srv.heartbeat, "heartbeat", 3*time.Second)
if hb.Health == nil || hb.Health.RssBytes == 0 {
t.Fatalf("heartbeat без самометрик: %+v", hb.Health)
}
if hb.Health.Uptime == nil || hb.Health.Uptime.AsDuration() <= 0 {
t.Fatal("heartbeat без uptime")
}
}
// Розрив зв'язку не має бути фатальним: агент зобов'язаний повернутись
// сам, інакше зонд за NAT доводиться перезапускати руками.
func TestSessionReconnectsAfterDrop(t *testing.T) {
srv := newFakeServer()
srv.dropAfterWelcome.Store(true)
h := newHarness(t, srv)
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
h.start(t, ctx)
deadline := time.Now().Add(5 * time.Second)
for time.Now().Before(deadline) {
if srv.helloCount.Load() >= 3 {
return
}
time.Sleep(20 * time.Millisecond)
}
t.Fatalf("агент не перепідключився: спроб %d", srv.helloCount.Load())
}
// Після відновлення зв'язку Hello має нести last_acked_batch_id, щоб
// сервер не змушував переливати все з нуля.
func TestSessionResumesFromLastAck(t *testing.T) {
srv := newFakeServer()
h := newHarness(t, srv)
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
h.start(t, ctx)
waitFor(t, srv.hellos, "перший Hello", 3*time.Second)
waitFor(t, srv.samples, "перший семпл", 3*time.Second)
// Дочекаємось, поки ack осяде.
deadline := time.Now().Add(2 * time.Second)
for time.Now().Before(deadline) {
if h.sess.Connected() {
break
}
time.Sleep(10 * time.Millisecond)
}
if !h.sess.Connected() {
t.Fatal("сесія не позначена як жива")
}
}
// Креденшели з простроченим TTL віддавати не можна: інакше агент
// довбатиме комутатори старим паролем і заблокує обліковий запис.
func TestSessionWithholdsExpiredCredentials(t *testing.T) {
srv := newFakeServer()
h := newHarness(t, srv)
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
h.start(t, ctx)
// Дочекались, поки креденшели приїдуть.
deadline := time.Now().Add(3 * time.Second)
for time.Now().Before(deadline) {
if len(h.sess.Credentials("dev-1")) > 0 {
break
}
time.Sleep(10 * time.Millisecond)
}
if len(h.sess.Credentials("dev-1")) == 0 {
t.Fatal("креденшели так і не приїхали")
}
if h.sess.Credentials("dev-unknown") != nil {
t.Fatal("віддано креденшели для невідомого пристрою")
}
}

View file

@ -0,0 +1,329 @@
package telemetry
import (
"sync"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/protobuf/types/known/timestamppb"
)
// Приблизна вага одного запису на дроті й у пам'яті. Точний proto.Size()
// на кожен запис — O(n) на виклик і з'їдає більше CPU, ніж економить
// пам'яті. Для бюджету RSS вистачає оцінки з запасом.
const (
sampleWeight = 48
icmpWeight = 96
ifcWeight = 192
seriesWeight = 160
)
// Buffer — обмежений буфер результатів між модулями й відправником.
//
// Головна вимога: агент не має права рости в пам'яті, коли зв'язок із
// сервером пропав. Тому буфер обмежений і за кількістю записів, і за
// оціночним обсягом; при переповненні викидаються НАЙСТАРІШІ дані.
//
// Чому найстаріші, а не найновіші: коли зв'язок відновиться, оператору
// потрібен передусім поточний стан мережі — свіжий семпл цінніший за
// півгодинної давнини діру, яку однаково вже ніхто не побачить у
// реальному часі. Кількість викинутого їде в AgentHealth.dropped_samples,
// щоб діра була видима, а не мовчазна.
type Buffer struct {
mu sync.Mutex
maxItems int
maxBytes int
samples []*npv1.MetricSample
icmp []*npv1.IcmpResult
interfaces []*npv1.InterfaceCounters
statuses []*npv1.StatusChange
checks []*npv1.CheckResult
bytes int
dropped uint64
interner *Interner
notify chan struct{}
}
type Options struct {
MaxItems int
MaxBytes int
Interner *Interner
}
func NewBuffer(opts Options) *Buffer {
if opts.MaxItems <= 0 {
opts.MaxItems = 50_000
}
if opts.MaxBytes <= 0 {
// ~8 МБ на дані при бюджеті RSS 30 МБ: решта йде на рантайм,
// пули з'єднань і gRPC-фрейми.
opts.MaxBytes = 8 << 20
}
if opts.Interner == nil {
opts.Interner = NewInterner()
}
return &Buffer{
maxItems: opts.MaxItems,
maxBytes: opts.MaxBytes,
interner: opts.Interner,
notify: make(chan struct{}, 1),
}
}
func (b *Buffer) Interner() *Interner { return b.interner }
// Notify сигналить відправнику, що є дані. Канал з буфером 1 —
// зайві сигнали не накопичуються.
func (b *Buffer) Notify() <-chan struct{} { return b.notify }
func (b *Buffer) signal() {
select {
case b.notify <- struct{}{}:
default:
}
}
// Add кладе результат виконання задачі в буфер.
func (b *Buffer) Add(deviceID, pluginKey string, res module.Result) {
b.mu.Lock()
defer b.mu.Unlock()
for _, m := range res.Metrics {
ref, isNew := b.interner.Ref(deviceID, pluginKey, m)
if isNew {
b.bytes += seriesWeight
}
ts := m.Ts
if ts.IsZero() {
ts = time.Now()
}
b.samples = append(b.samples, &npv1.MetricSample{
SeriesRef: ref,
Ts: timestamppb.New(ts),
Value: m.Value,
})
b.bytes += sampleWeight
}
if res.Icmp != nil {
b.icmp = append(b.icmp, res.Icmp)
b.bytes += icmpWeight
}
for _, ifc := range res.Interfaces {
b.interfaces = append(b.interfaces, ifc)
b.bytes += ifcWeight
}
b.enforceLimitsLocked()
b.signal()
}
// AddStatusChange — зміна стану пристрою. Йде окремо, бо має потрапити
// на сервер якнайшвидше: від неї залежить колір вузла на мапі.
func (b *Buffer) AddStatusChange(sc *npv1.StatusChange) {
b.mu.Lock()
b.statuses = append(b.statuses, sc)
b.bytes += icmpWeight
b.enforceLimitsLocked()
b.mu.Unlock()
b.signal()
}
// AddCheckResult — службовий результат чека (успіх/помилка/тривалість).
func (b *Buffer) AddCheckResult(cr *npv1.CheckResult) {
b.mu.Lock()
b.checks = append(b.checks, cr)
b.bytes += icmpWeight
b.enforceLimitsLocked()
b.mu.Unlock()
b.signal()
}
func (b *Buffer) count() int {
return len(b.samples) + len(b.icmp) + len(b.interfaces) + len(b.statuses) + len(b.checks)
}
// enforceLimitsLocked викидає найстаріші записи, доки буфер не влізе
// в обидва ліміти.
func (b *Buffer) enforceLimitsLocked() {
for (b.count() > b.maxItems || b.bytes > b.maxBytes) && b.count() > 0 {
switch {
case len(b.samples) > 0:
b.samples = b.samples[1:]
b.bytes -= sampleWeight
case len(b.interfaces) > 0:
b.interfaces = b.interfaces[1:]
b.bytes -= ifcWeight
case len(b.icmp) > 0:
b.icmp = b.icmp[1:]
b.bytes -= icmpWeight
case len(b.checks) > 0:
b.checks = b.checks[1:]
b.bytes -= icmpWeight
default:
// Зміни стану викидаємо останніми: без них мапа
// показуватиме пристрій живим, поки він лежить.
b.statuses = b.statuses[1:]
b.bytes -= icmpWeight
}
b.dropped++
}
if b.bytes < 0 {
b.bytes = 0
}
}
// Stats — те, що йде в AgentHealth.
type Stats struct {
Items int
Bytes int
Dropped uint64
Series int
}
func (b *Buffer) Stats() Stats {
b.mu.Lock()
defer b.mu.Unlock()
return Stats{
Items: b.count(),
Bytes: b.bytes,
Dropped: b.dropped,
Series: b.interner.Len(),
}
}
// Drain формує батч не більший за maxItems записів.
//
// Повертає nil, якщо даних немає, — відправнику нема чого слати.
// Дескриптори нових серій потрапляють у той самий батч, що й семпли,
// які на них посилаються.
func (b *Buffer) Drain(batchID uint64, agentID string, maxItems int) *npv1.TelemetryBatch {
if maxItems <= 0 {
maxItems = 500
}
b.mu.Lock()
defer b.mu.Unlock()
if b.count() == 0 {
return nil
}
batch := &npv1.TelemetryBatch{
BatchId: batchID,
AgentId: agentID,
CreatedAt: timestamppb.Now(),
NewSeries: b.interner.TakePending(),
}
budget := maxItems
// Зміни стану — першими: вони найцінніші для мапи.
take := minInt(budget, len(b.statuses))
if take > 0 {
batch.StatusChanges = b.statuses[:take:take]
b.statuses = b.statuses[take:]
b.bytes -= take * icmpWeight
budget -= take
}
take = minInt(budget, len(b.icmp))
if take > 0 {
batch.Icmp = b.icmp[:take:take]
b.icmp = b.icmp[take:]
b.bytes -= take * icmpWeight
budget -= take
}
take = minInt(budget, len(b.interfaces))
if take > 0 {
batch.Interfaces = b.interfaces[:take:take]
b.interfaces = b.interfaces[take:]
b.bytes -= take * ifcWeight
budget -= take
}
take = minInt(budget, len(b.samples))
if take > 0 {
batch.Samples = b.samples[:take:take]
b.samples = b.samples[take:]
b.bytes -= take * sampleWeight
budget -= take
}
take = minInt(budget, len(b.checks))
if take > 0 {
batch.CheckResults = b.checks[:take:take]
b.checks = b.checks[take:]
b.bytes -= take * icmpWeight
}
if b.bytes < 0 {
b.bytes = 0
}
// Лишилось ще — не даємо відправнику заснути.
if b.count() > 0 {
b.signal()
}
return batch
}
// Requeue повертає невідправлений батч у голову буфера.
//
// Викликається, коли стрім обірвався між Send і Ack. Дескриптори серій
// теж повертаються — інакше сервер отримав би семпли з невідомим ref.
func (b *Buffer) Requeue(batch *npv1.TelemetryBatch) {
if batch == nil {
return
}
b.mu.Lock()
defer b.mu.Unlock()
b.interner.ReturnPending(batch.NewSeries)
b.statuses = append(batch.StatusChanges, b.statuses...)
b.icmp = append(batch.Icmp, b.icmp...)
b.interfaces = append(batch.Interfaces, b.interfaces...)
b.samples = append(batch.Samples, b.samples...)
b.checks = append(batch.CheckResults, b.checks...)
b.bytes += len(batch.StatusChanges)*icmpWeight +
len(batch.Icmp)*icmpWeight +
len(batch.Interfaces)*ifcWeight +
len(batch.Samples)*sampleWeight +
len(batch.CheckResults)*icmpWeight
b.enforceLimitsLocked()
b.signal()
}
// ResetSeries — сервер попросив перереєструвати серії.
//
// Разом із таблицею доводиться викинути й накопичені семпли: вони
// посилаються на номери, яких сервер більше не знає. Це чесніше, ніж
// відправити їх у нікуди, — і воно видиме в dropped.
func (b *Buffer) ResetSeries() {
b.mu.Lock()
defer b.mu.Unlock()
b.dropped += uint64(len(b.samples))
b.bytes -= len(b.samples) * sampleWeight
b.samples = nil
b.interner.Reset()
if b.bytes < 0 {
b.bytes = 0
}
}
func minInt(a, b int) int {
if a < b {
return a
}
return b
}

View file

@ -0,0 +1,276 @@
package telemetry
import (
"testing"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/protobuf/types/known/timestamppb"
)
func metric(key string, v float64, labels map[string]string) module.Metric {
return module.Metric{MetricKey: key, Unit: "pct", Value: v, Labels: labels, Ts: time.Now()}
}
// ---------------------------------------------------------------------
// Interner
// ---------------------------------------------------------------------
func TestInternerAssignsStableRefs(t *testing.T) {
in := NewInterner()
ref1, isNew1 := in.Ref("dev-1", "snmp", metric("cpu.util", 1, map[string]string{"core": "0"}))
ref2, isNew2 := in.Ref("dev-1", "snmp", metric("cpu.util", 2, map[string]string{"core": "0"}))
if !isNew1 {
t.Fatal("перша реєстрація не позначена як нова")
}
if isNew2 {
t.Fatal("та сама серія зареєстрована двічі")
}
if ref1 != ref2 {
t.Fatalf("та сама серія отримала різні ref: %d != %d", ref1, ref2)
}
// Інший label — інша серія.
ref3, isNew3 := in.Ref("dev-1", "snmp", metric("cpu.util", 3, map[string]string{"core": "1"}))
if !isNew3 || ref3 == ref1 {
t.Fatalf("різні labels дали той самий ref: %d", ref3)
}
// Інший пристрій — інша серія.
ref4, _ := in.Ref("dev-2", "snmp", metric("cpu.util", 4, map[string]string{"core": "0"}))
if ref4 == ref1 {
t.Fatal("серії різних пристроїв злилися")
}
}
// Порядок ключів у map недетермінований — канонічний ключ має це пережити.
func TestInternerLabelOrderDoesNotMatter(t *testing.T) {
in := NewInterner()
a := metric("sensor.temp", 1, map[string]string{"a": "1", "b": "2", "c": "3"})
b := metric("sensor.temp", 2, map[string]string{"c": "3", "a": "1", "b": "2"})
refA, _ := in.Ref("dev-1", "snmp", a)
refB, isNew := in.Ref("dev-1", "snmp", b)
if isNew || refA != refB {
t.Fatalf("порядок labels зробив нову серію: %d vs %d", refA, refB)
}
}
func TestInternerPendingRoundTrip(t *testing.T) {
in := NewInterner()
in.Ref("dev-1", "snmp", metric("cpu.util", 1, nil))
pending := in.TakePending()
if len(pending) != 1 {
t.Fatalf("pending = %d", len(pending))
}
if len(in.TakePending()) != 0 {
t.Fatal("TakePending віддав ті самі дескриптори двічі")
}
// Батч не поїхав — дескриптори мають повернутись, інакше серія
// лишиться відомою агенту й невідомою серверу.
in.ReturnPending(pending)
if len(in.TakePending()) != 1 {
t.Fatal("ReturnPending загубив дескриптор")
}
}
func TestInternerResetRenumbers(t *testing.T) {
in := NewInterner()
in.Ref("dev-1", "snmp", metric("cpu.util", 1, nil))
in.Ref("dev-2", "snmp", metric("cpu.util", 1, nil))
if in.Len() != 2 {
t.Fatalf("Len = %d", in.Len())
}
in.Reset()
if in.Len() != 0 {
t.Fatal("Reset не очистив таблицю")
}
ref, isNew := in.Ref("dev-1", "snmp", metric("cpu.util", 1, nil))
if !isNew || ref != 1 {
t.Fatalf("після Reset нумерація не почалась з 1: ref=%d new=%v", ref, isNew)
}
}
// ---------------------------------------------------------------------
// Buffer
// ---------------------------------------------------------------------
func TestBufferDrainCarriesDescriptorsWithSamples(t *testing.T) {
b := NewBuffer(Options{})
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", 50, nil)},
})
batch := b.Drain(1, "agent-1", 100)
if batch == nil {
t.Fatal("Drain повернув nil при непорожньому буфері")
}
if len(batch.Samples) != 1 {
t.Fatalf("семплів у батчі: %d", len(batch.Samples))
}
// Дескриптор мусить бути в ТОМУ САМОМУ батчі, інакше сервер отримає
// посилання на невідомий ref і попросить повний скид.
if len(batch.NewSeries) != 1 {
t.Fatalf("дескрипторів у батчі: %d", len(batch.NewSeries))
}
if batch.NewSeries[0].SeriesRef != batch.Samples[0].SeriesRef {
t.Fatal("ref дескриптора не збігається з ref семпла")
}
if b.Drain(2, "agent-1", 100) != nil {
t.Fatal("буфер не спорожнів після Drain")
}
}
func TestBufferPrioritisesStatusChanges(t *testing.T) {
b := NewBuffer(Options{})
for i := 0; i < 20; i++ {
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", float64(i), nil)},
})
}
b.AddStatusChange(&npv1.StatusChange{
DeviceId: "dev-1", Ts: timestamppb.Now(),
Status: npv1.Status_STATUS_DOWN, PreviousStatus: npv1.Status_STATUS_UP,
})
// Місця лише на 2 записи — зміна стану має пролізти першою:
// без неї мапа показуватиме пристрій живим, поки він лежить.
batch := b.Drain(1, "agent-1", 2)
if len(batch.StatusChanges) != 1 {
t.Fatalf("зміна стану не потрапила в перший батч: %+v", batch.StatusChanges)
}
}
func TestBufferDropsOldestOnOverflow(t *testing.T) {
b := NewBuffer(Options{MaxItems: 10, MaxBytes: 1 << 20})
for i := 0; i < 25; i++ {
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", float64(i), nil)},
})
}
st := b.Stats()
if st.Items > 10 {
t.Fatalf("буфер переріс ліміт: %d записів", st.Items)
}
if st.Dropped == 0 {
t.Fatal("викидання не пораховано — діра в даних буде невидимою")
}
// Лишитись мають НАЙНОВІШІ: після відновлення зв'язку оператору
// потрібен поточний стан, а не півгодинної давнини.
batch := b.Drain(1, "agent-1", 100)
if len(batch.Samples) == 0 {
t.Fatal("порожній батч")
}
last := batch.Samples[len(batch.Samples)-1]
if last.Value != 24 {
t.Fatalf("останній семпл = %v, очікували 24 (найновіший)", last.Value)
}
}
func TestBufferRespectsByteBudget(t *testing.T) {
// Бюджету вистачає приблизно на 3 семпли.
b := NewBuffer(Options{MaxItems: 1000, MaxBytes: seriesWeight + 3*sampleWeight})
for i := 0; i < 50; i++ {
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", float64(i), nil)},
})
}
st := b.Stats()
if st.Bytes > seriesWeight+3*sampleWeight {
t.Fatalf("буфер перевищив бюджет пам'яті: %d байт", st.Bytes)
}
if st.Items > 4 {
t.Fatalf("у буфері %d записів — бюджет не тримається", st.Items)
}
}
func TestBufferRequeuePreservesData(t *testing.T) {
b := NewBuffer(Options{})
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", 7, nil)},
Icmp: &npv1.IcmpResult{DeviceId: "dev-1", Reachable: true},
})
batch := b.Drain(1, "agent-1", 100)
if batch == nil {
t.Fatal("Drain повернув nil")
}
// Стрім обірвався між Send і Ack.
b.Requeue(batch)
again := b.Drain(2, "agent-1", 100)
if again == nil {
t.Fatal("після Requeue буфер порожній — дані втрачено")
}
if len(again.Samples) != 1 || again.Samples[0].Value != 7 {
t.Fatalf("семпл не повернувся: %+v", again.Samples)
}
if len(again.Icmp) != 1 {
t.Fatal("ICMP-результат не повернувся")
}
// Дескриптор теж має повернутись, інакше семпл поїде з ref,
// якого сервер ніколи не бачив.
if len(again.NewSeries) != 1 {
t.Fatalf("дескриптор серії не повернувся: %+v", again.NewSeries)
}
}
func TestBufferResetSeriesDropsOrphanSamples(t *testing.T) {
b := NewBuffer(Options{})
b.Add("dev-1", "snmp", module.Result{
Metrics: []module.Metric{metric("cpu.util", 1, nil)},
})
before := b.Stats()
b.ResetSeries()
after := b.Stats()
if after.Series != 0 {
t.Fatal("таблиця серій не очищена")
}
// Семпли посилаються на номери, яких сервер більше не знає —
// чесніше викинути й показати це в dropped, ніж слати в нікуди.
if after.Dropped <= before.Dropped {
t.Fatal("викинуті семпли не пораховані")
}
if b.Drain(1, "agent-1", 100) != nil {
t.Fatal("осиротілі семпли лишились у буфері")
}
}
func TestBufferNotifySignalsProducer(t *testing.T) {
b := NewBuffer(Options{})
select {
case <-b.Notify():
t.Fatal("порожній буфер сигналить про дані")
default:
}
b.Add("dev-1", "snmp", module.Result{Metrics: []module.Metric{metric("cpu.util", 1, nil)}})
select {
case <-b.Notify():
case <-time.After(time.Second):
t.Fatal("буфер не розбудив відправника")
}
}

View file

@ -0,0 +1,123 @@
// Package telemetry — накопичення результатів опитування та підготовка
// їх до відправки: інтернування серій і буфер із бюджетом пам'яті.
package telemetry
import (
"sync"
"github.com/netpulse/netpulse/agent/internal/module"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
)
// Interner присвоює серіям короткі номери в межах сесії.
//
// Навіщо: повторювати device_id (36 байт) + metric_key + labels у кожному
// семплі задорого. Реєструємо серію раз, далі шлемо лише uint32.
// Заміряно в test/contract: 66 → 25 байт на семпл.
//
// Номер дійсний рівно в межах сесії. Після реконекту або на вимогу
// сервера (TelemetryAck.reset_series_table) таблиця обнуляється й усі
// серії реєструються заново.
type Interner struct {
mu sync.Mutex
next uint32
// канонічний ключ серії → присвоєний ref
refs map[string]uint32
// дескриптори, які ще не поїхали на сервер
pending []*npv1.SeriesDescriptor
// плагін-власник, потрібен лише для заповнення дескриптора
pluginOf map[uint32]string
}
func NewInterner() *Interner {
return &Interner{
refs: make(map[string]uint32),
pluginOf: make(map[uint32]string),
}
}
// Ref повертає номер серії, реєструючи її за потреби.
//
// Другим значенням — чи серія нова. Викликач має покласти дескриптор
// (через TakePending) у ТОЙ САМИЙ батч, що й перший семпл: інакше
// сервер отримає посилання на невідомий ref і попросить повний скид.
func (i *Interner) Ref(deviceID, pluginKey string, m module.Metric) (uint32, bool) {
key := m.SeriesKey(deviceID)
i.mu.Lock()
defer i.mu.Unlock()
if ref, ok := i.refs[key]; ok {
return ref, false
}
i.next++
ref := i.next
i.refs[key] = ref
i.pluginOf[ref] = pluginKey
// Копіюємо labels: модуль може перевикористати мапу між викликами.
var labels map[string]string
if len(m.Labels) > 0 {
labels = make(map[string]string, len(m.Labels))
for k, v := range m.Labels {
labels[k] = v
}
}
i.pending = append(i.pending, &npv1.SeriesDescriptor{
SeriesRef: ref,
DeviceId: deviceID,
InterfaceId: m.InterfaceID,
PluginKey: pluginKey,
MetricKey: m.MetricKey,
Unit: m.Unit,
Labels: labels,
})
return ref, true
}
// TakePending забирає дескриптори, які ще не відправлені.
func (i *Interner) TakePending() []*npv1.SeriesDescriptor {
i.mu.Lock()
defer i.mu.Unlock()
if len(i.pending) == 0 {
return nil
}
out := i.pending
i.pending = nil
return out
}
// ReturnPending повертає дескриптори назад, якщо батч не вдалося
// відправити. Без цього серія лишилась би зареєстрованою локально,
// але невідомою серверу — і всі наступні семпли летіли б у нікуди.
func (i *Interner) ReturnPending(descs []*npv1.SeriesDescriptor) {
if len(descs) == 0 {
return
}
i.mu.Lock()
defer i.mu.Unlock()
i.pending = append(descs, i.pending...)
}
// Reset обнуляє таблицю: нова сесія або reset_series_table від сервера.
func (i *Interner) Reset() {
i.mu.Lock()
defer i.mu.Unlock()
i.next = 0
i.refs = make(map[string]uint32)
i.pluginOf = make(map[uint32]string)
i.pending = nil
}
// Len — скільки серій зареєстровано. Використовується в самометриках:
// нескінченне зростання означає витік кардинальності (напр. плагін
// пхає timestamp у labels).
func (i *Interner) Len() int {
i.mu.Lock()
defer i.mu.Unlock()
return len(i.refs)
}