Netpulse_SasS/server/internal/store/events.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

115 lines
4 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"
"time"
"github.com/jackc/pgx/v5"
)
// Подія для UI. Пишеться в core.event_outbox тією ж транзакцією, що й
// зміна, яку описує, — інакше WebSocket міг би розповісти про перехід,
// якого в базі ще (або вже) немає.
type Event struct {
ID int64 `json:"id"`
TenantID string `json:"-"`
Topic string `json:"topic"`
Payload json.RawMessage `json:"payload"`
At time.Time `json:"at"`
}
// FetchEvents читає нові події після afterID.
//
// Транспорт навмисно простий: опитування таблиці замість LISTEN/NOTIFY.
// NOTIFY не переживає падіння підписника й обмежений 8 КБ на
// повідомлення, а тут потрібна гарантія, що жодна зміна статусу не
// загубиться між перезапусками API. Ціна — один дешевий запит за
// індексом раз на секунду.
func (s *Store) FetchEvents(ctx context.Context, afterID int64, limit int) ([]Event, error) {
if limit <= 0 || limit > 1000 {
limit = 500
}
rows, err := s.bg.Query(ctx, `
SELECT id, tenant_id::text, topic, payload::text, created_at
FROM core.event_outbox
WHERE id > $1
ORDER BY id
LIMIT $2
`, afterID, limit)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Event
for rows.Next() {
var e Event
var payload string
if err := rows.Scan(&e.ID, &e.TenantID, &e.Topic, &payload, &e.At); err != nil {
return nil, err
}
e.Payload = json.RawMessage(payload)
out = append(out, e)
}
return out, rows.Err()
}
// LatestEventID — з якого місця починати читати новому підписнику.
//
// Нова сесія WebSocket не має отримувати вчорашні події: клієнт щойно
// завантажив повний стан мапи, і все старіше в ньому вже враховано.
func (s *Store) LatestEventID(ctx context.Context) (int64, error) {
var id *int64
if err := s.bg.QueryRow(ctx,
`SELECT max(id) FROM core.event_outbox`).Scan(&id); err != nil {
return 0, err
}
if id == nil {
return 0, nil
}
return *id, nil
}
// MarkEventsPublished позначає події доставленими.
//
// Позначка не керує доставкою (підписники йдуть за id), вона потрібна
// прибиральнику: невідправлені події видаляти не можна, а відправлені —
// можна, і без цього поля таблиця росла б вічно.
func (s *Store) MarkEventsPublished(ctx context.Context, throughID int64) error {
_, err := s.bg.Exec(ctx, `
UPDATE core.event_outbox
SET published_at = now()
WHERE id <= $1 AND published_at IS NULL
`, throughID)
return err
}
// PruneEvents видаляє доставлені події, старші за вказаний вік.
func (s *Store) PruneEvents(ctx context.Context, olderThan time.Duration) (int64, error) {
tag, err := s.bg.Exec(ctx, `
DELETE FROM core.event_outbox
WHERE published_at IS NOT NULL AND created_at < now() - $1::interval
`, olderThan.String())
if err != nil {
return 0, err
}
return tag.RowsAffected(), nil
}
// PublishEvent кладе подію в outbox поза чужою транзакцією.
// Для змін, що вже мають власну транзакцію, подія пишеться прямо в ній.
func (s *Store) PublishEvent(ctx context.Context, tenantID, topic string, payload any) error {
data, err := json.Marshal(payload)
if err != nil {
return err
}
return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
_, err := tx.Exec(ctx, `
INSERT INTO core.event_outbox (tenant_id, topic, payload)
VALUES ($1, $2, $3::jsonb)
`, tenantID, topic, string(data))
return err
})
}