Приймач syslog на зонді: розбір RFC3164/5424, черга, ліміти
Модуль слухає UDP у мережі клієнта: комутатор у закритій мережі до сервера не достукається, а зонд уже має вихідний канал. Розбір трьома рівнями суворості — 5424, 3164 і «як є». Останній не запасний варіант, а робочий режим: пристрої, що шлють голий текст, існують, і втратити подію гірше, ніж зберегти саме повідомлення. Текст Cisco "%SYS-5-CONFIG_I: ..." свідомо НЕ ріжеться в поле tag — саме за ним шукають зміну конфігу. Ліміт частоти на джерело, а не спільний: комутатор, що зациклився на помилці, інакше витіснив би з журналу всю решту мережі — тобто рівно те, що треба бачити під час аварії. З тієї ж причини переповнена черга викидає найстаріше, а не найновіше. Транспорт до сервера ще не під'єднано: серверний StreamLogs уже є, лишається цикл відправки на зонді. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
parent
59a8043ecf
commit
0bc2335661
3 changed files with 808 additions and 0 deletions
252
agent/internal/modules/syslog/parse.go
Normal file
252
agent/internal/modules/syslog/parse.go
Normal file
|
|
@ -0,0 +1,252 @@
|
|||
package syslog
|
||||
|
||||
import (
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Message — розібрана подія.
|
||||
//
|
||||
// Порожні поля — норма, а не помилка розбору: половина заліза шле щось
|
||||
// віддалено схоже на RFC3164, а комутатор із дешевої серії може почати
|
||||
// рядок одразу з тексту. Втратити таку подію гірше, ніж зберегти її з
|
||||
// самим лише повідомленням.
|
||||
type Message struct {
|
||||
Time time.Time
|
||||
Facility uint32
|
||||
Severity uint32
|
||||
Hostname string
|
||||
Tag string
|
||||
Message string
|
||||
// Structured data з RFC5424, зведені в пласкі ключі "id.параметр".
|
||||
Parsed map[string]string
|
||||
}
|
||||
|
||||
// Parse розбирає рядок syslog.
|
||||
//
|
||||
// Три формати в порядку спадання суворості: RFC5424, RFC3164 і «як є».
|
||||
// Останній не запасний варіант на випадок помилки, а окремий робочий
|
||||
// режим: пристрої, що шлють голий текст, існують, і їхні повідомлення
|
||||
// теж потрібні в журналі.
|
||||
//
|
||||
// now передається зовні, щоб розбір лишався передбачуваним у тестах:
|
||||
// RFC3164 не містить року, і його доводиться домислювати.
|
||||
func Parse(raw string, now time.Time) Message {
|
||||
m := Message{Time: now, Severity: 6, Facility: 1} // info / user
|
||||
|
||||
rest, pri, ok := cutPriority(raw)
|
||||
if ok {
|
||||
m.Facility = pri / 8
|
||||
m.Severity = pri % 8
|
||||
}
|
||||
m.Message = strings.TrimSpace(rest)
|
||||
|
||||
if parse5424(&m, rest) {
|
||||
return m
|
||||
}
|
||||
parse3164(&m, rest, now)
|
||||
return m
|
||||
}
|
||||
|
||||
// cutPriority знімає "<190>" з початку.
|
||||
func cutPriority(s string) (rest string, pri uint32, ok bool) {
|
||||
if len(s) < 3 || s[0] != '<' {
|
||||
return s, 0, false
|
||||
}
|
||||
end := strings.IndexByte(s, '>')
|
||||
// Пріоритет — щонайбільше три цифри; довше означає, що кутова дужка
|
||||
// просто трапилась у тексті.
|
||||
if end < 2 || end > 4 {
|
||||
return s, 0, false
|
||||
}
|
||||
n, err := strconv.Atoi(s[1:end])
|
||||
if err != nil || n < 0 || n > 191 {
|
||||
return s, 0, false
|
||||
}
|
||||
return s[end+1:], uint32(n), true
|
||||
}
|
||||
|
||||
// parse5424 розбирає "1 TIMESTAMP HOST APP PROCID MSGID [SD] MSG".
|
||||
func parse5424(m *Message, s string) bool {
|
||||
if !strings.HasPrefix(s, "1 ") {
|
||||
return false
|
||||
}
|
||||
f := strings.SplitN(s[2:], " ", 6)
|
||||
if len(f) < 6 {
|
||||
return false
|
||||
}
|
||||
|
||||
ts, host, app, procID, msgID, tail := f[0], f[1], f[2], f[3], f[4], f[5]
|
||||
|
||||
if t, err := time.Parse(time.RFC3339Nano, ts); err == nil {
|
||||
m.Time = t
|
||||
} else if ts != "-" {
|
||||
// Мітка часу є, але нечитабельна — це вже не 5424.
|
||||
return false
|
||||
}
|
||||
|
||||
m.Hostname = nilDash(host)
|
||||
m.Tag = nilDash(app)
|
||||
if p := nilDash(procID); p != "" {
|
||||
m.Tag += "[" + p + "]"
|
||||
}
|
||||
|
||||
sd, msg := cutStructured(tail)
|
||||
m.Parsed = sd
|
||||
if id := nilDash(msgID); id != "" {
|
||||
if m.Parsed == nil {
|
||||
m.Parsed = map[string]string{}
|
||||
}
|
||||
m.Parsed["msgid"] = id
|
||||
}
|
||||
m.Message = strings.TrimSpace(msg)
|
||||
return true
|
||||
}
|
||||
|
||||
// cutStructured знімає з початку блоки [id key="value" ...].
|
||||
func cutStructured(s string) (map[string]string, string) {
|
||||
s = strings.TrimLeft(s, " ")
|
||||
if strings.HasPrefix(s, "-") {
|
||||
return nil, strings.TrimPrefix(s, "-")
|
||||
}
|
||||
|
||||
out := map[string]string{}
|
||||
for strings.HasPrefix(s, "[") {
|
||||
end := findElementEnd(s)
|
||||
if end < 0 {
|
||||
// Незакрита дужка: усе, що лишилось, — текст.
|
||||
return orNil(out), s
|
||||
}
|
||||
parseElement(out, s[1:end])
|
||||
s = s[end+1:]
|
||||
}
|
||||
return orNil(out), s
|
||||
}
|
||||
|
||||
// findElementEnd шукає ']', яка закриває елемент, не плутаючись у
|
||||
// лапках: значення параметра цілком може містити дужку.
|
||||
func findElementEnd(s string) int {
|
||||
inQuotes := false
|
||||
for i := 1; i < len(s); i++ {
|
||||
switch s[i] {
|
||||
case '\\':
|
||||
i++
|
||||
case '"':
|
||||
inQuotes = !inQuotes
|
||||
case ']':
|
||||
if !inQuotes {
|
||||
return i
|
||||
}
|
||||
}
|
||||
}
|
||||
return -1
|
||||
}
|
||||
|
||||
func parseElement(out map[string]string, body string) {
|
||||
id, params, _ := strings.Cut(body, " ")
|
||||
id = strings.TrimSpace(id)
|
||||
if id == "" {
|
||||
return
|
||||
}
|
||||
for _, kv := range splitParams(params) {
|
||||
k, v, ok := strings.Cut(kv, "=")
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
v = strings.Trim(v, `"`)
|
||||
v = strings.ReplaceAll(v, `\"`, `"`)
|
||||
out[id+"."+strings.TrimSpace(k)] = v
|
||||
}
|
||||
}
|
||||
|
||||
// splitParams ділить 'a="1" b="2 3"' на пари, не розриваючи лапки.
|
||||
func splitParams(s string) []string {
|
||||
var (
|
||||
out []string
|
||||
start int
|
||||
inQuotes bool
|
||||
)
|
||||
for i := 0; i < len(s); i++ {
|
||||
switch s[i] {
|
||||
case '\\':
|
||||
i++
|
||||
case '"':
|
||||
inQuotes = !inQuotes
|
||||
case ' ':
|
||||
if !inQuotes {
|
||||
if i > start {
|
||||
out = append(out, s[start:i])
|
||||
}
|
||||
start = i + 1
|
||||
}
|
||||
}
|
||||
}
|
||||
if start < len(s) {
|
||||
out = append(out, s[start:])
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// parse3164 розбирає "MMM d hh:mm:ss HOST tag[pid]: msg".
|
||||
func parse3164(m *Message, s string, now time.Time) {
|
||||
s = strings.TrimLeft(s, " ")
|
||||
if len(s) < 16 {
|
||||
return
|
||||
}
|
||||
|
||||
// Рік у форматі не передбачений — беремо поточний. У ніч на перше
|
||||
// січня це дає майбутню дату для грудневих подій, тому відкочуємо
|
||||
// на рік назад: подія з майбутнього псує сортування журналу
|
||||
// набагато помітніше, ніж зсув на добу.
|
||||
t, err := time.ParseInLocation(time.Stamp, s[:15], now.Location())
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
t = t.AddDate(now.Year(), 0, 0)
|
||||
if t.Sub(now) > 24*time.Hour {
|
||||
t = t.AddDate(-1, 0, 0)
|
||||
}
|
||||
m.Time = t
|
||||
|
||||
rest := strings.TrimLeft(s[15:], " ")
|
||||
host, tail, ok := strings.Cut(rest, " ")
|
||||
if !ok {
|
||||
m.Message = rest
|
||||
return
|
||||
}
|
||||
m.Hostname = host
|
||||
|
||||
// Тег закінчується двокрапкою або пробілом — але лише якщо він
|
||||
// схожий на тег. Cisco шле "%SYS-5-CONFIG_I: ...", і відрізати це
|
||||
// в поле tag не можна: саме за цим текстом шукають зміну конфігу.
|
||||
if tag, msg, ok := strings.Cut(tail, ": "); ok && isTag(tag) {
|
||||
m.Tag = tag
|
||||
m.Message = strings.TrimSpace(msg)
|
||||
return
|
||||
}
|
||||
m.Message = strings.TrimSpace(tail)
|
||||
}
|
||||
|
||||
// isTag відсіює довгі й дивні «теги»: у RFC3164 це коротке ім'я
|
||||
// програми, а не половина повідомлення.
|
||||
func isTag(s string) bool {
|
||||
if s == "" || len(s) > 48 || strings.ContainsAny(s, " %") {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func nilDash(s string) string {
|
||||
if s == "-" {
|
||||
return ""
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
func orNil(m map[string]string) map[string]string {
|
||||
if len(m) == 0 {
|
||||
return nil
|
||||
}
|
||||
return m
|
||||
}
|
||||
287
agent/internal/modules/syslog/receiver.go
Normal file
287
agent/internal/modules/syslog/receiver.go
Normal file
|
|
@ -0,0 +1,287 @@
|
|||
// Package syslog — приймач подій syslog на зонді.
|
||||
//
|
||||
// Слухає UDP у мережі клієнта й тунелює події назовні. Приймач саме на
|
||||
// зонді, а не на сервері: комутатор у закритій мережі до сервера не
|
||||
// достукається, а зонд уже має вихідний канал і не потребує жодного
|
||||
// відкритого порту ззовні.
|
||||
//
|
||||
// Журнал потрібен не сам по собі. Порт, який фліпає раз на годину,
|
||||
// пінгом не видно взагалі — хост живий. У логах видно одразу. А подія
|
||||
// «%SYS-5-CONFIG_I» дає позачерговий бекап конфігу за секунди замість
|
||||
// середніх дванадцяти годин очікування нічного cron.
|
||||
package syslog
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
"net"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
|
||||
"google.golang.org/protobuf/types/known/timestamppb"
|
||||
)
|
||||
|
||||
// DefaultAddr — стандартний порт syslog.
|
||||
//
|
||||
// Нижче 1024, тож у Linux потрібна CAP_NET_BIND_SERVICE. Порт
|
||||
// налаштовується саме тому: віддати зонду право на привілейований порт
|
||||
// можна не всюди, а перенаправити 514 на 5514 правилом фаєрвола можна
|
||||
// скрізь.
|
||||
const DefaultAddr = ":514"
|
||||
|
||||
// maxDatagram — стеля на одну UDP-датаграму.
|
||||
//
|
||||
// RFC5424 дозволяє скільки завгодно, практика — 2 КіБ. Вісім із запасом
|
||||
// на багатослівні трапи Juniper і не більше: буфер виділяється на кожен
|
||||
// прийом.
|
||||
const maxDatagram = 8192
|
||||
|
||||
// Receiver приймає й накопичує події до відправки на сервер.
|
||||
type Receiver struct {
|
||||
addr string
|
||||
log *slog.Logger
|
||||
|
||||
mu sync.Mutex
|
||||
queue []*npv1.SyslogEntry
|
||||
dropped uint64
|
||||
buckets map[string]*bucket
|
||||
|
||||
maxQueue int
|
||||
|
||||
// Ліміти приходять від сервера в LogAck: він бачить картину по всіх
|
||||
// зондах і краще знає, що вважати шумом.
|
||||
minSeverity atomic.Int32
|
||||
perSource atomic.Int64
|
||||
|
||||
resolve atomic.Pointer[func(ip string) string]
|
||||
|
||||
notify chan struct{}
|
||||
}
|
||||
|
||||
// New створює приймач. Порожня адреса означає DefaultAddr.
|
||||
func New(addr string, log *slog.Logger) *Receiver {
|
||||
if addr == "" {
|
||||
addr = DefaultAddr
|
||||
}
|
||||
r := &Receiver{
|
||||
addr: addr,
|
||||
log: log,
|
||||
buckets: map[string]*bucket{},
|
||||
maxQueue: 20_000,
|
||||
notify: make(chan struct{}, 1),
|
||||
}
|
||||
// 7 — debug: типово беремо все. Фільтрувати вирішує сервер.
|
||||
r.minSeverity.Store(7)
|
||||
return r
|
||||
}
|
||||
|
||||
// SetResolver задає спосіб знайти пристрій за адресою джерела.
|
||||
func (r *Receiver) SetResolver(f func(ip string) string) {
|
||||
r.resolve.Store(&f)
|
||||
}
|
||||
|
||||
// ApplyAck застосовує ліміти, надіслані сервером.
|
||||
func (r *Receiver) ApplyAck(minSeverity, perSourcePerSec uint32) {
|
||||
if minSeverity > 0 && minSeverity <= 7 {
|
||||
r.minSeverity.Store(int32(minSeverity))
|
||||
}
|
||||
r.perSource.Store(int64(perSourcePerSec))
|
||||
}
|
||||
|
||||
// Ready повідомляє про появу подій у черзі.
|
||||
func (r *Receiver) Ready() <-chan struct{} { return r.notify }
|
||||
|
||||
// Drain забирає з черги до limit подій разом із лічильником відкинутих.
|
||||
func (r *Receiver) Drain(limit int) ([]*npv1.SyslogEntry, uint64) {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
|
||||
if limit <= 0 || limit > len(r.queue) {
|
||||
limit = len(r.queue)
|
||||
}
|
||||
if limit == 0 {
|
||||
return nil, 0
|
||||
}
|
||||
|
||||
out := r.queue[:limit]
|
||||
r.queue = append([]*npv1.SyslogEntry(nil), r.queue[limit:]...)
|
||||
dropped := r.dropped
|
||||
r.dropped = 0
|
||||
return out, dropped
|
||||
}
|
||||
|
||||
// Requeue повертає невідправлені події на початок черги.
|
||||
//
|
||||
// Порядок має значення: журнал читають хронологічно, і пачка, що
|
||||
// повернулась у хвіст, показала б аварію після її ж наслідків.
|
||||
func (r *Receiver) Requeue(entries []*npv1.SyslogEntry) {
|
||||
if len(entries) == 0 {
|
||||
return
|
||||
}
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
|
||||
room := r.maxQueue - len(r.queue)
|
||||
if room <= 0 {
|
||||
r.dropped += uint64(len(entries))
|
||||
return
|
||||
}
|
||||
if len(entries) > room {
|
||||
// Свіжі події важливіші за старі: під час аварії саме вони
|
||||
// пояснюють, що відбувається зараз.
|
||||
r.dropped += uint64(len(entries) - room)
|
||||
entries = entries[len(entries)-room:]
|
||||
}
|
||||
r.queue = append(entries, r.queue...)
|
||||
}
|
||||
|
||||
// Run слухає порт, доки живий контекст.
|
||||
func (r *Receiver) Run(ctx context.Context) error {
|
||||
pc, err := net.ListenPacket("udp", r.addr)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer pc.Close()
|
||||
|
||||
go func() {
|
||||
<-ctx.Done()
|
||||
// Закриття сокета — єдиний спосіб перервати ReadFrom, який
|
||||
// блокується без урахування контексту.
|
||||
_ = pc.Close()
|
||||
}()
|
||||
|
||||
r.log.Info("приймач syslog слухає", "адреса", r.addr)
|
||||
|
||||
buf := make([]byte, maxDatagram)
|
||||
for {
|
||||
n, src, err := pc.ReadFrom(buf)
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
r.Handle(string(buf[:n]), addrIP(src), time.Now())
|
||||
}
|
||||
}
|
||||
|
||||
// Handle розбирає й кладе в чергу одну подію.
|
||||
//
|
||||
// Окремо від Run, щоб приймач можна було перевірити без сокета.
|
||||
func (r *Receiver) Handle(raw, srcIP string, now time.Time) {
|
||||
// Одна датаграма може містити кілька рядків: так поводяться деякі
|
||||
// реалізації при сплеску. Розбираємо кожен окремо.
|
||||
for _, line := range strings.Split(strings.TrimRight(raw, "\x00\n"), "\n") {
|
||||
if strings.TrimSpace(line) == "" {
|
||||
continue
|
||||
}
|
||||
m := Parse(line, now)
|
||||
|
||||
if int32(m.Severity) > r.minSeverity.Load() {
|
||||
continue
|
||||
}
|
||||
if !r.allow(srcIP, now) {
|
||||
r.drop()
|
||||
continue
|
||||
}
|
||||
|
||||
e := &npv1.SyslogEntry{
|
||||
Ts: timestamppb.New(m.Time),
|
||||
SourceIp: srcIP,
|
||||
Facility: m.Facility,
|
||||
Severity: m.Severity,
|
||||
Hostname: m.Hostname,
|
||||
Tag: m.Tag,
|
||||
Message: m.Message,
|
||||
Parsed: m.Parsed,
|
||||
}
|
||||
if f := r.resolve.Load(); f != nil {
|
||||
e.DeviceId = (*f)(srcIP)
|
||||
}
|
||||
r.push(e)
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Receiver) push(e *npv1.SyslogEntry) {
|
||||
r.mu.Lock()
|
||||
if len(r.queue) >= r.maxQueue {
|
||||
// Викидаємо найстаріше, а не найновіше. Черга переповнюється
|
||||
// під час аварії, і саме свіжі рядки пояснюють, що зараз
|
||||
// відбувається.
|
||||
r.queue = r.queue[1:]
|
||||
r.dropped++
|
||||
}
|
||||
r.queue = append(r.queue, e)
|
||||
r.mu.Unlock()
|
||||
|
||||
select {
|
||||
case r.notify <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Receiver) drop() {
|
||||
r.mu.Lock()
|
||||
r.dropped++
|
||||
r.mu.Unlock()
|
||||
}
|
||||
|
||||
// --- обмеження частоти ------------------------------------------------
|
||||
|
||||
// bucket — відро токенів на одне джерело.
|
||||
type bucket struct {
|
||||
tokens float64
|
||||
last time.Time
|
||||
}
|
||||
|
||||
// allow пропускає подію, якщо джерело не перевищило ліміт.
|
||||
//
|
||||
// Ліміт на джерело, а не спільний: один комутатор, що зациклився на
|
||||
// повідомленні про помилку, інакше витіснив би з журналу всю решту
|
||||
// мережі — тобто рівно те, що потрібно бачити під час аварії.
|
||||
func (r *Receiver) allow(ip string, now time.Time) bool {
|
||||
rate := float64(r.perSource.Load())
|
||||
if rate <= 0 {
|
||||
return true
|
||||
}
|
||||
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
|
||||
b, ok := r.buckets[ip]
|
||||
if !ok {
|
||||
// Прибирання разом зі створенням: окремий прибиральник заради
|
||||
// мапи, яка росте на одне джерело, — зайва горутина.
|
||||
if len(r.buckets) > 4096 {
|
||||
r.buckets = map[string]*bucket{}
|
||||
}
|
||||
r.buckets[ip] = &bucket{tokens: rate - 1, last: now}
|
||||
return true
|
||||
}
|
||||
|
||||
b.tokens += now.Sub(b.last).Seconds() * rate
|
||||
if b.tokens > rate {
|
||||
b.tokens = rate
|
||||
}
|
||||
b.last = now
|
||||
|
||||
if b.tokens < 1 {
|
||||
return false
|
||||
}
|
||||
b.tokens--
|
||||
return true
|
||||
}
|
||||
|
||||
func addrIP(a net.Addr) string {
|
||||
if u, ok := a.(*net.UDPAddr); ok {
|
||||
return u.IP.String()
|
||||
}
|
||||
host, _, err := net.SplitHostPort(a.String())
|
||||
if err != nil {
|
||||
return a.String()
|
||||
}
|
||||
return host
|
||||
}
|
||||
269
agent/internal/modules/syslog/syslog_test.go
Normal file
269
agent/internal/modules/syslog/syslog_test.go
Normal file
|
|
@ -0,0 +1,269 @@
|
|||
package syslog
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
var now = time.Date(2026, 8, 25, 12, 0, 0, 0, time.UTC)
|
||||
|
||||
func TestParseRFC3164(t *testing.T) {
|
||||
m := Parse("<189>Aug 25 11:59:01 sw-core-01 mgd[1234]: інтерфейс піднявся", now)
|
||||
|
||||
if m.Facility != 23 || m.Severity != 5 {
|
||||
t.Fatalf("пріоритет: facility=%d severity=%d", m.Facility, m.Severity)
|
||||
}
|
||||
if m.Hostname != "sw-core-01" {
|
||||
t.Fatalf("хост: %q", m.Hostname)
|
||||
}
|
||||
if m.Tag != "mgd[1234]" {
|
||||
t.Fatalf("тег: %q", m.Tag)
|
||||
}
|
||||
if m.Message != "інтерфейс піднявся" {
|
||||
t.Fatalf("текст: %q", m.Message)
|
||||
}
|
||||
if !m.Time.Equal(time.Date(2026, 8, 25, 11, 59, 1, 0, time.UTC)) {
|
||||
t.Fatalf("час: %s", m.Time)
|
||||
}
|
||||
}
|
||||
|
||||
// Cisco шле "%SYS-5-CONFIG_I: ..." без тега. Відрізати цей текст у поле
|
||||
// tag не можна: саме за ним шукають зміну конфігу.
|
||||
func TestParseCiscoConfigChange(t *testing.T) {
|
||||
m := Parse("<189>Aug 25 11:59:01 rtr-01 %SYS-5-CONFIG_I: Configured from console by admin", now)
|
||||
|
||||
if m.Tag != "" {
|
||||
t.Fatalf("текст Cisco потрапив у тег: %q", m.Tag)
|
||||
}
|
||||
if m.Message != "%SYS-5-CONFIG_I: Configured from console by admin" {
|
||||
t.Fatalf("текст зіпсовано: %q", m.Message)
|
||||
}
|
||||
if m.Hostname != "rtr-01" {
|
||||
t.Fatalf("хост: %q", m.Hostname)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseRFC5424(t *testing.T) {
|
||||
raw := `<34>1 2026-08-25T11:58:00.123Z sw-01 sshd 4321 ID47 ` +
|
||||
`[exampleSDID@32473 iut="3" eventSource="Application"] спроба входу`
|
||||
m := Parse(raw, now)
|
||||
|
||||
if m.Facility != 4 || m.Severity != 2 {
|
||||
t.Fatalf("пріоритет: %d/%d", m.Facility, m.Severity)
|
||||
}
|
||||
if m.Hostname != "sw-01" || m.Tag != "sshd[4321]" {
|
||||
t.Fatalf("хост/тег: %q %q", m.Hostname, m.Tag)
|
||||
}
|
||||
if m.Message != "спроба входу" {
|
||||
t.Fatalf("текст: %q", m.Message)
|
||||
}
|
||||
if m.Parsed["exampleSDID@32473.iut"] != "3" {
|
||||
t.Fatalf("structured data: %v", m.Parsed)
|
||||
}
|
||||
if m.Parsed["msgid"] != "ID47" {
|
||||
t.Fatalf("msgid загубився: %v", m.Parsed)
|
||||
}
|
||||
if !m.Time.Equal(time.Date(2026, 8, 25, 11, 58, 0, 123000000, time.UTC)) {
|
||||
t.Fatalf("час: %s", m.Time)
|
||||
}
|
||||
}
|
||||
|
||||
// Пробіл усередині значення не має ділити параметри, а ']' у лапках —
|
||||
// закривати елемент.
|
||||
func TestParseStructuredQuoting(t *testing.T) {
|
||||
raw := `<34>1 2026-08-25T11:58:00Z h app - - [np x="a b" y="c]d"] текст`
|
||||
m := Parse(raw, now)
|
||||
|
||||
if m.Parsed["np.x"] != "a b" {
|
||||
t.Fatalf("значення з пробілом: %v", m.Parsed)
|
||||
}
|
||||
if m.Parsed["np.y"] != "c]d" {
|
||||
t.Fatalf("дужка в лапках: %v", m.Parsed)
|
||||
}
|
||||
if m.Message != "текст" {
|
||||
t.Fatalf("текст: %q", m.Message)
|
||||
}
|
||||
}
|
||||
|
||||
// Голий рядок без пріоритету теж має доїхати: такі пристрої існують, і
|
||||
// втратити подію гірше, ніж зберегти її з самим повідомленням.
|
||||
func TestParsePlainText(t *testing.T) {
|
||||
m := Parse("щось зламалось", now)
|
||||
|
||||
if m.Message != "щось зламалось" {
|
||||
t.Fatalf("текст: %q", m.Message)
|
||||
}
|
||||
if !m.Time.Equal(now) {
|
||||
t.Fatalf("час має бути часом прийому: %s", m.Time)
|
||||
}
|
||||
}
|
||||
|
||||
// Грудневі події, прийняті в січні, не мають опинятись у майбутньому.
|
||||
func TestParse3164YearRollover(t *testing.T) {
|
||||
jan := time.Date(2027, 1, 1, 0, 30, 0, 0, time.UTC)
|
||||
m := Parse("<13>Dec 31 23:59:00 h tag: пізно", jan)
|
||||
|
||||
if m.Time.Year() != 2026 {
|
||||
t.Fatalf("рік не відкотився: %s", m.Time)
|
||||
}
|
||||
if m.Time.After(jan) {
|
||||
t.Fatalf("подія з майбутнього: %s", m.Time)
|
||||
}
|
||||
}
|
||||
|
||||
func quiet() *slog.Logger {
|
||||
return slog.New(slog.NewTextHandler(io.Discard, nil))
|
||||
}
|
||||
|
||||
func TestReceiverQueueAndResolve(t *testing.T) {
|
||||
r := New(":0", quiet())
|
||||
r.SetResolver(func(ip string) string {
|
||||
if ip == "10.0.0.1" {
|
||||
return "dev-1"
|
||||
}
|
||||
return ""
|
||||
})
|
||||
|
||||
r.Handle("<13>Aug 25 11:00:00 h tag: перше", "10.0.0.1", now)
|
||||
r.Handle("<13>Aug 25 11:00:01 h tag: друге", "10.0.0.9", now)
|
||||
|
||||
got, dropped := r.Drain(0)
|
||||
if len(got) != 2 || dropped != 0 {
|
||||
t.Fatalf("черга: %d подій, %d відкинуто", len(got), dropped)
|
||||
}
|
||||
if got[0].DeviceId != "dev-1" {
|
||||
t.Fatalf("пристрій не зіставлено: %q", got[0].DeviceId)
|
||||
}
|
||||
if got[1].DeviceId != "" {
|
||||
t.Fatalf("чужа адреса зіставилась: %q", got[1].DeviceId)
|
||||
}
|
||||
if got[0].SourceIp != "10.0.0.1" {
|
||||
t.Fatalf("адреса джерела: %q", got[0].SourceIp)
|
||||
}
|
||||
}
|
||||
|
||||
// Одна датаграма з кількома рядками має дати кілька подій.
|
||||
func TestReceiverSplitsLines(t *testing.T) {
|
||||
r := New(":0", quiet())
|
||||
r.Handle("<13>рядок один\n<13>рядок два\n", "10.0.0.1", now)
|
||||
|
||||
got, _ := r.Drain(0)
|
||||
if len(got) != 2 {
|
||||
t.Fatalf("очікував 2 події, маю %d", len(got))
|
||||
}
|
||||
}
|
||||
|
||||
func TestReceiverSeverityFilter(t *testing.T) {
|
||||
r := New(":0", quiet())
|
||||
r.ApplyAck(4, 0) // не нижче warning
|
||||
|
||||
r.Handle("<191>Aug 25 11:00:00 h t: debug", "10.0.0.1", now) // severity 7
|
||||
r.Handle("<187>Aug 25 11:00:00 h t: помилка", "10.0.0.1", now)
|
||||
|
||||
got, _ := r.Drain(0)
|
||||
if len(got) != 1 {
|
||||
t.Fatalf("фільтр severity пропустив %d подій", len(got))
|
||||
}
|
||||
if got[0].Severity != 3 {
|
||||
t.Fatalf("лишилась не та подія: severity=%d", got[0].Severity)
|
||||
}
|
||||
}
|
||||
|
||||
// Джерело, що зациклилось, не має витіснити з журналу решту мережі.
|
||||
func TestReceiverRateLimitPerSource(t *testing.T) {
|
||||
r := New(":0", quiet())
|
||||
r.ApplyAck(7, 5)
|
||||
|
||||
for i := 0; i < 20; i++ {
|
||||
r.Handle(fmt.Sprintf("<13>шум %d", i), "10.0.0.1", now)
|
||||
}
|
||||
r.Handle("<13>важливе", "10.0.0.2", now)
|
||||
|
||||
got, dropped := r.Drain(0)
|
||||
if dropped == 0 {
|
||||
t.Fatal("ліміт не спрацював")
|
||||
}
|
||||
|
||||
var fromOther int
|
||||
for _, e := range got {
|
||||
if e.SourceIp == "10.0.0.2" {
|
||||
fromOther++
|
||||
}
|
||||
}
|
||||
if fromOther != 1 {
|
||||
t.Fatalf("подія з тихого джерела загубилась: %d", fromOther)
|
||||
}
|
||||
}
|
||||
|
||||
// Повернені події лягають на початок: журнал читають хронологічно.
|
||||
func TestReceiverRequeueKeepsOrder(t *testing.T) {
|
||||
r := New(":0", quiet())
|
||||
r.Handle("<13>третє", "10.0.0.1", now)
|
||||
|
||||
first, _ := r.Drain(0)
|
||||
r.Handle("<13>четверте", "10.0.0.1", now)
|
||||
r.Requeue(first)
|
||||
|
||||
got, _ := r.Drain(0)
|
||||
if len(got) != 2 {
|
||||
t.Fatalf("подій: %d", len(got))
|
||||
}
|
||||
if got[0].Message != "третє" {
|
||||
t.Fatalf("порядок порушено: %q перед %q", got[0].Message, got[1].Message)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReceiverListens(t *testing.T) {
|
||||
r := New("127.0.0.1:0", quiet())
|
||||
|
||||
// Порт 0 віддає ядро, тож адресу треба дізнатись після прив'язки —
|
||||
// заради цього єдиного тесту слухаємо руками.
|
||||
pc, err := net.ListenPacket("udp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatalf("сокет: %v", err)
|
||||
}
|
||||
addr := pc.LocalAddr().String()
|
||||
pc.Close()
|
||||
|
||||
r.addr = addr
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- r.Run(ctx) }()
|
||||
|
||||
conn, err := net.Dial("udp", addr)
|
||||
if err != nil {
|
||||
t.Fatalf("під'єднання: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
deadline := time.After(3 * time.Second)
|
||||
for {
|
||||
// Помилку запису ігноруємо свідомо: доки Run не встиг
|
||||
// прив'язатись, loopback відповідає ECONNREFUSED на UDP.
|
||||
_, _ = conn.Write([]byte("<13>Aug 25 11:00:00 h tag: живий"))
|
||||
select {
|
||||
case <-r.Ready():
|
||||
got, _ := r.Drain(0)
|
||||
if len(got) == 0 || got[0].Message != "живий" {
|
||||
t.Fatalf("прийнято не те: %+v", got)
|
||||
}
|
||||
cancel()
|
||||
if err := <-done; err != nil {
|
||||
t.Fatalf("Run: %v", err)
|
||||
}
|
||||
return
|
||||
case <-deadline:
|
||||
t.Fatal("подія не дійшла за три секунди")
|
||||
case <-time.After(50 * time.Millisecond):
|
||||
// UDP на loopback втрачає пакети рідко, але перший може
|
||||
// прийти до того, як сокет почав читати.
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue