From 0bc2335661693f6302243a9b877b25ff99822c95 Mon Sep 17 00:00:00 2001 From: byrsapty Date: Tue, 25 Aug 2026 13:31:40 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9F=D1=80=D0=B8=D0=B9=D0=BC=D0=B0=D1=87=20sy?= =?UTF-8?q?slog=20=D0=BD=D0=B0=20=D0=B7=D0=BE=D0=BD=D0=B4=D1=96:=20=D1=80?= =?UTF-8?q?=D0=BE=D0=B7=D0=B1=D1=96=D1=80=20RFC3164/5424,=20=D1=87=D0=B5?= =?UTF-8?q?=D1=80=D0=B3=D0=B0,=20=D0=BB=D1=96=D0=BC=D1=96=D1=82=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Модуль слухає UDP у мережі клієнта: комутатор у закритій мережі до сервера не достукається, а зонд уже має вихідний канал. Розбір трьома рівнями суворості — 5424, 3164 і «як є». Останній не запасний варіант, а робочий режим: пристрої, що шлють голий текст, існують, і втратити подію гірше, ніж зберегти саме повідомлення. Текст Cisco "%SYS-5-CONFIG_I: ..." свідомо НЕ ріжеться в поле tag — саме за ним шукають зміну конфігу. Ліміт частоти на джерело, а не спільний: комутатор, що зациклився на помилці, інакше витіснив би з журналу всю решту мережі — тобто рівно те, що треба бачити під час аварії. З тієї ж причини переповнена черга викидає найстаріше, а не найновіше. Транспорт до сервера ще не під'єднано: серверний StreamLogs уже є, лишається цикл відправки на зонді. Co-Authored-By: Claude Opus 5 --- agent/internal/modules/syslog/parse.go | 252 ++++++++++++++++ agent/internal/modules/syslog/receiver.go | 287 +++++++++++++++++++ agent/internal/modules/syslog/syslog_test.go | 269 +++++++++++++++++ 3 files changed, 808 insertions(+) create mode 100644 agent/internal/modules/syslog/parse.go create mode 100644 agent/internal/modules/syslog/receiver.go create mode 100644 agent/internal/modules/syslog/syslog_test.go diff --git a/agent/internal/modules/syslog/parse.go b/agent/internal/modules/syslog/parse.go new file mode 100644 index 0000000..9fbb6eb --- /dev/null +++ b/agent/internal/modules/syslog/parse.go @@ -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 +} diff --git a/agent/internal/modules/syslog/receiver.go b/agent/internal/modules/syslog/receiver.go new file mode 100644 index 0000000..9a3432f --- /dev/null +++ b/agent/internal/modules/syslog/receiver.go @@ -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 +} diff --git a/agent/internal/modules/syslog/syslog_test.go b/agent/internal/modules/syslog/syslog_test.go new file mode 100644 index 0000000..7ecd1e9 --- /dev/null +++ b/agent/internal/modules/syslog/syslog_test.go @@ -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 втрачає пакети рідко, але перший може + // прийти до того, як сокет почав читати. + } + } +}