Один коміт, а не десяток тематичних, свідомо: теми переплетені в
спільних файлах (store.go, docker-compose.yml, deploy/README.md), і
розділити їх можна було б лише індексуванням шматків. Коміти, які не
збираються, гірші за один великий — тим паче що це рівно той стан, який
перевірявся разом.
ЩО ПРАЦЮЄ НА СТЕНДІ Й ПЕРЕВІРЕНО ТАМ
0058 подієві алерти: syslog, ncm, compliance спрацьовують у мить
події; правило з нереалізованим джерелом більше не зберігається
мовчки
0059 snmp.walk і прототипи шаблонів — таблиці з динамічним індексом
описуються шаблоном, а не Go
0060 відкат конфігу: план як різниця, маскування паролів із підписом
плану, обов'язковий контрольний збір, verifying при обриві
0061 кнопки Telegram: довге опитування, авторизація не з callback_data
0062 аудит і архів хостів; тест на AST, що падає на ключі без назви
0063 RLS: три ролі, окремий пул для фонових тактів
0064 строки зберігання даних і сторінка сховища
0065 приймач SNMP-трапів; перевірено справжніми пакетами по дроту,
переклад v1→v2 за RFC 3584 дає правильний OID
0066 ескалації сповіщень
0067 алерт про вичерпання диска
0068 поля заливки конфігу переїхали в каталог профілів
Плюс: 137 тестів вебу з нуля (їх не було взагалі), одинадцять справжніх
вад, знайдених ними й виправлених, і виправлення двох інтеграційних
тестів grpcapi, які мовчки пропускались півтора року.
ЩО ЩЕ НЕ ЗАПУСКАЛОСЬ
netpulse установник: одна команда замість 18 змінних і
593 рядків інструкції
RLS з першого запуску нова інсталяція під політиками одразу;
RLS-EXISTING-INSTALL.md лишається тільки для
старих інсталяцій
.forgejo + CI раннер не зареєстрований
Ці три перевірені компіляцією й міркуванням, але не виконанням.
ГОЛОВНИЙ ВИСНОВОК ДВОХ СЕСІЙ
Зелена перевірка доводить рівно те, що вона перевіряє. Тест ізоляції RLS
був правильний і зелений — і пропустив зламаний вхід, бо перевіряв «чи
не видно чужого», коли зламалось «чи видно своє». Інтеграційні тести
grpcapi були зелені, бо не виконувались. Схема, довідник і протокол
описували те, чого в коді не існувало, і виглядало це як готове.
Тому в кожному завданні цих сесій стояла вимога назвати НЕПОКРИТЕ, а
чотири задачі закінчились не можливістю, а відмовою: правило з
нереалізованим джерелом не зберігається, профіль без команд заливки
каже про це замість мовчазної кнопки, міграція RLS валить сама себе на
таблиці без політики, тест словника аудиту падає на ключі без назви.
Подробиці — HISTORY.md, розділи за 26 і 27 серпня.
328 lines
11 KiB
Go
328 lines
11 KiB
Go
package store
|
||
|
||
import (
|
||
"context"
|
||
"crypto/sha256"
|
||
"crypto/subtle"
|
||
"encoding/json"
|
||
"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.bg.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,
|
||
-- Черга планувальника зонда. Зонд шле її в кожному
|
||
-- heartbeat від початку, а в зведенні її не було: без
|
||
-- неї видно, скільки задач виконується, і не видно,
|
||
-- скільки не встигло початись.
|
||
'tasks_queued', $8::int
|
||
)
|
||
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(), int32(hb.GetTasksQueued()))
|
||
|
||
// Пакет виконується в тенантній транзакції, хоча перший запит іде в
|
||
// гіпертаблицю, де RLS немає. Річ у другому: core.agents під
|
||
// tenant_isolation, і без app.tenant_id цей UPDATE зачіпав би нуль
|
||
// рядків — мовчки. Зонд слав би heartbeat кожні кілька секунд, а в
|
||
// інтерфейсі назавжди лишався б offline із застиглим часом
|
||
// останнього зв'язку. Помилки при цьому не було б ніде: SendBatch
|
||
// не вважає «0 рядків» збоєм.
|
||
return s.InTenantTx(ctx, a.TenantID, func(tx pgx.Tx) error {
|
||
res := tx.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()
|
||
}
|
||
|
||
// AgentUpdate — що можна змінити в зонді з UI.
|
||
//
|
||
// Порожні поля не чіпаються: форма може надіслати лише ліміти, не
|
||
// торкаючись імені й модулів.
|
||
type AgentUpdate struct {
|
||
Name string
|
||
Modules []string
|
||
Limits *Limits
|
||
}
|
||
|
||
// UpdateAgent змінює налаштування зонда.
|
||
//
|
||
// Зміна доїжджає до живого зонда не миттєво: модулі й ліміти їдуть у
|
||
// Welcome при наступному підключенні, а план — звіркою раз на п'ять
|
||
// секунд. Це свідомо: тримати окрему чергу «керівних» повідомлень
|
||
// заради налаштувань, які міняють раз на місяць, дорожче за саму
|
||
// затримку.
|
||
func (s *Store) UpdateAgent(ctx context.Context, tenantID, agentID string, u AgentUpdate) error {
|
||
return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
|
||
if u.Name != "" {
|
||
if _, err := tx.Exec(ctx, `
|
||
UPDATE core.agents SET name = $3 WHERE id = $1 AND tenant_id = $2
|
||
`, agentID, tenantID, u.Name); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
if u.Modules != nil {
|
||
if _, err := tx.Exec(ctx, `
|
||
UPDATE core.agents SET enabled_modules = $3::core.slug[]
|
||
WHERE id = $1 AND tenant_id = $2
|
||
`, agentID, tenantID, u.Modules); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
if u.Limits != nil {
|
||
// Зливаємо в наявний jsonb, а не заміняємо: у limits можуть
|
||
// лежати поля, яких форма не знає, і затирати їх мовчки —
|
||
// найшвидший спосіб зламати те, чого не бачив.
|
||
patch := map[string]any{}
|
||
if u.Limits.MaxConcurrency > 0 {
|
||
patch["max_concurrency"] = u.Limits.MaxConcurrency
|
||
}
|
||
if u.Limits.IcmpRatePPS > 0 {
|
||
patch["icmp_rate_pps"] = u.Limits.IcmpRatePPS
|
||
}
|
||
if u.Limits.BatchSize > 0 {
|
||
patch["batch_size"] = u.Limits.BatchSize
|
||
}
|
||
if u.Limits.MaxInFlight > 0 {
|
||
patch["max_in_flight"] = u.Limits.MaxInFlight
|
||
}
|
||
if len(patch) > 0 {
|
||
b, err := json.Marshal(patch)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if _, err := tx.Exec(ctx, `
|
||
UPDATE core.agents SET limits = limits || $3::jsonb
|
||
WHERE id = $1 AND tenant_id = $2
|
||
`, agentID, tenantID, string(b)); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
}
|
||
return nil
|
||
})
|
||
}
|
||
|
||
// DeleteAgent прибирає зонд.
|
||
//
|
||
// Хости, які він опитував, лишаються без зонда (ON DELETE SET NULL) —
|
||
// видаляти їх разом означало б втратити історію через заміну заліза.
|
||
func (s *Store) DeleteAgent(ctx context.Context, tenantID, agentID string) error {
|
||
return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
|
||
ct, err := tx.Exec(ctx,
|
||
`DELETE FROM core.agents WHERE id = $1 AND tenant_id = $2`, agentID, tenantID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if ct.RowsAffected() == 0 {
|
||
return ErrNotFound
|
||
}
|
||
return nil
|
||
})
|
||
}
|