Netpulse_SasS/server/internal/store/telemetry.go
byrsapty ed8fc831bf Дві сесії роботи: 0058–0068, розгортання однією командою, тести
Один коміт, а не десяток тематичних, свідомо: теми переплетені в
спільних файлах (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 серпня.
2026-08-27 17:32:49 +03:00

328 lines
11 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package store
import (
"context"
"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.bg.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 {
// Історія й подія для UI пишуться в тій самій транзакції, що й
// сама зміна: інакше мапа могла б показати новий колір без
// запису в історії або навпаки.
_, 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
), changed AS (
SELECT old_status, new_status FROM upd
WHERE old_status IS DISTINCT FROM new_status
), hist AS (
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 changed
ON CONFLICT (ts, device_id) DO NOTHING
RETURNING 1
)
INSERT INTO core.event_outbox (tenant_id, topic, payload)
SELECT $2, 'device.status', jsonb_build_object(
'device_id', $1::text,
'status', new_status::text,
'previous_status', old_status::text,
'reason', $5::text,
'ts', $4::timestamptz
)
FROM changed
`, 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))
}
// Трапи складає traps.go: там же лежить і все, що з ними далі
// роблять (словник назв, невідомі відправники, читання журналу).
// Тримати форму запису тут, а форму читання там — найкоротший шлях
// до того, що вони розійдуться.
queueTraps(batch, a.TenantID, b.GetTraps())
if batch.Len() == 0 {
return 0, nil
}
res := s.bg.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
}