Netpulse_SasS/agent/internal/session/logs.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

191 lines
6.9 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 session
import (
"context"
"net"
"strings"
"time"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
)
// maxLogBatch — стеля подій в одній пачці.
//
// Журнал не телеметрія: під час аварії він приходить сплеском, і
// відправляти по одній події означає накрутити тисячі викликів рівно
// тоді, коли мережа й так не в порядку.
const maxLogBatch = 500
// logFlush — як довго накопичувати пачку.
//
// Півсекунди — компроміс: тригер бекапу за «%SYS-5-CONFIG_I» має
// спрацьовувати відчутно швидше за хвилинний цикл, а слати кожен рядок
// окремо немає сенсу.
const logFlush = 500 * time.Millisecond
// logsLoop віддає накопичені події серверу.
//
// Окремий стрім, а не контрольний канал: сплеск логів під час аварії не
// має заважати heartbeat і командам. Саме тому в контракті StreamLogs
// існує окремо від Control.
// Syslog і трапи їдуть ОДНИМ стрімом, а не двома.
//
// Спокуса завести другий є: приймачі різні, порти різні, розбір різний.
// Але за межами зонда це та сама подія з мережі — вона лягає в сусідні
// таблиці, звіряється з тими самими подієвими правилами й приїжджає в
// той самий момент аварії. Другий стрім означав би другий комплект
// підтверджень, лімітів і черг переповнення — і два різні місця, у
// яких по-різному вирішено, що робити з пачкою, яку не вдалося
// відправити.
func (s *Session) logsLoop(ctx context.Context, client npv1.AgentServiceClient) error {
sys, trp := s.cfg.Syslog, s.cfg.Traps
if sys == nil && trp == nil {
return nil
}
stream, err := client.StreamLogs(ctx)
if err != nil {
return err
}
// Читач підтверджень окремою горутиною: Send і Recv на одному
// стрімі — єдине, що gRPC дозволяє робити паралельно.
ackErr := make(chan error, 1)
go func() {
for {
ack, err := stream.Recv()
if err != nil {
ackErr <- err
return
}
if e := ack.GetError(); e != nil {
s.log.Warn("сервер не прийняв журнал",
"code", e.GetCode(), "msg", e.GetMessage())
continue
}
// Ліміти задає сервер: він бачить картину по всіх зондах і
// краще знає, що вважати шумом. Стеля на джерело спільна
// для обох приймачів — шумить не протокол, а пристрій.
if sys != nil {
sys.ApplyAck(ack.GetMinSeverity(), ack.GetRateLimitPerSource())
}
if trp != nil {
trp.ApplyAck(ack.GetRateLimitPerSource())
}
}
}()
// Один із приймачів може бути вимкнений, тому канали готовності
// беремо через nil-заглушку: читання з nil-каналу блокується
// назавжди, і саме це в select потрібно — гілка, яка ніколи не
// спрацює, замість гілки, якої немає.
var sysReady, trapReady <-chan struct{}
if sys != nil {
sysReady = sys.Ready()
}
if trp != nil {
trapReady = trp.Ready()
}
ticker := time.NewTicker(logFlush)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return nil
case err := <-ackErr:
return err
case <-sysReady:
// Подія з'явилась — але не летимо одразу: даємо тіку
// зібрати сусідів у ту саму пачку.
case <-trapReady:
case <-ticker.C:
}
var (
entries []*npv1.SyslogEntry
traps []*npv1.SnmpTrap
dropped uint64
)
if sys != nil {
var d uint64
entries, d = sys.Drain(maxLogBatch)
dropped += d
}
if trp != nil {
// Половина пачки на трапи — не арифметика, а пріоритет:
// трапів за секунду на порядок менше, ніж рядків журналу,
// і стеля тут потрібна лише на випадок шторму. Витіснити
// журнал вони не мають.
var d uint64
traps, d = trp.Drain(maxLogBatch / 2)
dropped += d
}
if len(entries) == 0 && len(traps) == 0 && dropped == 0 {
continue
}
batch := &npv1.LogBatch{
BatchId: s.nextBatch.Add(1),
AgentId: s.cfg.AgentID,
Syslog: entries,
Traps: traps,
Dropped: dropped,
}
if err := stream.Send(batch); err != nil {
// Невідправлене повертаємо в чергу: наступна сесія
// доставить. Порядок зберігається — журнал читають
// хронологічно.
if sys != nil {
sys.Requeue(entries)
}
if trp != nil {
trp.Requeue(traps)
}
return err
}
}
}
// resolveDeviceByIP шукає хост за адресою джерела.
//
// Зіставлення на зонді, а не на сервері: у сервера немає контексту
// мережі клієнта, а один і той самий приватний діапазон трапляється в
// десятках кабінетів. Зонд же має свіжий список своїх хостів.
//
// Не знайшли — не біда: подія все одно доїде з порожнім device_id і
// заповненим source_ip. Викидати журнал через незнайому адресу означало
// б утратити рівно те, що показує появу нового заліза в мережі.
func (s *Session) resolveDeviceByIP(ip string) string {
if ip == "" {
return ""
}
addr := net.ParseIP(ip)
s.devMu.RLock()
defer s.devMu.RUnlock()
for id, d := range s.devices {
target := strings.TrimSpace(d.GetAddress())
if target == "" {
continue
}
if target == ip {
return id
}
// Адреса хоста може бути записана з портом або як FQDN —
// порівнюємо ще й розібраний варіант.
if addr != nil {
if host, _, err := net.SplitHostPort(target); err == nil {
if net.ParseIP(host).Equal(addr) {
return id
}
}
if net.ParseIP(target).Equal(addr) {
return id
}
}
}
return ""
}