diff --git a/HISTORY.md b/HISTORY.md index a7d966a..a3ee325 100644 --- a/HISTORY.md +++ b/HISTORY.md @@ -189,10 +189,79 @@ Debian 13, Go 1.25.13. `go vet` чисто, `go test ./... -race` — усі п публічний репозиторій без автентифікації (`Credentials are incorrect`). Коміти лежать локально в `main` і чекають на токен. -### Далі (Етап 2, частина 3) +--- + +## 2026-08-14 — Етап 2 (частина 3): серверна сторона AgentService + +### Створено + +``` +netpulse/server/ +├── README.md рішення, параметри, стан перевірки +├── cmd/netpulse-server/ TLS, keepalive, m'яка зупинка, keyring із -dek +└── internal/ + ├── crypto/ AES-GCM-256, keyring із ротацією ключів + ├── store/ agents, plan, credentials, telemetry, discovery, ncm + └── grpcapi/ AgentService + перехоплювачі автентифікації +``` + +Плюс правки в агенті: `-token` і передача його в метаданих кожного виклику; +`Interner.MarkAllPending()` — перереєстрація серій на початку сесії. + +### Прийняті рішення + +1. **Ізоляція тенантів робиться двічі:** RLS (`SET LOCAL app.tenant_id`) плюс явний + предикат `tenant_id`. Не перестраховка: RLS не працює на гіпертаблях, а саме + туди йде вся телеметрія. +2. **`schedule_offset` — чиста функція від `check_id`**, тож будь-який вузол + сервера дає те саме значення. +3. **Хеш плану — лише з полів, що впливають на поведінку.** Зміна опису пристрою + не змушує переливати 50 000 задач. +4. **Токен каже, ЯКИЙ це зонд; сертифікат — що він має право говорити.** Чужий + `agent_id` при валідному токені → `PermissionDenied`. +5. **DEK не покидає сервер.** Дамп БД без ключів не дає жодного пароля. + Комплект креденшелів із TTL 1 год. +6. **Запис телеметрії — `ON CONFLICT DO NOTHING`** (at-least-once). +7. **Статус пристрою — один запит із умовним записом в історію**, інакше два + воркери наввипередки писали б неіснуючі переходи `up→up`. +8. **Впевненість зіставлення спадає за надійністю ознаки:** chassis-id 95 → + MAC 90 → IP 80 → sysName 60. Останнє низьке навмисно: sysName вводить людина. +9. **Перереєстрація серій замість обнулення нумерації.** Спершу агент мав би + скидати interner на реконекті, але це викидало б увесь накопичений за час + обриву буфер — саме ті дані, заради яких він накопичувався. + +### Перевірено на стенді + +11 інтеграційних тестів проти **живої БД зі схемою Етапу 1** і справжнього gRPC — +усі PASS з `-race`. Найцінніше: повтор батчу не дублює ані рядки, ані переходи в +історії; зустрічний звіт B→A не створює другий лінк; ручний (`is_pinned`) лінк не +затирається; тіло конфігу лежить зашифрованим. + +**Живий наскрізний прогін** (справжній агент + справжній сервер + БД, 40 с, +`icmp.ping` кожні 5 с проти 127.0.0.1): 5 ICMP-семплів, метрики `icmp.rtt_avg` +0.108 мс / jitter / loss, статус `unknown → up` з причиною `icmp` і рівно одним +переходом в історії, heartbeat із RSS 11.6 МБ і `dropped_samples=0`, зонд +позначений `offline` після зупинки, clock skew −0.9 мс. + +### Знайдено під час перевірки + +Приведення типу на місці (`$1::text`) **не розв'язує** конфлікт виведення типів у +Postgres, а нав'язує тип обом уживанням параметра. Коли `$1` потрібен і як `uuid` +для колонки, і як `text` для конкатенації, кастувати треба протилежне уживання: +`VALUES ($1::uuid, …, '/шлях/' || $1 || '.git')`. Та сама пастка двічі: у сіді +тесту (`VALUES ($1, $1, …)` для `core.slug` і `text`) і в `ncm.repos`. + +### Чого ще немає + +- Git-двигун (libgit2): тіло конфігу шифрується в `core.secrets`, + `commit_sha` тимчасово = hex контентного хеша. Дедуплікація й `prev_config_id` + працюють, тож diff будується вже зараз. +- `EnrollmentService` — зонди заводяться вставкою в `core.agents`. +- Сервер не надсилає `TaskDelta` (лише повний план) і не ініціює `ConfigJob`. + +### Далі -- Серверна сторона `AgentService` поверх схеми Етапу 1: запис телеметрії в - hypertables, резолвер `series_ref` → `ts.series.id`, планувальник, що - роздає `TaskPlan` із `schedule_offset`. - Модуль topology на агенті: LLDP/CDP/ARP/FDB → `topo.neighbors`. -- `EnrollmentService`: видача сертифікатів зондам. +- `EnrollmentService` + видача сертифікатів. +- Планувальник NCM: бекап за cron і за Syslog-подією. +- REST/WebSocket API для UI поверх тієї ж БД. diff --git a/agent/internal/config/config.go b/agent/internal/config/config.go index e83a311..713bbd8 100644 --- a/agent/internal/config/config.go +++ b/agent/internal/config/config.go @@ -22,6 +22,9 @@ type Config struct { Endpoint string AgentID string Hostname string + // Токен зонда, виданий при Enroll. Визначає, ЯКИЙ це зонд; + // сертифікат mTLS — що він узагалі має право говорити з сервером. + Token string // mTLS. Порожні шляхи + Insecure = локальний стенд без шифрування. CertFile string @@ -70,6 +73,7 @@ func Parse(args []string) (*Config, error) { fs.StringVar(&c.AgentID, "agent-id", envOr("NETPULSE_AGENT_ID", ""), "ідентифікатор зонда") fs.StringVar(&c.Hostname, "hostname", envOr("NETPULSE_HOSTNAME", host), "ім'я хоста для журналу сервера") + fs.StringVar(&c.Token, "token", envOr("NETPULSE_TOKEN", ""), "токен зонда з Enroll") 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 сервера") @@ -112,6 +116,9 @@ func (c *Config) validate() error { if c.AgentID == "" { return errors.New("не вказано -agent-id: зонд має бути зареєстрований через Enroll") } + if c.Token == "" { + return errors.New("не вказано -token (або NETPULSE_TOKEN)") + } if !c.Insecure && (c.CertFile == "" || c.KeyFile == "") { return errors.New("потрібні -cert і -key для mTLS (або -insecure для локального стенду)") } diff --git a/agent/internal/session/session.go b/agent/internal/session/session.go index 0e32b9c..7351207 100644 --- a/agent/internal/session/session.go +++ b/agent/internal/session/session.go @@ -205,6 +205,11 @@ func (s *Session) runOnce(ctx context.Context) error { return fmt.Errorf("замість Welcome прийшло %T", first.Payload) } s.applyWelcome(welcome) + // Таблиця series_ref на сервері живе рівно стільки, скільки сесія. + // Перереєстровуємо всі відомі серії, замість обнуляти нумерацію: + // інакше довелося б викинути буфер, накопичений за час обриву, — + // саме ті дані, заради яких він і накопичувався. + s.cfg.Buffer.Interner().MarkAllPending() s.connected.Store(true) s.log.Info("сесію встановлено", "session_id", welcome.SessionId, diff --git a/agent/internal/telemetry/interner.go b/agent/internal/telemetry/interner.go index 391de92..70a53a9 100644 --- a/agent/internal/telemetry/interner.go +++ b/agent/internal/telemetry/interner.go @@ -23,16 +23,17 @@ type Interner struct { next uint32 // канонічний ключ серії → присвоєний ref refs map[string]uint32 + // усі створені дескриптори за ref — потрібні, щоб після + // реконекту перереєструвати серії, не втрачаючи буфер + descs map[uint32]*npv1.SeriesDescriptor // дескриптори, які ще не поїхали на сервер pending []*npv1.SeriesDescriptor - // плагін-власник, потрібен лише для заповнення дескриптора - pluginOf map[uint32]string } func NewInterner() *Interner { return &Interner{ - refs: make(map[string]uint32), - pluginOf: make(map[uint32]string), + refs: make(map[string]uint32), + descs: make(map[uint32]*npv1.SeriesDescriptor), } } @@ -54,7 +55,6 @@ func (i *Interner) Ref(deviceID, pluginKey string, m module.Metric) (uint32, boo i.next++ ref := i.next i.refs[key] = ref - i.pluginOf[ref] = pluginKey // Копіюємо labels: модуль може перевикористати мапу між викликами. var labels map[string]string @@ -65,7 +65,7 @@ func (i *Interner) Ref(deviceID, pluginKey string, m module.Metric) (uint32, boo } } - i.pending = append(i.pending, &npv1.SeriesDescriptor{ + desc := &npv1.SeriesDescriptor{ SeriesRef: ref, DeviceId: deviceID, InterfaceId: m.InterfaceID, @@ -73,7 +73,9 @@ func (i *Interner) Ref(deviceID, pluginKey string, m module.Metric) (uint32, boo MetricKey: m.MetricKey, Unit: m.Unit, Labels: labels, - }) + } + i.descs[ref] = desc + i.pending = append(i.pending, desc) return ref, true } @@ -102,14 +104,35 @@ func (i *Interner) ReturnPending(descs []*npv1.SeriesDescriptor) { i.pending = append(descs, i.pending...) } -// Reset обнуляє таблицю: нова сесія або reset_series_table від сервера. +// MarkAllPending повертає ВСІ відомі серії в чергу на відправку. +// +// Викликається на початку кожної нової сесії. Таблиця series_ref живе +// на сервері рівно стільки, скільки сесія, тому після реконекту він +// нічого не пам'ятає. Альтернатива — обнулити нумерацію на агенті — +// означала б викинути весь накопичений за час обриву буфер: саме ті +// дані, заради яких він і накопичувався. +func (i *Interner) MarkAllPending() { + i.mu.Lock() + defer i.mu.Unlock() + + if len(i.descs) == 0 { + return + } + pending := make([]*npv1.SeriesDescriptor, 0, len(i.descs)) + for _, d := range i.descs { + pending = append(pending, d) + } + i.pending = pending +} + +// Reset обнуляє таблицю: сервер попросив перереєстрацію з нуля. 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.descs = make(map[uint32]*npv1.SeriesDescriptor) i.pending = nil } diff --git a/server/README.md b/server/README.md new file mode 100644 index 0000000..75bb02f --- /dev/null +++ b/server/README.md @@ -0,0 +1,146 @@ +# NetPulse Server — AgentService + +Приймальна сторона: зонди підключаються сюди самі, сервер до них не ходить. + +```bash +go build -trimpath -ldflags "-s -w" -o netpulse-server ./cmd/netpulse-server +``` + +```bash +./netpulse-server -listen :9443 -dsn "postgres://..." -dek "k1=" -cert srv.pem -key srv.key -client-ca agents-ca.pem +``` + +## Будова + +``` +cmd/netpulse-server/ точка входу, TLS, keepalive, m'яка зупинка +internal/ + crypto/ AES-GCM-256 для core.secrets, keyring із ротацією + store/ + store.go пул pgx, InTenantTx із SET LOCAL app.tenant_id + agents.go автентифікація зонда за токеном, heartbeat, статуси задач + plan.go побудова TaskPlan, детермінований schedule_offset, хеш плану + credentials.go розшифровка секретів і видача комплекту з TTL + telemetry.go резолвер серій + запис у гіпертаблиці + discovery.go сусіди → зіставлення з інвентарем → topo.links + ncm.go прийом конфігів, дедуплікація, шифроване тіло + grpcapi/ AgentService: Control, StreamTelemetry, StreamLogs, + ReportDiscovery, UploadConfig + перехоплювачі автентифікації +``` + +## Рішення, які варто розуміти + +**Ізоляція тенантів робиться двічі.** RLS у БД (через `SET LOCAL app.tenant_id`) плюс +явний предикат `tenant_id` у кожному запиті. Дублювання не зайве: RLS **не працює на +гіпертаблях** — TimescaleDB не поєднує row level security зі стисненням, а саме туди +йде вся телеметрія. Для неї другий механізм єдиний, тому він не «про всяк випадок». + +**`SET LOCAL`, а не `SET`.** Значення живе до кінця транзакції й не протікає на +наступний запит, який візьме те саме з'єднання з пулу. Протікання тут означало б показ +чужих даних. + +**`schedule_offset` рахує сервер, детерміновано від `check_id`** (`fnv64 % interval`). +Функція чиста, тому будь-який вузол сервера дає те саме значення, а зонд після +перезапуску повертається у свій слот. Інакше 5000 чеків з інтервалом 60 с раз на +хвилину били б сплеском. + +**Хеш плану рахується з полів, що впливають на поведінку** — id, тип, параметри, +інтервали. Зміна опису пристрою не змушує переливати 50 000 задач після кожного обриву +зв'язку. Агент шле хеш у `Hello`, сервер відповідає `task_plan_follows=false`, якщо він +збігся. + +**Токен каже, ЯКИЙ це зонд; сертифікат — що він має право говорити.** Обидва потрібні. +Зонд, який назвався чужим `agent_id` при валідному токені, отримує `PermissionDenied`: +інакше він забрав би чужий план і чужі креденшели. + +**Креденшели розшифровуються на сервері.** DEK не покидає сервер; у БД лише `key_id`, +тому дамп бази без доступу до ключів не дає жодного пароля від обладнання клієнта. +Комплект їде з TTL в годину — після відкликання доступу зонд перестане ним +користуватись сам, навіть якщо зв'язок із ним втрачено. + +**Один нечитабельний секрет не рушить весь комплект.** Решта пристроїв мусить +опитуватись далі. + +**Запис телеметрії — `ON CONFLICT DO NOTHING`.** Семантика доставки at-least-once, тож +повторний батч після реконекту не має ані падати, ані дублювати рядки. Первинні ключі +`(ts, device_id)`, `(ts, interface_id)`, `(ts, series_id)` роблять це безпечним. + +**Статус пристрою змінюється одним запитом із умовним записом в історію.** Не +read-modify-write: два воркери, що обробляють сусідні батчі, наввипередки писали б +неіснуючі переходи `up→up`. + +**Автовиявлення інтерпретує сервер.** Агент доповідає лише «на порту X бачу chassis Y». +Впевненість зіставлення спадає за надійністю ознаки: chassis-id (95) → MAC (90) → +IP керування (80) → sysName (60). Останнє низьке навмисно: `sysName` вводить людина, і +на двох комутаторах цілком може бути `switch`. Лінк із `is_pinned` автовиявлення не +чіпає — інакше кожен запуск затирав би ручні правки. + +## Параметри + +| Прапорець | Змінна | Призначення | +|-----------|--------|-------------| +| `-listen` | `NETPULSE_LISTEN` | адреса gRPC, типово `:9443` | +| `-dsn` | `NETPULSE_DSN` | PostgreSQL | +| `-dek` | `NETPULSE_DEK` | ключі шифрування `id=[,...]` | +| `-cert` / `-key` | `NETPULSE_CERT` / `_KEY` | сертифікат сервера | +| `-client-ca` | `NETPULSE_CLIENT_CA` | CA зондів; вмикає mTLS | +| `-insecure` | `NETPULSE_INSECURE=1` | без TLS, лише локальний стенд | + +## Стан перевірки + +Стенд Debian 13 / PostgreSQL 17.11 / TimescaleDB 2.29.1 / Go 1.25. +`go vet` чисто, `go test -race` — усі тести проходять. + +Інтеграційні тести працюють проти **справжньої БД зі схемою Етапу 1** і справжнього +gRPC (`NETPULSE_TEST_DSN`; без змінної пропускаються): + +| Тест | Що доводить | +|------|-------------| +| `TestControlHandshake` | план зібрано з `core.checks`, offset детермінований і в межах інтервалу, креденшели розшифрувались тим самим ключем і AAD, зонд позначений online | +| `TestControlRejectsMismatchedAgentID` | чужий `agent_id` при валідному токені відхилено | +| `TestUnauthenticatedRejected` | без токена сесії немає | +| `TestTelemetryPersisted` | семпл, ICMP і `util_out_pct` у гіпертаблицях; статус пристрою піднявся; **повтор батчу не продублював ані рядки, ані переходи в історії** | +| `TestTelemetryUnknownSeriesRef` | невідомий ref → `reset_series_table`, а не тихе відкидання | +| `TestDiscoveryResolvesLink` | сусід зіставлений за chassis-id, лінк створено з портами й `capacity_bps`; **зустрічний звіт B→A не створив дублікат** | +| `TestDiscoveryRespectsPinnedLink` | ручний лінк не затерто | +| `TestConfigUploadAndDedup` | чанки склеєні, тіло збережено **зашифрованим** і читається назад; повторний збір не створює версію | +| `TestConfigUploadRejectsBadChecksum` | зіпсований конфіг не потрапляє в базу | +| `TestHeartbeatPersisted` | самометрики в `ts.agent_health`, `dropped_samples` видно у зведенні зонда | +| `TestPlanHashSkipsResend` | збіг хеша → сервер не шле план | + +### Живий наскрізний прогін + +Справжній `netpulse-agent` проти справжнього `netpulse-server` і живої БД, +40 секунд, чек `icmp.ping` кожні 5 с проти `127.0.0.1`: + +| Що перевірено | Результат | +|---------------|-----------| +| Зонд підключився, отримав план | `tasks=1, plan_unchanged=false` | +| ICMP-семпли в `ts.icmp_samples` | 5 за 25 с — рівно за розкладом | +| Узагальнені метрики | `icmp.rtt_avg` 0.108 мс, `icmp.jitter`, `icmp.loss_pct` | +| Статус пристрою | `unknown → up`, причина `icmp`, **рівно один перехід в історії** | +| Heartbeat | записано, RSS зонда 11.6 МБ, `dropped_samples=0` | +| Версія/ОС/архітектура зонда | зафіксовані в `core.agents` | +| Розрив сесії | зонд позначений `offline` | +| Розсинхронізація годинника | −0.9 мс | + +### Знайдено під час перевірки + +Приведення типу на місці (`$1::text`) **не розв'язує** конфлікт виведення типів у +Postgres, а нав'язує тип обом уживанням параметра. Коли `$1` потрібен і як `uuid` для +колонки, і як `text` для конкатенації, кастувати треба протилежне уживання: +`VALUES ($1::uuid, …, '/шлях/' || $1 || '.git')`. + +## Чого ще немає + +- **Git-двигун (libgit2) не підключено.** Тіло конфігу зберігається зашифрованим у + `core.secrets`, а `ncm.configs.commit_sha` тимчасово містить hex контентного хеша. + Дедуплікація, підрахунок рядків і ланцюжок `prev_config_id` працюють уже зараз, тож + diff між версіями будується без Git. Коли двигун з'явиться, зміниться лише джерело + `commit_sha`. +- **`EnrollmentService` не реалізовано** — зонди поки заводяться вставкою в + `core.agents` з `sha256` токена. +- **`TaskDelta` сервер не надсилає**: при зміні плану поки йде повна синхронізація. + Механізм на боці агента вже є. +- **`ConfigJob` / `ConfigApplyJob` сервер не ініціює** — планувальник бекапів за cron + і Syslog-подією ще не написаний. diff --git a/server/go.mod b/server/go.mod new file mode 100644 index 0000000..fd922d1 --- /dev/null +++ b/server/go.mod @@ -0,0 +1,12 @@ +module github.com/netpulse/netpulse/server + +go 1.24 + +require ( + github.com/jackc/pgx/v5 v5.7.6 + github.com/netpulse/netpulse/gen/go v0.0.0 + google.golang.org/grpc v1.76.0 + google.golang.org/protobuf v1.36.12 +) + +replace github.com/netpulse/netpulse/gen/go => ../gen/go diff --git a/server/go.sum b/server/go.sum new file mode 100644 index 0000000..a9cc096 --- /dev/null +++ b/server/go.sum @@ -0,0 +1,64 @@ +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.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +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/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.7.6 h1:rWQc5FwZSPX58r1OQmkuaNicxdmExaEz5A2DO2hUuTk= +github.com/jackc/pgx/v5 v5.7.6/go.mod h1:aruU7o91Tc2q2cFp5h4uP3f6ztExVpyVv88Xl/8Vl8M= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +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/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk= +github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= +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/crypto v0.51.0 h1:IBPXwPfKxY7cWQZ38ZCIRPI50YLeevDLlLnyC5wRGTI= +golang.org/x/crypto v0.51.0/go.mod h1:8AdwkbraGNABw2kOX6YFPs3WM22XqI4EXEd8g+x7Oc8= +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/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= +golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +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/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/server/internal/crypto/secrets.go b/server/internal/crypto/secrets.go new file mode 100644 index 0000000..dce90ec --- /dev/null +++ b/server/internal/crypto/secrets.go @@ -0,0 +1,173 @@ +// Package crypto — шифрування секретів для core.secrets. +// +// AES-GCM-256. Ключі шифрування даних (DEK) живуть поза БД: у +// KMS/Vault або, для self-hosted, у файлі з правами 0600. У таблиці +// лежить лише key_id, тому дамп бази без доступу до ключів не дає +// жодного пароля від обладнання клієнта. +package crypto + +import ( + "crypto/aes" + "crypto/cipher" + "crypto/rand" + "errors" + "fmt" + "sync" +) + +const ( + KeySize = 32 // AES-256 + NonceSize = 12 // 96 біт, рекомендація NIST для GCM + TagSize = 16 +) + +var ( + ErrUnknownKey = errors.New("невідомий key_id") + ErrBadKeySize = errors.New("ключ має бути 32 байти") +) + +// Keyring — набір DEK за ідентифікаторами. Кілька ключів одночасно +// потрібні для ротації: нові секрети шифруються активним ключем, +// старі ще читаються попереднім. +type Keyring struct { + mu sync.RWMutex + keys map[string][]byte + active string +} + +func NewKeyring() *Keyring { + return &Keyring{keys: make(map[string][]byte)} +} + +// Add додає ключ. Перший доданий стає активним. +func (k *Keyring) Add(keyID string, key []byte) error { + if len(key) != KeySize { + return fmt.Errorf("%w: %d", ErrBadKeySize, len(key)) + } + k.mu.Lock() + defer k.mu.Unlock() + + cp := make([]byte, KeySize) + copy(cp, key) + k.keys[keyID] = cp + if k.active == "" { + k.active = keyID + } + return nil +} + +// SetActive призначає ключ, яким шифруються нові секрети. +func (k *Keyring) SetActive(keyID string) error { + k.mu.Lock() + defer k.mu.Unlock() + if _, ok := k.keys[keyID]; !ok { + return ErrUnknownKey + } + k.active = keyID + return nil +} + +func (k *Keyring) ActiveKeyID() string { + k.mu.RLock() + defer k.mu.RUnlock() + return k.active +} + +func (k *Keyring) aead(keyID string) (cipher.AEAD, error) { + k.mu.RLock() + key, ok := k.keys[keyID] + k.mu.RUnlock() + if !ok { + return nil, fmt.Errorf("%w: %q", ErrUnknownKey, keyID) + } + + block, err := aes.NewCipher(key) + if err != nil { + return nil, err + } + return cipher.NewGCM(block) +} + +// Secret — те, що лягає в рядок core.secrets. +type Secret struct { + KeyID string + Nonce []byte + Ciphertext []byte + AuthTag []byte +} + +// Encrypt шифрує активним ключем. +// +// aad (additional authenticated data) прив'язує шифротекст до +// контексту — зазвичай tenant_id||kind||object_id. Без цього рядок +// із секретом можна було б переставити на інший пристрій або тенант +// прямим UPDATE у БД, і розшифрування б відпрацювало. +func (k *Keyring) Encrypt(plaintext []byte, aad string) (*Secret, error) { + keyID := k.ActiveKeyID() + if keyID == "" { + return nil, errors.New("keyring порожній") + } + + aead, err := k.aead(keyID) + if err != nil { + return nil, err + } + + nonce := make([]byte, NonceSize) + if _, err := rand.Read(nonce); err != nil { + return nil, err + } + + sealed := aead.Seal(nil, nonce, plaintext, []byte(aad)) + if len(sealed) < TagSize { + return nil, errors.New("несподівано короткий шифротекст") + } + + // Схема тримає тег окремо від шифротексту, тому розділяємо: + // Go дописує його в кінець. + return &Secret{ + KeyID: keyID, + Nonce: nonce, + Ciphertext: sealed[:len(sealed)-TagSize], + AuthTag: sealed[len(sealed)-TagSize:], + }, nil +} + +// Decrypt розшифровує й перевіряє цілісність. +func (k *Keyring) Decrypt(s *Secret, aad string) ([]byte, error) { + if s == nil { + return nil, errors.New("порожній секрет") + } + if len(s.Nonce) != NonceSize { + return nil, fmt.Errorf("nonce має бути %d байт, отримано %d", NonceSize, len(s.Nonce)) + } + if len(s.AuthTag) != TagSize { + return nil, fmt.Errorf("auth_tag має бути %d байт, отримано %d", TagSize, len(s.AuthTag)) + } + + aead, err := k.aead(s.KeyID) + if err != nil { + return nil, err + } + + sealed := make([]byte, 0, len(s.Ciphertext)+TagSize) + sealed = append(sealed, s.Ciphertext...) + sealed = append(sealed, s.AuthTag...) + + out, err := aead.Open(nil, s.Nonce, sealed, []byte(aad)) + if err != nil { + // Не розкриваємо, що саме не зійшлось: помилка автентифікації + // GCM не має ставати оракулом. + return nil, errors.New("не вдалося розшифрувати секрет") + } + return out, nil +} + +// GenerateKey — новий випадковий DEK. +func GenerateKey() ([]byte, error) { + key := make([]byte, KeySize) + if _, err := rand.Read(key); err != nil { + return nil, err + } + return key, nil +} diff --git a/server/internal/grpcapi/integration_test.go b/server/internal/grpcapi/integration_test.go new file mode 100644 index 0000000..ce5b38e --- /dev/null +++ b/server/internal/grpcapi/integration_test.go @@ -0,0 +1,959 @@ +// Наскрізний тест серверної сторони: справжній PostgreSQL/TimescaleDB, +// справжній gRPC, справжня схема з Етапу 1. +// +// Запуск: +// +// NETPULSE_TEST_DSN="postgres://netpulse:netpulse@localhost/netpulse_it" go test ./... +// +// Без змінної тест пропускається — щоб `go test ./...` не падав там, +// де бази немає. +package grpcapi_test + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "log/slog" + "net" + "os" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgxpool" + "github.com/netpulse/netpulse/server/internal/crypto" + "github.com/netpulse/netpulse/server/internal/grpcapi" + "github.com/netpulse/netpulse/server/internal/store" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/metadata" + "google.golang.org/grpc/test/bufconn" + "google.golang.org/protobuf/types/known/durationpb" + "google.golang.org/protobuf/types/known/timestamppb" +) + +type fixture struct { + pool *pgxpool.Pool + store *store.Store + ring *crypto.Keyring + client npv1.AgentServiceClient + ctx context.Context + tenantID string + agentID string + deviceID string + ifaceID string + peerID string + peerIfID string + checkID string + token string +} + +func setup(t *testing.T) *fixture { + t.Helper() + + dsn := os.Getenv("NETPULSE_TEST_DSN") + if dsn == "" { + t.Skip("NETPULSE_TEST_DSN не задано — інтеграційний тест пропущено") + } + + ctx := context.Background() + + st, err := store.New(ctx, dsn) + if err != nil { + t.Fatalf("підключення до БД: %v", err) + } + t.Cleanup(st.Close) + + ring := crypto.NewKeyring() + key, err := crypto.GenerateKey() + if err != nil { + t.Fatalf("GenerateKey: %v", err) + } + if err := ring.Add("test-key", key); err != nil { + t.Fatalf("Add: %v", err) + } + + f := &fixture{pool: st.Pool(), store: st, ring: ring, ctx: ctx} + f.seed(t) + + // Сервер із тими самими перехоплювачами, що й у проді. + svc := grpcapi.New(st, ring, + slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelWarn}))) + + lis := bufconn.Listen(1 << 20) + srv := grpc.NewServer( + grpc.ChainUnaryInterceptor(svc.UnaryInterceptor), + grpc.ChainStreamInterceptor(svc.StreamInterceptor), + ) + npv1.RegisterAgentServiceServer(srv, svc) + go func() { _ = srv.Serve(lis) }() + + conn, err := grpc.NewClient("passthrough:///bufnet", + grpc.WithContextDialer(func(ctx context.Context, _ string) (net.Conn, error) { + return lis.DialContext(ctx) + }), + grpc.WithTransportCredentials(insecure.NewCredentials()), + ) + if err != nil { + t.Fatalf("клієнт: %v", err) + } + t.Cleanup(func() { + _ = conn.Close() + srv.Stop() + _ = lis.Close() + }) + + f.client = npv1.NewAgentServiceClient(conn) + return f +} + +// authCtx додає токен зонда так само, як це робить агент. +func (f *fixture) authCtx() context.Context { + return metadata.AppendToOutgoingContext(f.ctx, "authorization", "Bearer "+f.token) +} + +func (f *fixture) seed(t *testing.T) { + t.Helper() + ctx := f.ctx + + slug := fmt.Sprintf("it-%d", time.Now().UnixNano()) + f.token = "np_test_" + slug + sum := sha256.Sum256([]byte(f.token)) + + must := func(q string, args ...any) { + t.Helper() + if _, err := f.pool.Exec(ctx, q, args...); err != nil { + t.Fatalf("seed %q: %v", q, err) + } + } + scan := func(dest *string, q string, args ...any) { + t.Helper() + if err := f.pool.QueryRow(ctx, q, args...).Scan(dest); err != nil { + t.Fatalf("seed %q: %v", q, err) + } + } + + // slug має домен core.slug, name — text. Один і той самий $1 для + // обох Postgres вивести не може, тому передаємо окремо. + scan(&f.tenantID, ` + INSERT INTO core.tenants (slug, name, status) VALUES ($1, $2, 'active') + RETURNING id::text`, slug, slug) + + t.Cleanup(func() { + // Каскади прибирають усе, крім гіпертаблиць — у них немає FK. + _, _ = f.pool.Exec(context.Background(), + `DELETE FROM ts.icmp_samples WHERE tenant_id = $1`, f.tenantID) + _, _ = f.pool.Exec(context.Background(), + `DELETE FROM ts.if_counters WHERE tenant_id = $1`, f.tenantID) + _, _ = f.pool.Exec(context.Background(), + `DELETE FROM ts.device_status_history WHERE tenant_id = $1`, f.tenantID) + _, _ = f.pool.Exec(context.Background(), + `DELETE FROM ts.agent_health WHERE tenant_id = $1`, f.tenantID) + _, _ = f.pool.Exec(context.Background(), + `DELETE FROM ts.samples WHERE series_id IN (SELECT id FROM ts.series WHERE tenant_id = $1)`, f.tenantID) + _, _ = f.pool.Exec(context.Background(), + `DELETE FROM core.tenants WHERE id = $1`, f.tenantID) + }) + + scan(&f.agentID, ` + INSERT INTO core.agents (tenant_id, name, token_hash, status, enabled_modules) + VALUES ($1, 'probe-it', $2, 'online', '{icmp,snmp,topology}') + RETURNING id::text`, f.tenantID, sum[:]) + + scan(&f.deviceID, ` + INSERT INTO inv.devices (tenant_id, agent_id, name, address, kind, vendor, chassis_id, system_name) + VALUES ($1, $2, 'core-sw-01', '10.10.0.1', 'switch', 'cisco', '0011.2233.4455', 'core-sw-01') + RETURNING id::text`, f.tenantID, f.agentID) + + // Сусід, якого автовиявлення має знайти. + scan(&f.peerID, ` + INSERT INTO inv.devices (tenant_id, agent_id, name, address, kind, chassis_id, system_name) + VALUES ($1, $2, 'edge-rtr-01', '10.10.0.2', 'router', '0011.2233.6677', 'edge-rtr-01') + RETURNING id::text`, f.tenantID, f.agentID) + + scan(&f.ifaceID, ` + INSERT INTO inv.interfaces (tenant_id, device_id, if_index, name, speed_bps, oper_status) + VALUES ($1, $2, 1, 'GigabitEthernet0/1', 1000000000, 'up') + RETURNING id::text`, f.tenantID, f.deviceID) + + scan(&f.peerIfID, ` + INSERT INTO inv.interfaces (tenant_id, device_id, if_index, name, speed_bps, oper_status) + VALUES ($1, $2, 1, 'ether1', 1000000000, 'up') + RETURNING id::text`, f.tenantID, f.peerID) + + scan(&f.checkID, ` + INSERT INTO core.checks (tenant_id, device_id, check_type, params, interval_sec, timeout_ms, retries) + VALUES ($1, $2, 'icmp.ping', '{"count":3}'::jsonb, 60, 3000, 2) + RETURNING id::text`, f.tenantID, f.deviceID) + + // Креденшел зі справжнім шифруванням: тест має пройти той самий + // шлях, що й прод, включно з AAD. + aad := f.tenantID + "|snmp_v2c|" + f.deviceID + sec, err := f.ring.Encrypt([]byte("public-it"), aad) + if err != nil { + t.Fatalf("Encrypt: %v", err) + } + + var secretID string + scan(&secretID, ` + INSERT INTO core.secrets (tenant_id, kind, key_id, nonce, ciphertext, auth_tag, aad) + VALUES ($1, 'snmp_v3', $2, $3, $4, $5, $6) + RETURNING id::text`, + f.tenantID, sec.KeyID, sec.Nonce, sec.Ciphertext, sec.AuthTag, aad) + + var credID string + scan(&credID, ` + INSERT INTO inv.credentials (tenant_id, name, proto, username, port, secret_id) + VALUES ($1, 'snmp-ro', 'snmp_v2c', '', 161, $2) + RETURNING id::text`, f.tenantID, secretID) + + must(`INSERT INTO inv.device_credentials (device_id, credential_id, priority) + VALUES ($1, $2, 10)`, f.deviceID, credID) +} + +// --------------------------------------------------------------------- +// 1. Рукостискання, план задач, креденшели +// --------------------------------------------------------------------- + +func TestControlHandshake(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second) + defer cancel() + + stream, err := f.client.Control(ctx) + if err != nil { + t.Fatalf("Control: %v", err) + } + + if err := stream.Send(&npv1.ControlUp{ + Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{ + AgentId: f.agentID, + Hostname: "probe-it", + Build: &npv1.AgentBuild{Version: "test", Os: "linux", Arch: "amd64"}, + }}, + }); err != nil { + t.Fatalf("Hello: %v", err) + } + + down, err := stream.Recv() + if err != nil { + t.Fatalf("Welcome: %v", err) + } + w := down.GetWelcome() + if w == nil { + t.Fatalf("замість Welcome прийшло %T", down.Payload) + } + if !w.TaskPlanFollows { + t.Fatal("сервер не збирається слати план, хоча агент прийшов без хеша") + } + if w.TelemetryMaxBatchSize == 0 || w.MaxConcurrentChecks == 0 { + t.Fatalf("ліміти не заповнені: %+v", w) + } + + var ( + plan *npv1.TaskPlan + creds *npv1.CredentialBundle + mods *npv1.ModuleControl + ) + deadline := time.Now().Add(10 * time.Second) + for (plan == nil || creds == nil || mods == nil) && time.Now().Before(deadline) { + msg, err := stream.Recv() + if err != nil { + t.Fatalf("Recv: %v", err) + } + switch p := msg.Payload.(type) { + case *npv1.ControlDown_TaskPlan: + plan = p.TaskPlan + case *npv1.ControlDown_Credentials: + creds = p.Credentials + case *npv1.ControlDown_ModuleControl: + mods = p.ModuleControl + } + } + + if plan == nil { + t.Fatal("план задач не надійшов") + } + if len(plan.Tasks) != 1 || plan.Tasks[0].CheckId != f.checkID { + t.Fatalf("план зібрано неправильно: %+v", plan.Tasks) + } + task := plan.Tasks[0] + if task.CheckType != "icmp.ping" { + t.Fatalf("check_type = %q", task.CheckType) + } + if task.Interval.AsDuration() != time.Minute { + t.Fatalf("interval = %v", task.Interval.AsDuration()) + } + // Offset має бути в межах інтервалу — інакше задача або стартує + // одразу разом з усіма, або з'їде в наступний цикл. + off := task.ScheduleOffset.AsDuration() + if off < 0 || off >= time.Minute { + t.Fatalf("schedule_offset поза інтервалом: %v", off) + } + // І він має бути детермінованим: та сама функція на будь-якому вузлі. + if want := store.ScheduleOffset(f.checkID, time.Minute); off != want { + t.Fatalf("offset не детермінований: %v != %v", off, want) + } + + if len(plan.Devices) != 1 || plan.Devices[0].Address != "10.10.0.1" { + t.Fatalf("пристрої в плані: %+v", plan.Devices) + } + if len(plan.PlanHash) == 0 { + t.Fatal("план без хеша — агент не зможе уникнути перезаливки") + } + + if mods == nil || len(mods.Modules) == 0 { + t.Fatal("сервер не сказав, які модулі вмикати") + } + + if creds == nil { + t.Fatal("креденшели не надійшли") + } + list := creds.ByDevice[f.deviceID] + if list == nil || len(list.Credentials) != 1 { + t.Fatalf("креденшели пристрою: %+v", creds.ByDevice) + } + // Найважливіше: секрет розшифрувався тим самим ключем і AAD. + if got := list.Credentials[0].GetCommunity(); got != "public-it" { + t.Fatalf("community = %q, очікували public-it", got) + } + if creds.ExpiresAt == nil || creds.ExpiresAt.AsTime().Before(time.Now()) { + t.Fatal("креденшели без дійсного TTL") + } + + // Зонд позначений онлайн у БД. + var status string + if err := f.pool.QueryRow(f.ctx, + `SELECT status::text FROM core.agents WHERE id = $1`, f.agentID).Scan(&status); err != nil { + t.Fatalf("статус зонда: %v", err) + } + if status != "online" { + t.Fatalf("статус зонда = %q", status) + } +} + +// Токен визначає зонда. Назватись чужим id не можна навіть із валідним +// токеном — інакше один зонд отримав би план іншого. +func TestControlRejectsMismatchedAgentID(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.authCtx(), 10*time.Second) + defer cancel() + + stream, err := f.client.Control(ctx) + if err != nil { + t.Fatalf("Control: %v", err) + } + if err := stream.Send(&npv1.ControlUp{ + Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{ + AgentId: "00000000-0000-0000-0000-0000000000ff", + }}, + }); err != nil { + t.Fatalf("Hello: %v", err) + } + if _, err := stream.Recv(); err == nil { + t.Fatal("сервер прийняв чужий agent_id") + } +} + +func TestUnauthenticatedRejected(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.ctx, 10*time.Second) + defer cancel() + + stream, err := f.client.Control(ctx) + if err != nil { + return // клієнт міг відмовити одразу + } + _ = stream.Send(&npv1.ControlUp{Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{}}}) + if _, err := stream.Recv(); err == nil { + t.Fatal("сервер пустив сесію без токена") + } +} + +// --------------------------------------------------------------------- +// 2. Телеметрія доходить до гіпертаблиць +// --------------------------------------------------------------------- + +func TestTelemetryPersisted(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second) + defer cancel() + + stream, err := f.client.StreamTelemetry(ctx) + if err != nil { + t.Fatalf("StreamTelemetry: %v", err) + } + + now := time.Now().Truncate(time.Millisecond) + + batch := &npv1.TelemetryBatch{ + BatchId: 1, + AgentId: f.agentID, + CreatedAt: timestamppb.New(now), + NewSeries: []*npv1.SeriesDescriptor{{ + SeriesRef: 1, + DeviceId: f.deviceID, + PluginKey: "snmp", + MetricKey: "cpu.util", + Unit: "pct", + Labels: map[string]string{"core": "0"}, + }}, + Samples: []*npv1.MetricSample{ + {SeriesRef: 1, Ts: timestamppb.New(now), Value: 37.5}, + }, + Icmp: []*npv1.IcmpResult{{ + DeviceId: f.deviceID, CheckId: f.checkID, Ts: timestamppb.New(now), + RttAvgMs: 1.25, RttMinMs: 1.1, RttMaxMs: 1.5, JitterMs: 0.2, + LossPct: 0, PacketsSent: 3, PacketsRecv: 3, Reachable: true, + }}, + Interfaces: []*npv1.InterfaceCounters{{ + DeviceId: f.deviceID, InterfaceId: f.ifaceID, Ts: timestamppb.New(now), + InOctets: 1 << 30, OutOctets: 1 << 29, + InBps: 420e6, OutBps: 780e6, + UtilInPct: 42, UtilOutPct: 78, OperUp: true, + Interval: durationpb.New(time.Minute), + }}, + } + + if err := stream.Send(batch); err != nil { + t.Fatalf("Send: %v", err) + } + ack, err := stream.Recv() + if err != nil { + t.Fatalf("Recv ack: %v", err) + } + if ack.AckedThroughBatchId != 1 { + t.Fatalf("acked = %d, помилка: %+v", ack.AckedThroughBatchId, ack.Error) + } + + // Семпл у ts.samples із розв'язаним series_id. + var value float64 + var metricKey, unit string + if err := f.pool.QueryRow(f.ctx, ` + SELECT s.value, se.metric_key, se.unit + FROM ts.samples s JOIN ts.series se ON se.id = s.series_id + WHERE se.tenant_id = $1 AND se.metric_key = 'cpu.util' + ORDER BY s.ts DESC LIMIT 1`, f.tenantID).Scan(&value, &metricKey, &unit); err != nil { + t.Fatalf("семпл не записався: %v", err) + } + if value != 37.5 || unit != "pct" { + t.Fatalf("семпл спотворений: value=%v unit=%q", value, unit) + } + + // ICMP у широкій таблиці — те, що читає мапа. + var rtt float32 + var reachable bool + if err := f.pool.QueryRow(f.ctx, ` + SELECT rtt_avg_ms, reachable FROM ts.icmp_samples + WHERE device_id = $1 ORDER BY ts DESC LIMIT 1`, f.deviceID).Scan(&rtt, &reachable); err != nil { + t.Fatalf("ICMP не записався: %v", err) + } + if !reachable || rtt < 1.2 || rtt > 1.3 { + t.Fatalf("ICMP спотворений: rtt=%v reachable=%v", rtt, reachable) + } + + // Завантаження інтерфейсу — джерело швидкості анімації на мапі. + var utilOut float32 + if err := f.pool.QueryRow(f.ctx, ` + SELECT util_out_pct FROM ts.if_counters + WHERE interface_id = $1 ORDER BY ts DESC LIMIT 1`, f.ifaceID).Scan(&utilOut); err != nil { + t.Fatalf("лічильники не записались: %v", err) + } + if utilOut != 78 { + t.Fatalf("util_out_pct = %v", utilOut) + } + + // Статус пристрою піднявся з ICMP і лишив слід в історії. + var devStatus string + if err := f.pool.QueryRow(f.ctx, + `SELECT status::text FROM inv.devices WHERE id = $1`, f.deviceID).Scan(&devStatus); err != nil { + t.Fatalf("статус пристрою: %v", err) + } + if devStatus != "up" { + t.Fatalf("статус пристрою = %q, очікували up", devStatus) + } + + var transitions int + if err := f.pool.QueryRow(f.ctx, ` + SELECT count(*) FROM ts.device_status_history + WHERE device_id = $1 AND status = 'up'`, f.deviceID).Scan(&transitions); err != nil { + t.Fatalf("історія станів: %v", err) + } + if transitions != 1 { + t.Fatalf("переходів в історії: %d, очікували 1", transitions) + } + + // Повтор того самого батчу (at-least-once після реконекту) не має + // ані падати, ані дублювати рядки. + batch.IsRetransmit = true + batch.BatchId = 2 + batch.NewSeries = nil + if err := stream.Send(batch); err != nil { + t.Fatalf("повторний Send: %v", err) + } + if _, err := stream.Recv(); err != nil { + t.Fatalf("повторний ack: %v", err) + } + + var icmpRows int + if err := f.pool.QueryRow(f.ctx, + `SELECT count(*) FROM ts.icmp_samples WHERE device_id = $1`, f.deviceID).Scan(&icmpRows); err != nil { + t.Fatalf("підрахунок: %v", err) + } + if icmpRows != 1 { + t.Fatalf("повтор батчу продублював рядки: %d", icmpRows) + } + + // І жодного зайвого переходу в історії від повтору. + if err := f.pool.QueryRow(f.ctx, ` + SELECT count(*) FROM ts.device_status_history WHERE device_id = $1`, + f.deviceID).Scan(&transitions); err != nil { + t.Fatalf("історія: %v", err) + } + if transitions != 1 { + t.Fatalf("повтор додав перехід в історію: %d", transitions) + } +} + +// Невідомий series_ref → чесний запит на перереєстрацію. +func TestTelemetryUnknownSeriesRef(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second) + defer cancel() + + stream, err := f.client.StreamTelemetry(ctx) + if err != nil { + t.Fatalf("StreamTelemetry: %v", err) + } + + if err := stream.Send(&npv1.TelemetryBatch{ + BatchId: 5, AgentId: f.agentID, + Samples: []*npv1.MetricSample{{SeriesRef: 404, Ts: timestamppb.Now(), Value: 1}}, + }); err != nil { + t.Fatalf("Send: %v", err) + } + + ack, err := stream.Recv() + if err != nil { + t.Fatalf("Recv: %v", err) + } + if !ack.ResetSeriesTable { + t.Fatal("сервер проковтнув невідомий series_ref") + } + if ack.Error == nil || ack.Error.Code != "unknown_series_ref" { + t.Fatalf("немає діагностики: %+v", ack.Error) + } +} + +// --------------------------------------------------------------------- +// 3. Автовиявлення зводить сусідів у лінк +// --------------------------------------------------------------------- + +func TestDiscoveryResolvesLink(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second) + defer cancel() + + ack, err := f.client.ReportDiscovery(ctx, &npv1.DiscoveryReport{ + AgentId: f.agentID, + Final: true, + Neighbors: []*npv1.NeighborRecord{{ + DeviceId: f.deviceID, + LocalInterfaceId: f.ifaceID, + Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_LLDP, + RemoteChassisId: "0011.2233.6677", + RemoteSystemName: "edge-rtr-01", + RemotePortId: "ether1", + RemoteCapabilities: []string{"bridge", "router"}, + SeenAt: timestamppb.Now(), + }}, + }) + if err != nil { + t.Fatalf("ReportDiscovery: %v", err) + } + if !ack.Accepted { + t.Fatalf("звіт відхилено: %+v", ack.Error) + } + if ack.NeighborsResolved != 1 { + t.Fatalf("зіставлено сусідів: %d", ack.NeighborsResolved) + } + if ack.LinksCreated != 1 { + t.Fatalf("створено лінків: %d", ack.LinksCreated) + } + + // Сусід зберігся з резолвленими посиланнями й високою впевненістю: + // збіг за chassis-id надійніший за збіг за іменем. + var resolvedDev, resolvedIf string + var confidence int + if err := f.pool.QueryRow(f.ctx, ` + SELECT resolved_device_id::text, resolved_interface_id::text, confidence + FROM topo.neighbors WHERE device_id = $1`, f.deviceID). + Scan(&resolvedDev, &resolvedIf, &confidence); err != nil { + t.Fatalf("сусід не записався: %v", err) + } + if resolvedDev != f.peerID { + t.Fatalf("сусід зіставлений не з тим пристроєм: %s", resolvedDev) + } + if resolvedIf != f.peerIfID { + t.Fatalf("порт сусіда не зіставлений: %s", resolvedIf) + } + if confidence < 90 { + t.Fatalf("впевненість збігу за chassis-id = %d, очікували >= 90", confidence) + } + + // Лінк створився з портами й пропускною здатністю з інтерфейсу. + var capacity int64 + var discoveredBy string + if err := f.pool.QueryRow(f.ctx, ` + SELECT COALESCE(capacity_bps, 0), discovered_by::text + FROM topo.links WHERE tenant_id = $1`, f.tenantID).Scan(&capacity, &discoveredBy); err != nil { + t.Fatalf("лінк не створився: %v", err) + } + if capacity != 1000000000 { + t.Fatalf("capacity_bps = %d — анімація трафіку не матиме знаменника", capacity) + } + if discoveredBy != "lldp" { + t.Fatalf("discovered_by = %q", discoveredBy) + } + + // Повторний звіт не має плодити другий лінк — навіть якщо сусід + // прийшов з іншого боку (B→A). + if _, err := f.client.ReportDiscovery(ctx, &npv1.DiscoveryReport{ + AgentId: f.agentID, Final: true, + Neighbors: []*npv1.NeighborRecord{{ + DeviceId: f.peerID, + LocalInterfaceId: f.peerIfID, + Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_LLDP, + RemoteChassisId: "0011.2233.4455", + RemotePortId: "GigabitEthernet0/1", + SeenAt: timestamppb.Now(), + }}, + }); err != nil { + t.Fatalf("повторний звіт: %v", err) + } + + var links int + if err := f.pool.QueryRow(f.ctx, + `SELECT count(*) FROM topo.links WHERE tenant_id = $1`, f.tenantID).Scan(&links); err != nil { + t.Fatalf("підрахунок лінків: %v", err) + } + if links != 1 { + t.Fatalf("зустрічний звіт створив дублікат: лінків %d", links) + } +} + +// Підтверджений людиною лінк не має перезаписуватись автовиявленням. +func TestDiscoveryRespectsPinnedLink(t *testing.T) { + f := setup(t) + + if _, err := f.pool.Exec(f.ctx, ` + INSERT INTO topo.links + (tenant_id, a_device_id, a_interface_id, b_device_id, b_interface_id, + kind, discovered_by, confidence, is_pinned) + VALUES ($1,$2,$3,$4,$5,'physical','manual',100,true) + `, f.tenantID, f.deviceID, f.ifaceID, f.peerID, f.peerIfID); err != nil { + t.Fatalf("seed лінка: %v", err) + } + + ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second) + defer cancel() + + if _, err := f.client.ReportDiscovery(ctx, &npv1.DiscoveryReport{ + AgentId: f.agentID, Final: true, + Neighbors: []*npv1.NeighborRecord{{ + DeviceId: f.deviceID, + LocalInterfaceId: f.ifaceID, + Proto: npv1.DiscoveryProto_DISCOVERY_PROTO_CDP, + RemoteChassisId: "0011.2233.6677", + RemotePortId: "ether1", + SeenAt: timestamppb.Now(), + }}, + }); err != nil { + t.Fatalf("ReportDiscovery: %v", err) + } + + var discoveredBy string + var pinned bool + if err := f.pool.QueryRow(f.ctx, ` + SELECT discovered_by::text, is_pinned FROM topo.links WHERE tenant_id = $1 + `, f.tenantID).Scan(&discoveredBy, &pinned); err != nil { + t.Fatalf("лінк: %v", err) + } + if !pinned || discoveredBy != "manual" { + t.Fatalf("автовиявлення затерло ручний лінк: discovered_by=%q pinned=%v", + discoveredBy, pinned) + } +} + +// --------------------------------------------------------------------- +// 4. NCM: збір конфігу та дедуплікація +// --------------------------------------------------------------------- + +func TestConfigUploadAndDedup(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second) + defer cancel() + + body := []byte("!\nhostname core-sw-01\ninterface Gi0/1\n description uplink\n!\n") + sum := sha256.Sum256(body) + + upload := func(jobID string) *npv1.ConfigReceipt { + t.Helper() + stream, err := f.client.UploadConfig(ctx) + if err != nil { + t.Fatalf("UploadConfig: %v", err) + } + if err := stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Header{ + Header: &npv1.ConfigHeader{ + JobId: jobID, AgentId: f.agentID, DeviceId: f.deviceID, + ConfigType: "running", CollectedAt: timestamppb.Now(), Encoding: "none", + }}}); err != nil { + t.Fatalf("header: %v", err) + } + // Два чанки, щоб перевірити склеювання. + mid := len(body) / 2 + for i, part := range [][]byte{body[:mid], body[mid:]} { + if err := stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Chunk{ + Chunk: &npv1.ConfigChunk{Sequence: uint32(i), Data: part}, + }}); err != nil { + t.Fatalf("chunk %d: %v", i, err) + } + } + if err := stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Trailer{ + Trailer: &npv1.ConfigTrailer{ + Success: true, ContentSha256: sum[:], + SizeBytes: uint64(len(body)), ChunkCount: 2, + }}}); err != nil { + t.Fatalf("trailer: %v", err) + } + receipt, err := stream.CloseAndRecv() + if err != nil { + t.Fatalf("CloseAndRecv: %v", err) + } + return receipt + } + + first := upload("job-1") + if !first.Accepted || first.Unchanged { + t.Fatalf("перший конфіг: accepted=%v unchanged=%v err=%+v", + first.Accepted, first.Unchanged, first.Error) + } + if first.ConfigId == "" { + t.Fatal("немає config_id") + } + + // Другий збір того самого конфігу — не нова версія. + second := upload("job-2") + if !second.Accepted { + t.Fatalf("другий конфіг відхилено: %+v", second.Error) + } + if !second.Unchanged { + t.Fatal("незмінений конфіг створив нову версію — історія засмітиться") + } + + var versions int + if err := f.pool.QueryRow(f.ctx, + `SELECT count(*) FROM ncm.configs WHERE device_id = $1`, f.deviceID).Scan(&versions); err != nil { + t.Fatalf("підрахунок версій: %v", err) + } + if versions != 1 { + t.Fatalf("версій у ncm.configs: %d, очікували 1", versions) + } + + // Тіло збережене зашифрованим і читається назад тим самим ключем. + var keyID, aad string + var nonce, ct, tag []byte + if err := f.pool.QueryRow(f.ctx, ` + SELECT s.key_id, s.nonce, s.ciphertext, s.auth_tag, COALESCE(s.aad,'') + FROM ncm.configs c JOIN core.secrets s ON s.id = c.body_secret_id + WHERE c.device_id = $1`, f.deviceID).Scan(&keyID, &nonce, &ct, &tag, &aad); err != nil { + t.Fatalf("зашифроване тіло: %v", err) + } + plain, err := f.ring.Decrypt(&crypto.Secret{ + KeyID: keyID, Nonce: nonce, Ciphertext: ct, AuthTag: tag, + }, aad) + if err != nil { + t.Fatalf("розшифрування тіла: %v", err) + } + if string(plain) != string(body) { + t.Fatal("тіло конфігу спотворилось при збереженні") + } + + // І воно точно не лежить відкритим текстом. + if string(ct) == string(body) { + t.Fatal("конфіг збережено без шифрування") + } +} + +func TestConfigUploadRejectsBadChecksum(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second) + defer cancel() + + stream, err := f.client.UploadConfig(ctx) + if err != nil { + t.Fatalf("UploadConfig: %v", err) + } + + body := []byte("hostname core-sw-01\n") + bogus := sha256.Sum256([]byte("щось інше")) + + _ = stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Header{ + Header: &npv1.ConfigHeader{JobId: "job-x", DeviceId: f.deviceID, ConfigType: "running"}}}) + _ = stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Chunk{ + Chunk: &npv1.ConfigChunk{Sequence: 0, Data: body}}}) + _ = stream.Send(&npv1.ConfigUpload{Part: &npv1.ConfigUpload_Trailer{ + Trailer: &npv1.ConfigTrailer{Success: true, ContentSha256: bogus[:], + SizeBytes: uint64(len(body)), ChunkCount: 1}}}) + + receipt, err := stream.CloseAndRecv() + if err != nil { + t.Fatalf("CloseAndRecv: %v", err) + } + if receipt.Accepted { + t.Fatal("прийнято конфіг із розбіжною сумою") + } + if receipt.Error == nil || receipt.Error.Code != "checksum_mismatch" { + t.Fatalf("немає причини відмови: %+v", receipt.Error) + } + + var versions int + _ = f.pool.QueryRow(f.ctx, + `SELECT count(*) FROM ncm.configs WHERE device_id = $1`, f.deviceID).Scan(&versions) + if versions != 0 { + t.Fatalf("зіпсований конфіг усе одно збережено: версій %d", versions) + } +} + +// --------------------------------------------------------------------- +// 5. Heartbeat лягає в самометрики +// --------------------------------------------------------------------- + +func TestHeartbeatPersisted(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second) + defer cancel() + + stream, err := f.client.Control(ctx) + if err != nil { + t.Fatalf("Control: %v", err) + } + if err := stream.Send(&npv1.ControlUp{ + Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{AgentId: f.agentID}}, + }); err != nil { + t.Fatalf("Hello: %v", err) + } + if _, err := stream.Recv(); err != nil { + t.Fatalf("Welcome: %v", err) + } + + if err := stream.Send(&npv1.ControlUp{ + Payload: &npv1.ControlUp_Heartbeat{Heartbeat: &npv1.Heartbeat{ + Ts: timestamppb.Now(), + Health: &npv1.AgentHealth{ + CpuPct: 2.5, RssBytes: 12 << 20, Goroutines: 48, + QueueDepth: 3, DroppedSamples: 7, + }, + TasksRunning: 1, + }}, + }); err != nil { + t.Fatalf("Heartbeat: %v", err) + } + + // Дочекаємось Ping у відповідь — він приходить після запису. + deadline := time.Now().Add(10 * time.Second) + for time.Now().Before(deadline) { + msg, err := stream.Recv() + if err != nil { + t.Fatalf("Recv: %v", err) + } + if msg.GetPing() != nil { + break + } + } + + var rss int64 + if err := f.pool.QueryRow(f.ctx, ` + SELECT rss_bytes FROM ts.agent_health WHERE agent_id = $1 + ORDER BY ts DESC LIMIT 1`, f.agentID).Scan(&rss); err != nil { + t.Fatalf("самометрики не записались: %v", err) + } + if rss != 12<<20 { + t.Fatalf("rss_bytes = %d", rss) + } + + // Втрачені семпли мають бути видимі оператору у зведенні зонда, + // а не тільки в графіку. + var health []byte + if err := f.pool.QueryRow(f.ctx, + `SELECT health::text FROM core.agents WHERE id = $1`, f.agentID).Scan(&health); err != nil { + t.Fatalf("зведення: %v", err) + } + var h map[string]any + if err := json.Unmarshal(health, &h); err != nil { + t.Fatalf("health не JSON: %v", err) + } + if h["dropped_samples"] == nil || fmt.Sprint(h["dropped_samples"]) != "7" { + t.Fatalf("dropped_samples не потрапило у зведення: %v", h) + } +} + +// --------------------------------------------------------------------- +// 6. Хеш плану дозволяє не переливати задачі після реконекту +// --------------------------------------------------------------------- + +func TestPlanHashSkipsResend(t *testing.T) { + f := setup(t) + + ctx, cancel := context.WithTimeout(f.authCtx(), 20*time.Second) + defer cancel() + + // Перше підключення: дізнаємось хеш. + plan, err := f.store.BuildPlan(f.ctx, &store.Agent{ID: f.agentID, TenantID: f.tenantID}) + if err != nil { + t.Fatalf("BuildPlan: %v", err) + } + if len(plan.PlanHash) == 0 { + t.Fatal("план без хеша") + } + + stream, err := f.client.Control(ctx) + if err != nil { + t.Fatalf("Control: %v", err) + } + if err := stream.Send(&npv1.ControlUp{ + Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{ + AgentId: f.agentID, TaskPlanHash: plan.PlanHash, + }}, + }); err != nil { + t.Fatalf("Hello: %v", err) + } + + down, err := stream.Recv() + if err != nil { + t.Fatalf("Welcome: %v", err) + } + if down.GetWelcome().GetTaskPlanFollows() { + t.Fatal("сервер збирається слати план, хоча хеш збігся") + } + + // Хеш детермінований: та сама база — той самий хеш. + again, err := f.store.BuildPlan(f.ctx, &store.Agent{ID: f.agentID, TenantID: f.tenantID}) + if err != nil { + t.Fatalf("BuildPlan: %v", err) + } + if hex.EncodeToString(again.PlanHash) != hex.EncodeToString(plan.PlanHash) { + t.Fatal("хеш плану не детермінований") + } +} diff --git a/server/internal/grpcapi/service.go b/server/internal/grpcapi/service.go new file mode 100644 index 0000000..1051a3d --- /dev/null +++ b/server/internal/grpcapi/service.go @@ -0,0 +1,363 @@ +// Package grpcapi — серверна реалізація AgentService. +// +// Сервер ніколи не ініціює з'єднання до зонда: у мережі клієнта немає +// ані відкритих портів, ані прокидання NAT. Тому все, що виглядає як +// «команда з сервера», надсилається у зустрічному напрямку стріму +// Control, який відкрив сам агент. +package grpcapi + +import ( + "context" + "errors" + "io" + "log/slog" + "strings" + "sync" + "time" + + "github.com/netpulse/netpulse/server/internal/crypto" + "github.com/netpulse/netpulse/server/internal/store" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/metadata" + "google.golang.org/grpc/status" + "google.golang.org/protobuf/types/known/durationpb" + "google.golang.org/protobuf/types/known/timestamppb" +) + +type ctxKey string + +const agentCtxKey ctxKey = "netpulse.agent" + +// Service реалізує npv1.AgentServiceServer. +type Service struct { + npv1.UnimplementedAgentServiceServer + + store *store.Store + ring *crypto.Keyring + log *slog.Logger + + // Живі сесії за agent_id. Потрібні, щоб штовхнути зонду + // TaskDelta або ConfigJob, коли щось змінилось в UI. + mu sync.RWMutex + sessions map[string]*agentSession +} + +type agentSession struct { + agent *store.Agent + out chan *npv1.ControlDown + series *store.SeriesTable + opened time.Time +} + +func New(st *store.Store, ring *crypto.Keyring, log *slog.Logger) *Service { + if log == nil { + log = slog.Default() + } + return &Service{ + store: st, + ring: ring, + log: log, + sessions: make(map[string]*agentSession), + } +} + +// --------------------------------------------------------------------- +// Автентифікація +// --------------------------------------------------------------------- + +// StreamInterceptor і UnaryInterceptor автентифікують зонда за токеном +// у метаданих. У проді поверх цього ще mTLS: токен визначає, ЯКИЙ це +// зонд, сертифікат — що він узагалі має право говорити з сервером. +func (s *Service) StreamInterceptor(srv any, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error { + agent, err := s.authenticate(ss.Context()) + if err != nil { + return err + } + return handler(srv, &wrappedStream{ServerStream: ss, + ctx: context.WithValue(ss.Context(), agentCtxKey, agent)}) +} + +func (s *Service) UnaryInterceptor(ctx context.Context, req any, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (any, error) { + agent, err := s.authenticate(ctx) + if err != nil { + return nil, err + } + return handler(context.WithValue(ctx, agentCtxKey, agent), req) +} + +type wrappedStream struct { + grpc.ServerStream + ctx context.Context +} + +func (w *wrappedStream) Context() context.Context { return w.ctx } + +func (s *Service) authenticate(ctx context.Context) (*store.Agent, error) { + md, ok := metadata.FromIncomingContext(ctx) + if !ok { + return nil, status.Error(codes.Unauthenticated, "немає метаданих") + } + + var token string + for _, v := range md.Get("authorization") { + if after, found := strings.CutPrefix(v, "Bearer "); found { + token = after + break + } + } + if token == "" { + return nil, status.Error(codes.Unauthenticated, "немає токена зонда") + } + + agent, err := s.store.AuthenticateAgent(ctx, token) + if errors.Is(err, store.ErrAgentNotFound) { + return nil, status.Error(codes.Unauthenticated, "зонда не знайдено") + } + if err != nil { + return nil, status.Error(codes.Internal, "помилка автентифікації") + } + return agent, nil +} + +func agentFrom(ctx context.Context) (*store.Agent, error) { + a, ok := ctx.Value(agentCtxKey).(*store.Agent) + if !ok || a == nil { + return nil, status.Error(codes.Unauthenticated, "сесія без зонда") + } + return a, nil +} + +// --------------------------------------------------------------------- +// Control +// --------------------------------------------------------------------- + +func (s *Service) Control(stream npv1.AgentService_ControlServer) error { + ctx := stream.Context() + agent, err := agentFrom(ctx) + if err != nil { + return err + } + + first, err := stream.Recv() + if err != nil { + return err + } + hello := first.GetHello() + if hello == nil { + return status.Error(codes.InvalidArgument, "перше повідомлення має бути Hello") + } + if hello.GetAgentId() != "" && hello.GetAgentId() != agent.ID { + // Токен належить одному зонду, а він назвався іншим. + return status.Error(codes.PermissionDenied, "agent_id не збігається з токеном") + } + + if err := s.store.MarkAgentOnline(ctx, agent, hello); err != nil { + s.log.Warn("не вдалося позначити зонд онлайн", "agent", agent.ID, "err", err) + } + + sess := &agentSession{ + agent: agent, + out: make(chan *npv1.ControlDown, 64), + series: store.NewSeriesTable(), + opened: time.Now(), + } + s.register(sess) + defer func() { + s.unregister(agent.ID) + if err := s.store.MarkAgentOffline(context.WithoutCancel(ctx), agent); err != nil { + s.log.Warn("не вдалося позначити зонд офлайн", "agent", agent.ID, "err", err) + } + }() + + plan, err := s.store.BuildPlan(ctx, agent) + if err != nil { + return status.Errorf(codes.Internal, "побудова плану: %v", err) + } + + // Якщо хеш збігся — план у агента вже правильний, і 50 000 задач + // після кожного обриву зв'язку переливати не треба. + planUnchanged := len(hello.GetTaskPlanHash()) > 0 && + equalBytes(hello.GetTaskPlanHash(), plan.GetPlanHash()) + + welcome := &npv1.Welcome{ + SessionId: agent.ID + "@" + time.Now().UTC().Format(time.RFC3339Nano), + ServerTime: timestamppb.Now(), + HeartbeatInterval: durationpb.New(agent.Limits.HeartbeatEvery), + TelemetryMaxBatchSize: uint32(agent.Limits.BatchSize), + TelemetryMaxBatchInterval: durationpb.New(agent.Limits.BatchInterval), + TelemetryMaxInFlight: uint32(agent.Limits.MaxInFlight), + MaxConcurrentChecks: uint32(agent.Limits.MaxConcurrency), + IcmpRatePps: uint32(agent.Limits.IcmpRatePPS), + TaskPlanFollows: !planUnchanged, + } + if err := stream.Send(&npv1.ControlDown{ + Seq: 1, Payload: &npv1.ControlDown_Welcome{Welcome: welcome}, + }); err != nil { + return err + } + + // Єдиний писар: gRPC не допускає паралельних Send на одному стрімі, + // а слати треба і план, і креденшели, і пуші з UI. + writerDone := make(chan struct{}) + go func() { + defer close(writerDone) + var seq uint64 = 1 + for { + select { + case <-ctx.Done(): + return + case msg, ok := <-sess.out: + if !ok { + return + } + seq++ + msg.Seq = seq + if err := stream.Send(msg); err != nil { + return + } + } + } + }() + + s.push(sess, &npv1.ControlDown{ + Payload: &npv1.ControlDown_ModuleControl{ + ModuleControl: store.ModulesForPlan(plan, agent.Modules)}, + }) + + bundle, err := s.store.BuildCredentialBundle(ctx, agent, s.ring) + if err != nil { + s.log.Warn("не вдалося зібрати креденшели", "agent", agent.ID, "err", err) + } else if len(bundle.ByDevice) > 0 { + s.push(sess, &npv1.ControlDown{ + Payload: &npv1.ControlDown_Credentials{Credentials: bundle}, + }) + } + + if !planUnchanged { + s.push(sess, &npv1.ControlDown{ + Payload: &npv1.ControlDown_TaskPlan{TaskPlan: plan}, + }) + } + + s.log.Info("зонд підключився", + "agent", agent.ID, "tenant", agent.TenantID, + "tasks", len(plan.GetTasks()), "plan_unchanged", planUnchanged) + + err = s.readControl(ctx, stream, agent, sess) + close(sess.out) + <-writerDone + return err +} + +func (s *Service) readControl(ctx context.Context, stream npv1.AgentService_ControlServer, agent *store.Agent, sess *agentSession) error { + 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: + if err := s.store.RecordHeartbeat(ctx, agent, p.Heartbeat); err != nil { + s.log.Warn("heartbeat не записався", "agent", agent.ID, "err", err) + } + // Ping у відповідь дає обом сторонам свіжий clock skew. + s.push(sess, &npv1.ControlDown{ + Payload: &npv1.ControlDown_Ping{Ping: &npv1.Ping{ + PingId: uint64(time.Now().UnixNano()), ServerTime: timestamppb.Now(), + }}, + }) + + case *npv1.ControlUp_TaskStatus: + if err := s.store.RecordTaskStatus(ctx, agent, p.TaskStatus); err != nil { + s.log.Warn("статус задачі не записався", "agent", agent.ID, "err", err) + } + + case *npv1.ControlUp_CredentialRequest: + bundle, err := s.store.BuildCredentialBundle(ctx, agent, s.ring) + if err != nil { + s.log.Warn("повторна видача креденшелів", "agent", agent.ID, "err", err) + continue + } + s.push(sess, &npv1.ControlDown{ + Payload: &npv1.ControlDown_Credentials{Credentials: bundle}, + }) + + case *npv1.ControlUp_Event: + s.log.Info("подія зонда", + "agent", agent.ID, "kind", p.Event.GetKind().String(), + "message", p.Event.GetMessage()) + + case *npv1.ControlUp_Pong: + // Нічого не робимо: сам факт відповіді підтверджує, що + // канал живий, а це вже зафіксував транспорт. + } + } +} + +// push кладе повідомлення в чергу писаря. Черга переповнилась — +// зонд не встигає читати; краще втратити пуш, ніж заблокувати +// обробку всіх інших зондів на цьому воркері. +func (s *Service) push(sess *agentSession, msg *npv1.ControlDown) { + select { + case sess.out <- msg: + default: + s.log.Warn("черга контрольного каналу переповнена", "agent", sess.agent.ID) + } +} + +func (s *Service) register(sess *agentSession) { + s.mu.Lock() + defer s.mu.Unlock() + if old, ok := s.sessions[sess.agent.ID]; ok { + // Той самий зонд відкрив другу сесію: стара мертва, просто + // ще не помітила. Закривати її чергу тут не можна — це + // зробить її власний defer. + s.log.Warn("друга сесія того самого зонда", "agent", sess.agent.ID, + "попередня_відкрита", old.opened) + } + s.sessions[sess.agent.ID] = sess +} + +func (s *Service) unregister(agentID string) { + s.mu.Lock() + defer s.mu.Unlock() + delete(s.sessions, agentID) +} + +// SessionCount — скільки зондів онлайн просто зараз. +func (s *Service) SessionCount() int { + s.mu.RLock() + defer s.mu.RUnlock() + return len(s.sessions) +} + +// PushToAgent надсилає команду живій сесії зонда. Використовується +// з боку API, коли в UI щось змінили. +func (s *Service) PushToAgent(agentID string, msg *npv1.ControlDown) bool { + s.mu.RLock() + sess, ok := s.sessions[agentID] + s.mu.RUnlock() + if !ok { + return false + } + s.push(sess, msg) + return true +} + +func equalBytes(a, b []byte) bool { + if len(a) != len(b) { + return false + } + for i := range a { + if a[i] != b[i] { + return false + } + } + return true +} diff --git a/server/internal/grpcapi/streams.go b/server/internal/grpcapi/streams.go new file mode 100644 index 0000000..ceb152b --- /dev/null +++ b/server/internal/grpcapi/streams.go @@ -0,0 +1,268 @@ +package grpcapi + +import ( + "context" + "errors" + "io" + + "github.com/netpulse/netpulse/server/internal/store" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +// --------------------------------------------------------------------- +// Телеметрія +// --------------------------------------------------------------------- + +func (s *Service) StreamTelemetry(stream npv1.AgentService_StreamTelemetryServer) error { + ctx := stream.Context() + agent, err := agentFrom(ctx) + if err != nil { + return err + } + + // Таблиця series_ref живе рівно стільки, скільки цей стрім. + // Агент знає про це й перереєстровує серії на початку кожної сесії. + table := store.NewSeriesTable() + + var acked uint64 + + for { + batch, err := stream.Recv() + if errors.Is(err, io.EOF) { + return nil + } + if err != nil { + return err + } + + st, writeErr := s.store.WriteBatch(ctx, agent, batch, table) + + var unknown *store.ErrUnknownSeriesRef + if errors.As(writeErr, &unknown) { + // Ми загубили відповідність — чесно просимо перереєстрацію + // замість тихо викинути дані. + s.log.Warn("невідомий series_ref", "agent", agent.ID, "ref", unknown.Ref) + if err := stream.Send(&npv1.TelemetryAck{ + AckedThroughBatchId: acked, + ResetSeriesTable: true, + MaxInFlight: uint32(agent.Limits.MaxInFlight), + Error: &npv1.Error{ + Code: "unknown_series_ref", + Message: unknown.Error(), + Retryable: true, + }, + }); err != nil { + return err + } + continue + } + + if writeErr != nil { + s.log.Error("батч не записався", "agent", agent.ID, + "batch", batch.GetBatchId(), "err", writeErr) + // Батч не підтверджуємо: агент перешле його після + // реконекту, а ON CONFLICT DO NOTHING зробить це безпечним. + if err := stream.Send(&npv1.TelemetryAck{ + AckedThroughBatchId: acked, + NackBatchIds: []uint64{batch.GetBatchId()}, + MaxInFlight: uint32(agent.Limits.MaxInFlight), + Error: &npv1.Error{ + Code: "write_failed", Message: writeErr.Error(), Retryable: true, + }, + }); err != nil { + return err + } + continue + } + + if batch.GetBatchId() > acked { + acked = batch.GetBatchId() + } + + s.log.Debug("батч записано", "agent", agent.ID, "batch", batch.GetBatchId(), + "samples", st.Samples, "icmp", st.Icmp, "interfaces", st.Interfaces) + + if err := stream.Send(&npv1.TelemetryAck{ + AckedThroughBatchId: acked, + MaxInFlight: uint32(agent.Limits.MaxInFlight), + }); err != nil { + return err + } + } +} + +// --------------------------------------------------------------------- +// Журнали +// --------------------------------------------------------------------- + +func (s *Service) StreamLogs(stream npv1.AgentService_StreamLogsServer) error { + ctx := stream.Context() + agent, err := agentFrom(ctx) + if err != nil { + return err + } + + var acked uint64 + for { + batch, err := stream.Recv() + if errors.Is(err, io.EOF) { + return nil + } + if err != nil { + return err + } + + if _, err := s.store.WriteLogs(ctx, agent, batch); err != nil { + s.log.Error("журнали не записались", "agent", agent.ID, "err", err) + if err := stream.Send(&npv1.LogAck{ + AckedThroughBatchId: acked, + Error: &npv1.Error{Code: "write_failed", Message: err.Error(), Retryable: true}, + }); err != nil { + return err + } + continue + } + + if batch.GetDropped() > 0 { + s.log.Warn("зонд відкинув події через ліміт", + "agent", agent.ID, "dropped", batch.GetDropped()) + } + if batch.GetBatchId() > acked { + acked = batch.GetBatchId() + } + + if err := stream.Send(&npv1.LogAck{AckedThroughBatchId: acked}); err != nil { + return err + } + } +} + +// --------------------------------------------------------------------- +// Автовиявлення +// --------------------------------------------------------------------- + +func (s *Service) ReportDiscovery(ctx context.Context, rep *npv1.DiscoveryReport) (*npv1.DiscoveryAck, error) { + agent, err := agentFrom(ctx) + if err != nil { + return nil, err + } + + st, err := s.store.ApplyDiscovery(ctx, agent, rep) + if err != nil { + s.log.Error("звіт автовиявлення не застосувався", "agent", agent.ID, "err", err) + return &npv1.DiscoveryAck{ + Accepted: false, + Error: &npv1.Error{Code: "apply_failed", Message: err.Error(), Retryable: true}, + }, nil + } + + s.log.Info("автовиявлення застосовано", + "agent", agent.ID, + "сусідів", st.NeighborsSeen, "зіставлено", st.NeighborsResolved, + "лінків_створено", st.LinksCreated, "лінків_оновлено", st.LinksUpdated) + + return &npv1.DiscoveryAck{ + Accepted: true, + NeighborsResolved: uint32(st.NeighborsResolved), + LinksCreated: uint32(st.LinksCreated), + }, nil +} + +// --------------------------------------------------------------------- +// Конфігурації +// --------------------------------------------------------------------- + +// UploadConfig збирає конфіг із чанків: header → chunk* → trailer. +func (s *Service) UploadConfig(stream npv1.AgentService_UploadConfigServer) error { + ctx := stream.Context() + agent, err := agentFrom(ctx) + if err != nil { + return err + } + + var ( + header *npv1.ConfigHeader + body []byte + wantSeq uint32 + ) + + for { + msg, err := stream.Recv() + if errors.Is(err, io.EOF) { + return status.Error(codes.InvalidArgument, "стрім завершився без trailer") + } + if err != nil { + return err + } + + switch p := msg.Part.(type) { + case *npv1.ConfigUpload_Header: + header = p.Header + + case *npv1.ConfigUpload_Chunk: + if header == nil { + return status.Error(codes.InvalidArgument, "chunk раніше за header") + } + // Порядок важливий: конфіг склеюється байт-у-байт, і + // переставлені чанки дали б тихо зіпсований бекап. + if p.Chunk.GetSequence() != wantSeq { + return status.Errorf(codes.InvalidArgument, + "чанки не по порядку: чекали %d, отримали %d", wantSeq, p.Chunk.GetSequence()) + } + wantSeq++ + body = append(body, p.Chunk.GetData()...) + + case *npv1.ConfigUpload_Trailer: + if header == nil { + return status.Error(codes.InvalidArgument, "trailer без header") + } + tr := p.Trailer + + if !tr.GetSuccess() { + s.log.Warn("зонд не зміг зібрати конфіг", + "agent", agent.ID, "job", header.GetJobId(), + "err", tr.GetError().GetMessage()) + return stream.SendAndClose(&npv1.ConfigReceipt{ + JobId: header.GetJobId(), Accepted: false, Error: tr.GetError(), + }) + } + + outcome, err := s.store.StoreConfig(ctx, agent, store.ConfigSubmission{ + JobID: header.GetJobId(), + DeviceID: header.GetDeviceId(), + ConfigType: header.GetConfigType(), + Body: body, + ClaimedHash: tr.GetContentSha256(), + }, s.ring) + if err != nil { + return status.Errorf(codes.Internal, "збереження конфігу: %v", err) + } + if errors.Is(outcome.Err, store.ErrChecksumMismatch) { + s.log.Warn("конфіг із розбіжною контрольною сумою відхилено", + "agent", agent.ID, "device", header.GetDeviceId()) + return stream.SendAndClose(&npv1.ConfigReceipt{ + JobId: header.GetJobId(), + Accepted: false, + Error: &npv1.Error{ + Code: "checksum_mismatch", Retryable: true, + Message: "тіло не відповідає заявленому sha256", + }, + }) + } + + s.log.Info("конфіг прийнято", + "agent", agent.ID, "device", header.GetDeviceId(), + "розмір", len(body), "без_змін", outcome.Unchanged) + + return stream.SendAndClose(&npv1.ConfigReceipt{ + JobId: header.GetJobId(), + Accepted: outcome.Accepted, + Unchanged: outcome.Unchanged, + ConfigId: outcome.ConfigID, + CommitSha: outcome.CommitSHA, + }) + } + } +} diff --git a/server/internal/store/agents.go b/server/internal/store/agents.go new file mode 100644 index 0000000..d63cde6 --- /dev/null +++ b/server/internal/store/agents.go @@ -0,0 +1,227 @@ +package store + +import ( + "context" + "crypto/sha256" + "crypto/subtle" + "errors" + "fmt" + "time" + + "github.com/jackc/pgx/v5" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" +) + +var ErrAgentNotFound = errors.New("зонд не знайдено або токен відкликано") + +// Agent — ідентичність зонда після автентифікації. +type Agent struct { + ID string + TenantID string + Name string + SiteID string + Modules []string + Limits Limits +} + +// Limits — ліміти опитування, які сервер диктує зонду в Welcome. +type Limits struct { + MaxConcurrency int + IcmpRatePPS int + BatchSize int + BatchInterval time.Duration + MaxInFlight int + HeartbeatEvery time.Duration +} + +func defaultLimits() Limits { + return Limits{ + MaxConcurrency: 256, + IcmpRatePPS: 500, + BatchSize: 500, + BatchInterval: 5 * time.Second, + MaxInFlight: 4, + HeartbeatEvery: 30 * time.Second, + } +} + +// AuthenticateAgent знаходить зонда за токеном. +// +// У БД лежить лише sha256 — токен у відкритому вигляді не зберігається +// ніде. Порівняння хешів робиться в SQL за унікальним індексом; на +// знайденому рядку додатково звіряємо constant-time, щоб форма запиту +// не залежала від того, як індекс порівнював байти. +func (s *Store) AuthenticateAgent(ctx context.Context, token string) (*Agent, error) { + if token == "" { + return nil, ErrAgentNotFound + } + sum := sha256.Sum256([]byte(token)) + + var ( + a Agent + siteID *string + hash []byte + status string + limits map[string]any + ) + + err := s.pool.QueryRow(ctx, ` + SELECT id::text, tenant_id::text, name, site_id::text, token_hash, + status::text, enabled_modules, limits + FROM core.agents + WHERE token_hash = $1 + `, sum[:]).Scan(&a.ID, &a.TenantID, &a.Name, &siteID, &hash, &status, &a.Modules, &limits) + + if errors.Is(err, pgx.ErrNoRows) { + return nil, ErrAgentNotFound + } + if err != nil { + return nil, err + } + if subtle.ConstantTimeCompare(hash, sum[:]) != 1 { + return nil, ErrAgentNotFound + } + if status == "disabled" { + return nil, fmt.Errorf("%w: зонд вимкнено", ErrAgentNotFound) + } + if siteID != nil { + a.SiteID = *siteID + } + + a.Limits = defaultLimits() + applyLimitOverrides(&a.Limits, limits) + return &a, nil +} + +func applyLimitOverrides(l *Limits, m map[string]any) { + num := func(key string) (int, bool) { + v, ok := m[key] + if !ok { + return 0, false + } + switch n := v.(type) { + case float64: + return int(n), true + case int64: + return int(n), true + } + return 0, false + } + if v, ok := num("max_concurrency"); ok && v > 0 { + l.MaxConcurrency = v + } + if v, ok := num("icmp_rate_pps"); ok && v > 0 { + l.IcmpRatePPS = v + } + if v, ok := num("batch_size"); ok && v > 0 { + l.BatchSize = v + } + if v, ok := num("max_in_flight"); ok && v > 0 { + l.MaxInFlight = v + } +} + +// MarkAgentOnline фіксує підключення зонда та його версію. +func (s *Store) MarkAgentOnline(ctx context.Context, a *Agent, hello *npv1.Hello) error { + build := hello.GetBuild() + return s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + _, err := tx.Exec(ctx, ` + UPDATE core.agents + SET status = 'online', + last_heartbeat_at = now(), + version = COALESCE(NULLIF($2,''), version), + os = COALESCE(NULLIF($3,''), os), + arch = COALESCE(NULLIF($4,''), arch), + hostname= COALESCE(NULLIF($5,''), hostname), + updated_at = now() + WHERE id = $1 AND tenant_id = $6 + `, a.ID, build.GetVersion(), build.GetOs(), build.GetArch(), + hello.GetHostname(), a.TenantID) + return err + }) +} + +// MarkAgentOffline викликається при розриві сесії. +func (s *Store) MarkAgentOffline(ctx context.Context, a *Agent) error { + return s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + _, err := tx.Exec(ctx, ` + UPDATE core.agents SET status = 'offline', updated_at = now() + WHERE id = $1 AND tenant_id = $2 + `, a.ID, a.TenantID) + return err + }) +} + +// RecordHeartbeat пише самометрики зонда в ts.agent_health і оновлює +// зведення в core.agents.health для швидкого показу в UI. +func (s *Store) RecordHeartbeat(ctx context.Context, a *Agent, hb *npv1.Heartbeat) error { + h := hb.GetHealth() + ts := hb.GetTs().AsTime() + if ts.IsZero() { + ts = time.Now() + } + + batch := &pgx.Batch{} + batch.Queue(` + INSERT INTO ts.agent_health + (ts, agent_id, tenant_id, cpu_pct, rss_bytes, goroutines, + queue_depth, checks_per_sec, errors_per_min) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9) + ON CONFLICT (ts, agent_id) DO NOTHING + `, ts, a.ID, a.TenantID, h.GetCpuPct(), int64(h.GetRssBytes()), + int32(h.GetGoroutines()), int32(h.GetQueueDepth()), + h.GetChecksPerSec(), h.GetErrorsPerMin()) + + // dropped_samples у зведенні — це видима діра в даних; вона має + // бути помітна оператору, а не тільки в графіку самометрик. + batch.Queue(` + UPDATE core.agents + SET last_heartbeat_at = now(), + status = 'online', + health = jsonb_build_object( + 'rss_bytes', $3::bigint, + 'queue_depth', $4::int, + 'dropped_samples', $5::bigint, + 'tasks_running', $6::int, + 'clock_skew_ms', $7::bigint + ) + WHERE id = $1 AND tenant_id = $2 + `, a.ID, a.TenantID, int64(h.GetRssBytes()), int32(h.GetQueueDepth()), + int64(h.GetDroppedSamples()), int32(hb.GetTasksRunning()), + h.GetClockSkew().AsDuration().Milliseconds()) + + res := s.pool.SendBatch(ctx, batch) + defer res.Close() + for i := 0; i < batch.Len(); i++ { + if _, err := res.Exec(); err != nil { + return fmt.Errorf("heartbeat[%d]: %w", i, err) + } + } + return nil +} + +// RecordTaskStatus оновлює core.checks за доповіддю агента. +func (s *Store) RecordTaskStatus(ctx context.Context, a *Agent, u *npv1.TaskStatusUpdate) error { + var errText any + if u.GetError() != nil { + errText = u.GetError().GetCode() + ": " + u.GetError().GetMessage() + } + + return s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + _, err := tx.Exec(ctx, ` + UPDATE core.checks + SET last_run_at = COALESCE($3, now()), + last_error = $4, + updated_at = now() + WHERE id = $1 AND tenant_id = $2 + `, u.GetCheckId(), a.TenantID, tsOrNil(u), errText) + return err + }) +} + +func tsOrNil(u *npv1.TaskStatusUpdate) any { + if u.GetTs() == nil { + return nil + } + return u.GetTs().AsTime() +} diff --git a/server/internal/store/credentials.go b/server/internal/store/credentials.go new file mode 100644 index 0000000..2333145 --- /dev/null +++ b/server/internal/store/credentials.go @@ -0,0 +1,205 @@ +package store + +import ( + "context" + "encoding/json" + "time" + + "github.com/jackc/pgx/v5" + "github.com/netpulse/netpulse/server/internal/crypto" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" + "google.golang.org/protobuf/types/known/timestamppb" +) + +// CredentialTTL — скільки живе комплект на зонді. +// +// Короткий TTL — це не паранойя, а керованість: після відкликання +// доступу зонд перестане ним користуватись за годину без жодних +// додаткових дій, навіть якщо зв'язок із ним втрачено. +const CredentialTTL = time.Hour + +// snmpV3Payload — домовленість про формат відкритого тексту секрету +// kind='snmp_v3': два паролі в одному секреті, бо вони завжди +// використовуються разом і ротуються разом. +type snmpV3Payload struct { + AuthPassword string `json:"auth_password"` + PrivPassword string `json:"priv_password"` +} + +type credOptions struct { + SecLevel string `json:"sec_level"` + AuthProto string `json:"auth_proto"` + PrivProto string `json:"priv_proto"` + Context string `json:"context"` + SecurityName string `json:"security_name"` +} + +// BuildCredentialBundle збирає розшифровані креденшели для всіх +// пристроїв цього зонда. +// +// Розшифровка відбувається саме тут, а не на агенті: DEK не покидає +// сервер. Далі вони їдуть виключно всередині mTLS-каналу й живуть у +// пам'яті агента до expires_at. +func (s *Store) BuildCredentialBundle(ctx context.Context, a *Agent, ring *crypto.Keyring) (*npv1.CredentialBundle, error) { + bundle := &npv1.CredentialBundle{ + ByDevice: make(map[string]*npv1.CredentialList), + ExpiresAt: timestamppb.New(time.Now().Add(CredentialTTL)), + } + + err := s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + rows, err := tx.Query(ctx, ` + SELECT dc.device_id::text, + c.id::text, + c.proto::text, + COALESCE(c.username, ''), + COALESCE(c.port, 0), + COALESCE(c.options::text, '{}'), + COALESCE(s.kind::text, ''), + s.key_id, s.nonce, s.ciphertext, s.auth_tag, COALESCE(s.aad, '') + FROM inv.device_credentials dc + JOIN inv.credentials c ON c.id = dc.credential_id + JOIN inv.devices d ON d.id = dc.device_id + LEFT JOIN core.secrets s ON s.id = c.secret_id + WHERE c.tenant_id = $1 + AND d.agent_id = $2 + AND d.enabled + AND d.deleted_at IS NULL + ORDER BY dc.device_id, dc.priority + `, a.TenantID, a.ID) + if err != nil { + return err + } + defer rows.Close() + + for rows.Next() { + var ( + deviceID, credID, proto, username, optionsJSON string + port int32 + kind, keyID, aad string + nonce, ciphertext, authTag []byte + ) + var ( + keyIDPtr *string + aadPtr *string + ) + if err := rows.Scan(&deviceID, &credID, &proto, &username, &port, + &optionsJSON, &kind, &keyIDPtr, &nonce, &ciphertext, &authTag, &aadPtr); err != nil { + return err + } + if keyIDPtr != nil { + keyID = *keyIDPtr + } + if aadPtr != nil { + aad = *aadPtr + } + + cred := &npv1.Credential{ + CredentialId: credID, + Transport: transportFromProto(proto), + Username: username, + Port: uint32(port), + ExpiresAt: bundle.ExpiresAt, + } + + var plaintext []byte + if len(ciphertext) > 0 { + pt, decErr := ring.Decrypt(&crypto.Secret{ + KeyID: keyID, Nonce: nonce, Ciphertext: ciphertext, AuthTag: authTag, + }, aad) + if decErr != nil { + // Один нечитабельний секрет не має рушити весь + // комплект: решта пристроїв мусить опитуватись. + continue + } + plaintext = pt + } + + var opts credOptions + _ = json.Unmarshal([]byte(optionsJSON), &opts) + + if err := fillSecret(cred, kind, proto, plaintext, opts); err != nil { + continue + } + + list := bundle.ByDevice[deviceID] + if list == nil { + list = &npv1.CredentialList{} + bundle.ByDevice[deviceID] = list + } + list.Credentials = append(list.Credentials, cred) + } + return rows.Err() + }) + + return bundle, err +} + +func fillSecret(cred *npv1.Credential, kind, proto string, plaintext []byte, opts credOptions) error { + switch { + case proto == "snmp_v2c": + cred.Secret = &npv1.Credential_Community{Community: string(plaintext)} + + case proto == "snmp_v3": + var p snmpV3Payload + if len(plaintext) > 0 { + if err := json.Unmarshal(plaintext, &p); err != nil { + return err + } + } + cred.SnmpV3 = &npv1.SnmpV3Options{ + Level: secLevel(opts.SecLevel), + AuthProtocol: opts.AuthProto, + AuthPassword: p.AuthPassword, + PrivProtocol: opts.PrivProto, + PrivPassword: p.PrivPassword, + ContextName: opts.Context, + SecurityName: opts.SecurityName, + } + + case kind == "ssh_key": + cred.Secret = &npv1.Credential_PrivateKey{PrivateKey: plaintext} + + case proto == "http" || proto == "https" || proto == "api": + cred.Secret = &npv1.Credential_Token{Token: string(plaintext)} + + default: + cred.Secret = &npv1.Credential_Password{Password: string(plaintext)} + } + return nil +} + +func transportFromProto(p string) npv1.Transport { + switch p { + case "ssh": + return npv1.Transport_TRANSPORT_SSH + case "telnet": + return npv1.Transport_TRANSPORT_TELNET + case "snmp_v2c": + return npv1.Transport_TRANSPORT_SNMP_V2C + case "snmp_v3": + return npv1.Transport_TRANSPORT_SNMP_V3 + case "http": + return npv1.Transport_TRANSPORT_HTTP + case "https": + return npv1.Transport_TRANSPORT_HTTPS + case "api": + return npv1.Transport_TRANSPORT_API + case "modbus": + return npv1.Transport_TRANSPORT_MODBUS + default: + return npv1.Transport_TRANSPORT_UNSPECIFIED + } +} + +func secLevel(s string) npv1.SnmpV3Options_SecurityLevel { + switch s { + case "authPriv": + return npv1.SnmpV3Options_SECURITY_LEVEL_AUTH_PRIV + case "authNoPriv": + return npv1.SnmpV3Options_SECURITY_LEVEL_AUTH_NO_PRIV + case "noAuthNoPriv": + return npv1.SnmpV3Options_SECURITY_LEVEL_NO_AUTH_NO_PRIV + default: + return npv1.SnmpV3Options_SECURITY_LEVEL_UNSPECIFIED + } +} diff --git a/server/internal/store/discovery.go b/server/internal/store/discovery.go new file mode 100644 index 0000000..de720af --- /dev/null +++ b/server/internal/store/discovery.go @@ -0,0 +1,415 @@ +package store + +import ( + "context" + "encoding/json" + "errors" + "strconv" + "strings" + + "github.com/jackc/pgx/v5" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" +) + +// Рівні впевненості зіставлення сусіда з інвентарем. +// +// Порядок не довільний: chassis-id унікальний за стандартом LLDP, +// MAC — майже, IP керування може повторюватись у різних VRF, а +// sysName взагалі вводить людина й на двох комутаторах цілком може +// бути "switch". Тому лінк, зведений за іменем, отримує низьку +// впевненість і не має мовчки заміщати те, що підтвердила людина. +const ( + confChassis = 95 + confMAC = 90 + confMgmtIP = 80 + confSysName = 60 +) + +// DiscoveryStats — підсумок обробки звіту, їде в DiscoveryAck. +type DiscoveryStats struct { + NeighborsSeen int + NeighborsResolved int + LinksCreated int + LinksUpdated int + InterfacesUpdated int +} + +// ApplyDiscovery приймає звіт агента: зберігає сирих сусідів, +// зіставляє їх з інвентарем і зводить у topo.links. +// +// Агент не вирішує, хто з ким з'єднаний, — він доповідає лише +// «на порту X бачу chassis Y, port Z». Уся інтерпретація тут, тому +// правила зіставлення можна міняти без оновлення зондів у полі. +func (s *Store) ApplyDiscovery(ctx context.Context, a *Agent, rep *npv1.DiscoveryReport) (DiscoveryStats, error) { + var st DiscoveryStats + + if err := s.upsertInterfaces(ctx, a, rep.GetInterfaces(), &st); err != nil { + return st, err + } + + for _, n := range rep.GetNeighbors() { + st.NeighborsSeen++ + if err := s.applyNeighbor(ctx, a, n, &st); err != nil { + return st, err + } + } + + return st, nil +} + +func (s *Store) upsertInterfaces(ctx context.Context, a *Agent, ifs []*npv1.InterfaceRecord, st *DiscoveryStats) error { + if len(ifs) == 0 { + return nil + } + + return s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + for _, r := range ifs { + _, err := tx.Exec(ctx, ` + INSERT INTO inv.interfaces + (tenant_id, device_id, if_index, name, alias, mac, mtu, type, + speed_bps, duplex, admin_status, oper_status) + VALUES ($1,$2,$3,$4,$5,$6::macaddr,$7,$8,$9,$10, + $11::inv.if_admin_status, $12::inv.if_oper_status) + ON CONFLICT (device_id, if_index) WHERE if_index IS NOT NULL + DO UPDATE SET + name = EXCLUDED.name, + alias = EXCLUDED.alias, + mac = COALESCE(EXCLUDED.mac, inv.interfaces.mac), + speed_bps = COALESCE(NULLIF(EXCLUDED.speed_bps, 0), inv.interfaces.speed_bps), + admin_status = EXCLUDED.admin_status, + oper_status = EXCLUDED.oper_status, + updated_at = now() + `, a.TenantID, r.GetDeviceId(), r.GetIfIndex(), r.GetName(), + nullString(r.GetAlias()), nullString(r.GetMac()), + nullInt(int(r.GetMtu())), nullString(r.GetType()), + nullInt64(int64(r.GetSpeedBps())), nullString(r.GetDuplex()), + statusOr(r.GetAdminStatus(), "unknown"), + statusOr(r.GetOperStatus(), "unknown")) + if err != nil { + return err + } + st.InterfacesUpdated++ + } + return nil + }) +} + +func (s *Store) applyNeighbor(ctx context.Context, a *Agent, n *npv1.NeighborRecord, st *DiscoveryStats) error { + return s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + localIfID, err := s.resolveLocalInterface(ctx, tx, a.TenantID, n) + if err != nil { + return err + } + + remoteDevID, conf, err := s.resolveRemoteDevice(ctx, tx, a.TenantID, n) + if err != nil { + return err + } + + var remoteIfID any + if remoteDevID != nil { + remoteIfID, err = s.resolveRemoteInterface(ctx, tx, a.TenantID, *remoteDevID, n.GetRemotePortId(), n.GetRemotePortDescr()) + if err != nil { + return err + } + } + + raw := n.GetRawJson() + if len(raw) == 0 { + raw = []byte("{}") + } else if !json.Valid(raw) { + raw = []byte("{}") + } + + _, err = tx.Exec(ctx, ` + INSERT INTO topo.neighbors + (tenant_id, device_id, interface_id, proto, + remote_chassis_id, remote_system_name, remote_port_id, remote_port_descr, + remote_mgmt_ip, remote_mac, remote_platform, remote_capabilities, + resolved_device_id, resolved_interface_id, confidence, last_seen_at, raw) + VALUES ($1,$2,$3,$4::topo.discovery_proto,$5,$6,$7,$8, + $9::inet,$10::macaddr,$11,$12,$13,$14,$15,now(),$16::jsonb) + ON CONFLICT (device_id, COALESCE(interface_id, '00000000-0000-0000-0000-000000000000'::uuid), + proto, COALESCE(remote_chassis_id,''), COALESCE(remote_port_id,'')) + DO UPDATE SET + remote_system_name = EXCLUDED.remote_system_name, + remote_port_descr = EXCLUDED.remote_port_descr, + remote_mgmt_ip = EXCLUDED.remote_mgmt_ip, + remote_mac = EXCLUDED.remote_mac, + remote_platform = EXCLUDED.remote_platform, + remote_capabilities = EXCLUDED.remote_capabilities, + resolved_device_id = EXCLUDED.resolved_device_id, + resolved_interface_id = EXCLUDED.resolved_interface_id, + confidence = EXCLUDED.confidence, + last_seen_at = now(), + raw = EXCLUDED.raw + `, a.TenantID, n.GetDeviceId(), localIfID, protoName(n.GetProto()), + nullString(n.GetRemoteChassisId()), nullString(n.GetRemoteSystemName()), + nullString(n.GetRemotePortId()), nullString(n.GetRemotePortDescr()), + nullString(n.GetRemoteMgmtIp()), nullString(n.GetRemoteMac()), + nullString(n.GetRemotePlatform()), n.GetRemoteCapabilities(), + remoteDevID, remoteIfID, conf, string(raw)) + if err != nil { + return err + } + + if remoteDevID == nil { + return nil + } + st.NeighborsResolved++ + + created, err := s.upsertLink(ctx, tx, a.TenantID, linkSpec{ + aDevice: n.GetDeviceId(), + aInterface: localIfID, + bDevice: *remoteDevID, + bInterface: remoteIfID, + proto: protoName(n.GetProto()), + confidence: conf, + }) + if err != nil { + return err + } + if created { + st.LinksCreated++ + } else { + st.LinksUpdated++ + } + return nil + }) +} + +func (s *Store) resolveLocalInterface(ctx context.Context, tx pgx.Tx, tenantID string, n *npv1.NeighborRecord) (any, error) { + if n.GetLocalInterfaceId() != "" { + return n.GetLocalInterfaceId(), nil + } + + var id string + var err error + switch { + case n.GetLocalIfIndex() > 0: + err = tx.QueryRow(ctx, ` + SELECT id::text FROM inv.interfaces + WHERE tenant_id = $1 AND device_id = $2 AND if_index = $3 + `, tenantID, n.GetDeviceId(), n.GetLocalIfIndex()).Scan(&id) + case n.GetLocalPortName() != "": + err = tx.QueryRow(ctx, ` + SELECT id::text FROM inv.interfaces + WHERE tenant_id = $1 AND device_id = $2 AND lower(name) = lower($3) + `, tenantID, n.GetDeviceId(), n.GetLocalPortName()).Scan(&id) + default: + return nil, nil + } + + if errors.Is(err, pgx.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, err + } + return id, nil +} + +// resolveRemoteDevice шукає пристрій за спаданням надійності ознаки. +func (s *Store) resolveRemoteDevice(ctx context.Context, tx pgx.Tx, tenantID string, n *npv1.NeighborRecord) (*string, int, error) { + type probe struct { + query string + arg string + conf int + } + + probes := []probe{ + {`SELECT id::text FROM inv.devices + WHERE tenant_id=$1 AND deleted_at IS NULL AND chassis_id = $2 LIMIT 1`, + n.GetRemoteChassisId(), confChassis}, + + {`SELECT id::text FROM inv.devices + WHERE tenant_id=$1 AND deleted_at IS NULL AND base_mac = $2::macaddr LIMIT 1`, + n.GetRemoteMac(), confMAC}, + + // MAC сусіда може збігтися не з шасі, а з конкретним портом. + {`SELECT i.device_id::text FROM inv.interfaces i + WHERE i.tenant_id=$1 AND i.mac = $2::macaddr LIMIT 1`, + n.GetRemoteMac(), confMAC}, + + {`SELECT id::text FROM inv.devices + WHERE tenant_id=$1 AND deleted_at IS NULL AND host(address) = $2 LIMIT 1`, + n.GetRemoteMgmtIp(), confMgmtIP}, + + {`SELECT id::text FROM inv.devices + WHERE tenant_id=$1 AND deleted_at IS NULL + AND (system_name = $2 OR lower(name) = lower($2)) LIMIT 1`, + n.GetRemoteSystemName(), confSysName}, + } + + for _, p := range probes { + if p.arg == "" { + continue + } + var id string + err := tx.QueryRow(ctx, p.query, tenantID, p.arg).Scan(&id) + if errors.Is(err, pgx.ErrNoRows) { + continue + } + if err != nil { + // Некоректний MAC/IP від пристрою не має валити весь звіт. + continue + } + return &id, p.conf, nil + } + return nil, 0, nil +} + +// resolveRemoteInterface: port-id у LLDP буває чим завгодно — +// ім'ям порту, описом, ifIndex або MAC. Пробуємо все. +func (s *Store) resolveRemoteInterface(ctx context.Context, tx pgx.Tx, tenantID, deviceID, portID, portDescr string) (any, error) { + candidates := []string{portID, portDescr} + + for _, c := range candidates { + if strings.TrimSpace(c) == "" { + continue + } + var id string + + err := tx.QueryRow(ctx, ` + SELECT id::text FROM inv.interfaces + WHERE tenant_id=$1 AND device_id=$2 + AND (lower(name) = lower($3) OR lower(alias) = lower($3)) + LIMIT 1 + `, tenantID, deviceID, c).Scan(&id) + if err == nil { + return id, nil + } + if !errors.Is(err, pgx.ErrNoRows) { + return nil, err + } + + if idx, convErr := strconv.ParseInt(strings.TrimSpace(c), 10, 64); convErr == nil { + err = tx.QueryRow(ctx, ` + SELECT id::text FROM inv.interfaces + WHERE tenant_id=$1 AND device_id=$2 AND if_index=$3 LIMIT 1 + `, tenantID, deviceID, idx).Scan(&id) + if err == nil { + return id, nil + } + if !errors.Is(err, pgx.ErrNoRows) { + return nil, err + } + } + } + return nil, nil +} + +type linkSpec struct { + aDevice string + aInterface any + bDevice string + bInterface any + proto string + confidence int +} + +// upsertLink створює або оновлює лінк, нормалізуючи пару. +// +// Пошук іде за тим самим виразом LEAST/GREATEST, що й унікальний +// індекс links_pair_uniq, тому A→B і B→A знаходять один рядок. +// Лінк із is_pinned не перезаписується: підтверджене людиною має +// пріоритет над автовиявленням, інакше кожен запуск discovery +// затирав би ручні правки. +func (s *Store) upsertLink(ctx context.Context, tx pgx.Tx, tenantID string, spec linkSpec) (created bool, err error) { + const zero = "00000000-0000-0000-0000-000000000000" + + var ( + linkID string + isPinned bool + ) + err = tx.QueryRow(ctx, ` + SELECT id::text, is_pinned FROM topo.links + WHERE tenant_id = $1 + AND LEAST(a_device_id, b_device_id) = LEAST($2::uuid, $3::uuid) + AND GREATEST(a_device_id, b_device_id) = GREATEST($2::uuid, $3::uuid) + AND LEAST(COALESCE(a_interface_id, $6::uuid), COALESCE(b_interface_id, $6::uuid)) + = LEAST(COALESCE($4::uuid, $6::uuid), COALESCE($5::uuid, $6::uuid)) + AND GREATEST(COALESCE(a_interface_id, $6::uuid), COALESCE(b_interface_id, $6::uuid)) + = GREATEST(COALESCE($4::uuid, $6::uuid), COALESCE($5::uuid, $6::uuid)) + `, tenantID, spec.aDevice, spec.bDevice, spec.aInterface, spec.bInterface, zero). + Scan(&linkID, &isPinned) + + switch { + case err == nil: + if isPinned { + // Лише позначаємо, що лінк живий. + _, err = tx.Exec(ctx, ` + UPDATE topo.links SET last_seen_at = now() WHERE id = $1 + `, linkID) + return false, err + } + _, err = tx.Exec(ctx, ` + UPDATE topo.links + SET last_seen_at = now(), + confidence = GREATEST(confidence, $2), + discovered_by = $3::topo.discovery_proto + WHERE id = $1 + `, linkID, spec.confidence, spec.proto) + return false, err + + case errors.Is(err, pgx.ErrNoRows): + // Швидкість каналу беремо з порту: вона стане знаменником + // для util_pct і, зрештою, швидкістю анімації на мапі. + _, err = tx.Exec(ctx, ` + INSERT INTO topo.links + (tenant_id, a_device_id, a_interface_id, b_device_id, b_interface_id, + kind, capacity_bps, discovered_by, confidence, status, last_seen_at) + VALUES ($1,$2,$3,$4,$5,'physical', + (SELECT speed_bps FROM inv.interfaces WHERE id = $3::uuid), + $6::topo.discovery_proto, $7, 'unknown', now()) + ON CONFLICT DO NOTHING + `, tenantID, spec.aDevice, spec.aInterface, spec.bDevice, spec.bInterface, + spec.proto, spec.confidence) + return err == nil, err + + default: + return false, err + } +} + +func protoName(p npv1.DiscoveryProto) string { + switch p { + case npv1.DiscoveryProto_DISCOVERY_PROTO_LLDP: + return "lldp" + case npv1.DiscoveryProto_DISCOVERY_PROTO_CDP: + return "cdp" + case npv1.DiscoveryProto_DISCOVERY_PROTO_ARP: + return "arp" + case npv1.DiscoveryProto_DISCOVERY_PROTO_FDB: + return "fdb" + case npv1.DiscoveryProto_DISCOVERY_PROTO_STP: + return "stp" + case npv1.DiscoveryProto_DISCOVERY_PROTO_ROUTING: + return "routing" + case npv1.DiscoveryProto_DISCOVERY_PROTO_SNMP_TOPO: + return "snmp_topo" + default: + return "manual" + } +} + +func statusOr(s, def string) string { + if s == "" { + return def + } + return s +} + +func nullInt(v int) any { + if v == 0 { + return nil + } + return v +} + +func nullInt64(v int64) any { + if v == 0 { + return nil + } + return v +} diff --git a/server/internal/store/ncm.go b/server/internal/store/ncm.go new file mode 100644 index 0000000..9febe21 --- /dev/null +++ b/server/internal/store/ncm.go @@ -0,0 +1,182 @@ +package store + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "strings" + + "github.com/jackc/pgx/v5" + "github.com/netpulse/netpulse/server/internal/crypto" +) + +// ConfigSubmission — зібраний із чанків конфіг. +type ConfigSubmission struct { + JobID string + DeviceID string + ConfigType string + Body []byte + // sha256 тіла, як його порахував агент. + ClaimedHash []byte +} + +// ConfigOutcome — що сервер із цим зробив. +type ConfigOutcome struct { + Accepted bool + Unchanged bool + ConfigID string + CommitSHA string + Err error +} + +var ErrChecksumMismatch = errors.New("контрольна сума не збіглася") + +// StoreConfig приймає конфіг: звіряє суму, визначає, чи є зміна, +// і зберігає нову версію. +// +// СТАН РЕАЛІЗАЦІЇ: Git-двигун (libgit2) ще не підключено. Тіло +// зберігається зашифрованим у core.secrets, а commit_sha тимчасово +// містить hex контентного хеша. Усе інше — дедуплікація, підрахунок +// рядків, ланцюжок prev_config_id — працює вже зараз, тому diff між +// версіями будується без Git. Коли двигун з'явиться, зміниться лише +// джерело commit_sha. +func (s *Store) StoreConfig(ctx context.Context, a *Agent, sub ConfigSubmission, ring *crypto.Keyring) (ConfigOutcome, error) { + sum := sha256.Sum256(sub.Body) + if len(sub.ClaimedHash) > 0 && !equalBytes(sum[:], sub.ClaimedHash) { + return ConfigOutcome{Err: ErrChecksumMismatch}, nil + } + + out := ConfigOutcome{} + + err := s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + configType := sub.ConfigType + if configType == "" { + configType = "running" + } + + // Чи змінилось? Порівняння за хешем нормалізованого тіла — + // саме тому scrub робиться на сервері: інакше кожен збір + // виглядав би зміною через рядок з uptime. + var ( + prevID string + prevHash []byte + ) + err := tx.QueryRow(ctx, ` + SELECT id::text, content_hash + FROM ncm.configs + WHERE device_id = $1 AND tenant_id = $2 AND config_type = $3 + ORDER BY collected_at DESC + LIMIT 1 + `, sub.DeviceID, a.TenantID, configType).Scan(&prevID, &prevHash) + + hasPrev := err == nil + if err != nil && !errors.Is(err, pgx.ErrNoRows) { + return err + } + + if hasPrev && equalBytes(prevHash, sum[:]) { + out.Accepted = true + out.Unchanged = true + out.ConfigID = prevID + return nil + } + + var repoID string + err = tx.QueryRow(ctx, ` + SELECT id::text FROM ncm.repos WHERE tenant_id = $1 ORDER BY created_at LIMIT 1 + `, a.TenantID).Scan(&repoID) + if errors.Is(err, pgx.ErrNoRows) { + // Репозиторій заводиться при першому ж бекапі. + // $1 тут потрібен і як uuid (колонка tenant_id), і як text + // (частина шляху). Приведення на місці не допомагає — воно + // не розв'язує конфлікт, а нав'язує один тип обом уживанням. + // Тому лишаємо параметр текстовим і кастуємо там, де треба uuid. + err = tx.QueryRow(ctx, ` + INSERT INTO ncm.repos (tenant_id, name, storage_path) + VALUES ($1::uuid, 'default', '/var/lib/netpulse/git/' || $1 || '.git') + RETURNING id::text + `, a.TenantID).Scan(&repoID) + } + if err != nil { + return err + } + + // Тіло — секрет: у конфігах живуть ключі, community-рядки + // і хеші паролів. Кладемо зашифрованим тим самим механізмом, + // що й креденшели. + aad := a.TenantID + "|config|" + sub.DeviceID + sec, err := ring.Encrypt(sub.Body, aad) + if err != nil { + return err + } + + var secretID string + err = tx.QueryRow(ctx, ` + INSERT INTO core.secrets (tenant_id, kind, key_id, nonce, ciphertext, auth_tag, aad) + VALUES ($1, 'generic', $2, $3, $4, $5, $6) + RETURNING id::text + `, a.TenantID, sec.KeyID, sec.Nonce, sec.Ciphertext, sec.AuthTag, aad).Scan(&secretID) + if err != nil { + return fmt.Errorf("збереження зашифрованого тіла: %w", err) + } + + hexSum := hex.EncodeToString(sum[:]) + lines := strings.Count(string(sub.Body), "\n") + 1 + + var deviceName string + if err := tx.QueryRow(ctx, ` + SELECT name FROM inv.devices WHERE id = $1 AND tenant_id = $2 + `, sub.DeviceID, a.TenantID).Scan(&deviceName); err != nil { + return fmt.Errorf("пристрій %s: %w", sub.DeviceID, err) + } + + path := fmt.Sprintf("%s/%s.cfg", deviceName, configType) + + var prevArg any + if hasPrev { + prevArg = prevID + } + + err = tx.QueryRow(ctx, ` + INSERT INTO ncm.configs + (tenant_id, device_id, repo_id, commit_sha, blob_sha, branch, path, + config_type, size_bytes, line_count, content_hash, body_secret_id, + prev_config_id, is_change) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13::uuid,$14) + RETURNING id::text + `, a.TenantID, sub.DeviceID, repoID, hexSum, hexSum, + "device/"+deviceName, path, configType, + len(sub.Body), lines, sum[:], secretID, prevArg, hasPrev).Scan(&out.ConfigID) + if err != nil { + return fmt.Errorf("вставка версії конфігу: %w", err) + } + + if _, err := tx.Exec(ctx, ` + UPDATE ncm.device_policies + SET last_backup_at = now() + WHERE device_id = $1 AND tenant_id = $2 + `, sub.DeviceID, a.TenantID); err != nil { + return err + } + + out.Accepted = true + out.CommitSHA = hexSum + return nil + }) + + return out, err +} + +func equalBytes(a, b []byte) bool { + if len(a) != len(b) { + return false + } + for i := range a { + if a[i] != b[i] { + return false + } + } + return true +} diff --git a/server/internal/store/plan.go b/server/internal/store/plan.go new file mode 100644 index 0000000..01d833d --- /dev/null +++ b/server/internal/store/plan.go @@ -0,0 +1,155 @@ +package store + +import ( + "context" + "crypto/sha256" + "fmt" + "hash/fnv" + "time" + + "github.com/jackc/pgx/v5" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" + "google.golang.org/protobuf/types/known/durationpb" +) + +// ScheduleOffset розводить задачі всередині інтервалу. +// +// Рахує СЕРВЕР, і рахує детерміновано від check_id: інакше після +// кожного перезапуску агента 5000 чеків з інтервалом 60 с з'їжджали б +// у нову випадкову фазу, а раз на хвилину мережа отримувала б сплеск. +// Функція чиста — те саме значення на будь-якому вузлі сервера. +func ScheduleOffset(checkID string, interval time.Duration) time.Duration { + if interval <= 0 { + return 0 + } + h := fnv.New64a() + _, _ = h.Write([]byte(checkID)) + return time.Duration(h.Sum64() % uint64(interval)) +} + +// BuildPlan збирає повний план задач для зонда. +func (s *Store) BuildPlan(ctx context.Context, a *Agent) (*npv1.TaskPlan, error) { + var plan *npv1.TaskPlan + + err := s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error { + rows, err := tx.Query(ctx, ` + SELECT c.id::text, + c.device_id::text, + COALESCE(c.interface_id::text, ''), + c.check_type, + c.params::text, + c.interval_sec, + c.timeout_ms, + c.retries, + d.name, + COALESCE(host(d.address), COALESCE(d.fqdn, '')), + COALESCE(d.chassis_id, ''), + COALESCE(d.system_name, '') + FROM core.checks c + JOIN inv.devices d ON d.id = c.device_id + WHERE c.tenant_id = $1 + AND d.agent_id = $2 + AND c.enabled + AND d.enabled + AND d.deleted_at IS NULL + ORDER BY c.id + `, a.TenantID, a.ID) + if err != nil { + return err + } + defer rows.Close() + + p := &npv1.TaskPlan{Final: true} + devices := make(map[string]*npv1.DeviceTarget) + hasher := sha256.New() + + for rows.Next() { + var ( + checkID, deviceID, ifaceID, checkType, params string + intervalSec, timeoutMs, retries int32 + devName, devAddr, chassisID, sysName string + ) + if err := rows.Scan(&checkID, &deviceID, &ifaceID, &checkType, ¶ms, + &intervalSec, &timeoutMs, &retries, + &devName, &devAddr, &chassisID, &sysName); err != nil { + return err + } + + interval := time.Duration(intervalSec) * time.Second + task := &npv1.Task{ + CheckId: checkID, + DeviceId: deviceID, + InterfaceId: ifaceID, + CheckType: checkType, + ParamsJson: []byte(params), + Interval: durationpb.New(interval), + Timeout: durationpb.New(time.Duration(timeoutMs) * time.Millisecond), + Retries: uint32(retries), + Enabled: true, + ScheduleOffset: durationpb.New(ScheduleOffset(checkID, interval)), + } + p.Tasks = append(p.Tasks, task) + + // Хеш плану рахуємо з полів, які реально впливають на + // поведінку агента. Зміна опису пристрою не має змушувати + // переливати 50 000 задач. + fmt.Fprintf(hasher, "%s|%s|%s|%s|%d|%d|%d|%s\n", + checkID, deviceID, ifaceID, checkType, + intervalSec, timeoutMs, retries, params) + + if _, ok := devices[deviceID]; !ok { + devices[deviceID] = &npv1.DeviceTarget{ + DeviceId: deviceID, + Name: devName, + Address: devAddr, + ChassisId: chassisID, + SystemName: sysName, + } + } + } + if err := rows.Err(); err != nil { + return err + } + + for _, d := range devices { + p.Devices = append(p.Devices, d) + } + p.PlanHash = hasher.Sum(nil) + plan = p + return nil + }) + + return plan, err +} + +// ModulesForPlan визначає, які модулі треба активувати на зонді. +// +// Виводиться з типів чеків у плані плюс явно дозволених у +// core.agents.enabled_modules: зонд не має вмикати нічого, що йому +// не знадобиться, — це і пам'ять, і зайва поверхня атаки. +func ModulesForPlan(plan *npv1.TaskPlan, allowed []string) *npv1.ModuleControl { + needed := make(map[string]bool) + for _, t := range plan.GetTasks() { + if key, _, ok := cutDot(t.GetCheckType()); ok { + needed[key] = true + } + } + for _, m := range allowed { + needed[m] = true + } + + mc := &npv1.ModuleControl{Exclusive: true} + for key := range needed { + mc.Modules = append(mc.Modules, &npv1.ModuleSpec{Key: key, Enabled: true}) + } + return mc +} + +func cutDot(s string) (before, after string, found bool) { + for i := 0; i < len(s); i++ { + if s[i] == '.' { + return s[:i], s[i+1:], true + } + } + return s, "", false +} diff --git a/server/internal/store/store.go b/server/internal/store/store.go new file mode 100644 index 0000000..9664352 --- /dev/null +++ b/server/internal/store/store.go @@ -0,0 +1,89 @@ +// Package store — доступ до PostgreSQL/TimescaleDB. +// +// Головне правило шару: жоден запит не виконується без tenant_id. +// Ізоляція забезпечується двома незалежними механізмами, і це навмисно: +// +// 1. RLS у БД — політика tenant_isolation на кожній звичайній таблиці +// з колонкою tenant_id. Вмикається через SET LOCAL app.tenant_id. +// 2. Явний предикат tenant_id у кожному запиті цього пакета. +// +// Дублювання потрібне, бо RLS НЕ працює на гіпертаблях: TimescaleDB +// не поєднує row level security зі стисненням. Саме туди йде вся +// телеметрія, тому для неї другий механізм — єдиний. +package store + +import ( + "context" + "fmt" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" +) + +type Store struct { + pool *pgxpool.Pool +} + +func New(ctx context.Context, dsn string) (*Store, error) { + cfg, err := pgxpool.ParseConfig(dsn) + if err != nil { + return nil, fmt.Errorf("розбір DSN: %w", err) + } + cfg.MaxConns = 16 + cfg.MinConns = 2 + cfg.MaxConnLifetime = time.Hour + cfg.MaxConnIdleTime = 10 * time.Minute + + pool, err := pgxpool.NewWithConfig(ctx, cfg) + if err != nil { + return nil, err + } + if err := pool.Ping(ctx); err != nil { + pool.Close() + return nil, fmt.Errorf("ping: %w", err) + } + return &Store{pool: pool}, nil +} + +func (s *Store) Close() { s.pool.Close() } + +func (s *Store) Pool() *pgxpool.Pool { return s.pool } + +// InTenantTx виконує fn у транзакції з виставленим app.tenant_id. +// +// SET LOCAL, а не SET: значення живе рівно до кінця транзакції й не +// протікає на наступний запит, який візьме те саме з'єднання з пулу. +// Протікання тут означало б показ чужих даних, тому це не оптимізація, +// а вимога. +func (s *Store) InTenantTx(ctx context.Context, tenantID string, fn func(pgx.Tx) error) error { + tx, err := s.pool.Begin(ctx) + if err != nil { + return err + } + defer func() { _ = tx.Rollback(ctx) }() + + if _, err := tx.Exec(ctx, "SELECT set_config('app.tenant_id', $1, true)", tenantID); err != nil { + return fmt.Errorf("встановити app.tenant_id: %w", err) + } + if err := fn(tx); err != nil { + return err + } + return tx.Commit(ctx) +} + +// nullUUID перетворює порожній рядок на NULL — protobuf не має +// nullable-рядків, тому "" з дроту означає «не задано». +func nullUUID(s string) any { + if s == "" { + return nil + } + return s +} + +func nullString(s string) any { + if s == "" { + return nil + } + return s +} diff --git a/server/internal/store/telemetry.go b/server/internal/store/telemetry.go new file mode 100644 index 0000000..62297c3 --- /dev/null +++ b/server/internal/store/telemetry.go @@ -0,0 +1,319 @@ +package store + +import ( + "context" + "encoding/json" + "fmt" + "sync" + "time" + + "github.com/jackc/pgx/v5" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" +) + +// SeriesTable — мапа series_ref → ts.series.id у межах ОДНІЄЇ сесії. +// +// Агент нумерує серії локально й після реконекту починає з 1, тому +// таблиця не переживає сесію. Якщо сервер отримав семпл із номером, +// якого не знає, — це не привід мовчки викинути дані: він відповідає +// reset_series_table і агент реєструє все заново. +type SeriesTable struct { + mu sync.RWMutex + byRef map[uint32]int64 +} + +func NewSeriesTable() *SeriesTable { + return &SeriesTable{byRef: make(map[uint32]int64)} +} + +func (t *SeriesTable) get(ref uint32) (int64, bool) { + t.mu.RLock() + defer t.mu.RUnlock() + id, ok := t.byRef[ref] + return id, ok +} + +func (t *SeriesTable) put(ref uint32, id int64) { + t.mu.Lock() + t.byRef[ref] = id + t.mu.Unlock() +} + +func (t *SeriesTable) Len() int { + t.mu.RLock() + defer t.mu.RUnlock() + return len(t.byRef) +} + +// ErrUnknownSeriesRef — семпл посилається на незареєстровану серію. +type ErrUnknownSeriesRef struct{ Ref uint32 } + +func (e *ErrUnknownSeriesRef) Error() string { + return fmt.Sprintf("series_ref %d не зареєстровано в цій сесії", e.Ref) +} + +// RegisterSeries створює або знаходить серії й наповнює таблицю сесії. +func (s *Store) RegisterSeries(ctx context.Context, tenantID string, descs []*npv1.SeriesDescriptor, table *SeriesTable) error { + if len(descs) == 0 { + return nil + } + + return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { + for _, d := range descs { + labels, err := json.Marshal(nonNilLabels(d.GetLabels())) + if err != nil { + return err + } + + var id int64 + err = tx.QueryRow(ctx, ` + INSERT INTO ts.series + (tenant_id, device_id, interface_id, plugin_key, metric_key, unit, labels) + VALUES ($1, $2, $3, $4, $5, $6, $7::jsonb) + ON CONFLICT (tenant_id, device_id, metric_key, labels_hash) + DO UPDATE SET + unit = COALESCE(NULLIF(EXCLUDED.unit, ''), ts.series.unit), + interface_id = COALESCE(EXCLUDED.interface_id, ts.series.interface_id) + RETURNING id + `, tenantID, + nullUUID(d.GetDeviceId()), + nullUUID(d.GetInterfaceId()), + nullString(d.GetPluginKey()), + d.GetMetricKey(), + d.GetUnit(), + string(labels), + ).Scan(&id) + if err != nil { + return fmt.Errorf("реєстрація серії %s: %w", d.GetMetricKey(), err) + } + table.put(d.GetSeriesRef(), id) + } + return nil + }) +} + +func nonNilLabels(m map[string]string) map[string]string { + if m == nil { + return map[string]string{} + } + return m +} + +// BatchStats — що саме записалось; іде в журнал і метрики сервера. +type BatchStats struct { + Samples int + Icmp int + Interfaces int + Statuses int + Checks int +} + +// WriteBatch записує телеметрію. +// +// Семантика доставки — at-least-once, тому кожен запис іде з +// ON CONFLICT DO NOTHING: повторний батч після реконекту не має ані +// падати, ані дублювати рядки. Первинні ключі (ts, device_id), +// (ts, interface_id) і (ts, series_id) роблять це безпечним. +func (s *Store) WriteBatch(ctx context.Context, a *Agent, b *npv1.TelemetryBatch, table *SeriesTable) (BatchStats, error) { + var st BatchStats + + // Спершу — нові серії, бо семпли в цьому ж батчі на них посилаються. + if err := s.RegisterSeries(ctx, a.TenantID, b.GetNewSeries(), table); err != nil { + return st, err + } + + batch := &pgx.Batch{} + + for _, r := range b.GetIcmp() { + batch.Queue(` + INSERT INTO ts.icmp_samples + (ts, device_id, tenant_id, agent_id, rtt_avg_ms, rtt_min_ms, rtt_max_ms, + jitter_ms, loss_pct, packets_sent, packets_recv, reachable) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12) + ON CONFLICT (ts, device_id) DO NOTHING + `, r.GetTs().AsTime(), r.GetDeviceId(), a.TenantID, a.ID, + r.GetRttAvgMs(), r.GetRttMinMs(), r.GetRttMaxMs(), r.GetJitterMs(), + r.GetLossPct(), int16(r.GetPacketsSent()), int16(r.GetPacketsRecv()), + r.GetReachable()) + st.Icmp++ + } + + for _, r := range b.GetInterfaces() { + batch.Queue(` + INSERT INTO ts.if_counters + (ts, interface_id, device_id, tenant_id, + in_octets, out_octets, in_ucast_pkts, out_ucast_pkts, + in_errors, out_errors, in_discards, out_discards, + in_bps, out_bps, in_pps, out_pps, + util_in_pct, util_out_pct, oper_up) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19) + ON CONFLICT (ts, interface_id) DO NOTHING + `, r.GetTs().AsTime(), r.GetInterfaceId(), r.GetDeviceId(), a.TenantID, + int64(r.GetInOctets()), int64(r.GetOutOctets()), + int64(r.GetInUcastPkts()), int64(r.GetOutUcastPkts()), + int64(r.GetInErrors()), int64(r.GetOutErrors()), + int64(r.GetInDiscards()), int64(r.GetOutDiscards()), + r.GetInBps(), r.GetOutBps(), r.GetInPps(), r.GetOutPps(), + r.GetUtilInPct(), r.GetUtilOutPct(), r.GetOperUp()) + st.Interfaces++ + } + + for _, smp := range b.GetSamples() { + seriesID, ok := table.get(smp.GetSeriesRef()) + if !ok { + return st, &ErrUnknownSeriesRef{Ref: smp.GetSeriesRef()} + } + batch.Queue(` + INSERT INTO ts.samples (ts, series_id, value) + VALUES ($1,$2,$3) + ON CONFLICT (ts, series_id) DO NOTHING + `, smp.GetTs().AsTime(), seriesID, smp.GetValue()) + st.Samples++ + } + + if batch.Len() > 0 { + res := s.pool.SendBatch(ctx, batch) + for i := 0; i < batch.Len(); i++ { + if _, err := res.Exec(); err != nil { + res.Close() + return st, fmt.Errorf("запис телеметрії[%d]: %w", i, err) + } + } + if err := res.Close(); err != nil { + return st, err + } + } + + // Стан пристроїв оновлюємо окремо: там є умовний запис в історію, + // який не виражається одним INSERT. + for _, sc := range b.GetStatusChanges() { + if err := s.applyDeviceStatus(ctx, a.TenantID, sc.GetDeviceId(), + statusName(sc.GetStatus()), sc.GetTs().AsTime(), sc.GetReason()); err != nil { + return st, err + } + st.Statuses++ + } + + // Стан із ICMP — головне джерело кольору вузла на мапі. + for _, r := range b.GetIcmp() { + want := "down" + if r.GetReachable() { + want = "up" + } + if err := s.applyDeviceStatus(ctx, a.TenantID, r.GetDeviceId(), want, + r.GetTs().AsTime(), "icmp"); err != nil { + return st, err + } + } + + st.Checks = len(b.GetCheckResults()) + return st, nil +} + +// applyDeviceStatus змінює стан пристрою й пише в історію ЛИШЕ при +// фактичному переході. +// +// Один запит замість read-modify-write: інакше два воркери, що +// обробляють сусідні батчі, наввипередки писали б у історію +// неіснуючі переходи up→up. +func (s *Store) applyDeviceStatus(ctx context.Context, tenantID, deviceID, status string, ts time.Time, reason string) error { + if deviceID == "" || status == "" { + return nil + } + if ts.IsZero() { + ts = time.Now() + } + + return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error { + _, err := tx.Exec(ctx, ` + WITH prev AS ( + SELECT id, status + FROM inv.devices + WHERE id = $1 AND tenant_id = $2 + FOR UPDATE + ), upd AS ( + UPDATE inv.devices d + SET status = $3::inv.device_status, + status_changed_at = CASE + WHEN d.status IS DISTINCT FROM $3::inv.device_status THEN $4 + ELSE d.status_changed_at END, + last_seen_at = GREATEST(COALESCE(d.last_seen_at, $4), $4) + FROM prev + WHERE d.id = prev.id + RETURNING prev.status AS old_status, d.status AS new_status + ) + INSERT INTO ts.device_status_history (ts, device_id, tenant_id, status, prev_status, reason) + SELECT $4, $1, $2, new_status, old_status, $5 + FROM upd + WHERE old_status IS DISTINCT FROM new_status + ON CONFLICT (ts, device_id) DO NOTHING + `, deviceID, tenantID, status, ts, nullString(reason)) + return err + }) +} + +func statusName(s npv1.Status) string { + switch s { + case npv1.Status_STATUS_UP: + return "up" + case npv1.Status_STATUS_DOWN: + return "down" + case npv1.Status_STATUS_WARNING: + return "warning" + case npv1.Status_STATUS_MAINTENANCE: + return "maintenance" + case npv1.Status_STATUS_UNKNOWN: + return "unknown" + default: + return "" + } +} + +// WriteLogs зберігає syslog і трапи. +func (s *Store) WriteLogs(ctx context.Context, a *Agent, b *npv1.LogBatch) (int, error) { + batch := &pgx.Batch{} + + for _, e := range b.GetSyslog() { + parsed, _ := json.Marshal(nonNilLabels(e.GetParsed())) + batch.Queue(` + INSERT INTO ts.syslog + (ts, tenant_id, device_id, source_ip, facility, severity, hostname, tag, message, parsed) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10::jsonb) + `, tsOrNow(e.GetTs().AsTime()), a.TenantID, nullUUID(e.GetDeviceId()), + nullString(e.GetSourceIp()), int16(e.GetFacility()), int16(e.GetSeverity()), + nullString(e.GetHostname()), nullString(e.GetTag()), e.GetMessage(), string(parsed)) + } + + for _, t := range b.GetTraps() { + vb := make(map[string]string, len(t.GetVarbinds())) + for _, v := range t.GetVarbinds() { + vb[v.GetOid()] = v.GetValue() + } + payload, _ := json.Marshal(vb) + batch.Queue(` + INSERT INTO ts.snmp_traps (ts, tenant_id, device_id, source_ip, trap_oid, varbinds) + VALUES ($1,$2,$3,$4,$5,$6::jsonb) + `, tsOrNow(t.GetTs().AsTime()), a.TenantID, nullUUID(t.GetDeviceId()), + nullString(t.GetSourceIp()), nullString(t.GetTrapOid()), string(payload)) + } + + if batch.Len() == 0 { + return 0, nil + } + + res := s.pool.SendBatch(ctx, batch) + defer res.Close() + for i := 0; i < batch.Len(); i++ { + if _, err := res.Exec(); err != nil { + return i, fmt.Errorf("запис журналу[%d]: %w", i, err) + } + } + return batch.Len(), nil +} + +func tsOrNow(t time.Time) time.Time { + if t.IsZero() { + return time.Now() + } + return t +}