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

555 lines
20 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"
"errors"
"time"
"github.com/jackc/pgx/v5"
)
// Виміри черг системи — сировина для сторінки «Черги».
//
// Питання, на яке відповідає ця сторінка, одне: чи все встигає, і якщо
// ні, то де саме. Тому знімок один на всі черги, а не ендпоїнт на кожну:
// відповідь, зібрана з семи запитів у різні секунди, суперечила б сама
// собі рівно тоді, коли щось ламається.
//
// Тут лише числа. Висновок із них робить queues_verdict.go — окремо,
// бо пороги — це продуктове рішення, яке хочеться читати й міняти, не
// перечитуючи SQL.
//
// Про предикат tenant_id у кожному запиті. Він лишається й після 0063,
// коли політики RLS справді почали діяти, — і не як перестраховка.
// Половина чисел нижче береться з гіпертаблиць, де політик немає й бути
// не може; а там, де вони є, предикат — це ще й те, завдяки чому план
// запиту лишається передбачуваним: політика додає умову, індекс
// використовує саме предикат.
// QueueFacts — усі виміри на один момент часу.
type QueueFacts struct {
At time.Time
WindowSec int
Jobs JobsQueueFacts
Commands CommandsQueueFacts
Outbox OutboxFacts
Agents []AgentQueueFacts
Checks ChecksFacts
Notify NotifyFacts
Pool PoolFacts
}
// JobsQueueFacts — черга збору конфігів (ncm.jobs).
type JobsQueueFacts struct {
Queued int
Running int
OldestQueuedSec int
OldestRunningSec int
// Stuck — ті, що вже перевищили десять хвилин і найближчим тіком
// прибиральника стануть невдалими (див. ReapStuckJobs).
Stuck int
Arrived int
Finished int
Failed int
// Reaped — скільки завдань за вікно вбив саме прибиральник. Це і є
// втрата: бекап, якого не сталося.
Reaped int
ByAgent []AgentBacklog
}
// AgentBacklog — скільки чекає на конкретний зонд. Потрібно, щоб
// відповісти не «черга росте», а «черга росте ось на цьому зонді».
type AgentBacklog struct {
AgentName string
Queued int
OldestSec int
}
// CommandsQueueFacts — масове виконання команд (ncm.command_targets).
//
// 'pending' тут не хвороба, а сам механізм: прогін навмисно тримає
// решту хостів, доки не звільниться слот паралельності. Тому глибина
// сама по собі ні про що не говорить, а от вік найстарішого 'queued' —
// говорить: його вже віддали диспетчеру, і він мав поїхати за секунди.
type CommandsQueueFacts struct {
ActiveRuns int
Concurrency int
Pending int
Queued int
Running int
OldestQueuedSec int
OldestPendingSec int
Finished int
// TimedOut — хости, які прибиральник закрив як «зонд не відповів».
// Команда на них не виконалась, і другої спроби не буде.
TimedOut int
}
// OutboxFacts — вихідна черга подій для WebSocket (core.event_outbox).
type OutboxFacts struct {
Unpublished int
// Capped — підрахунок уперся в стелю вибірки; справжнє число більше.
Capped bool
OldestSec int
Created int
}
// AgentQueueFacts — буфер і планувальник одного зонда.
//
// Величини з core.agents.health — це останній heartbeat. Вказівники, а
// не нулі: «зонд не доповів цього поля» і «зонд доповів нуль» — різні
// стани, і плутати їх на сторінці про втрати не можна.
type AgentQueueFacts struct {
ID string
Name string
Status string
Devices int
HeartbeatSec int
QueueDepth *int
Dropped *int64
TasksRunning *int
TasksQueued *int
// Тренд глибини буфера з ts.agent_health: те саме число, але
// півгодини тому й у піку. Одне поточне значення не відрізняє
// «стабільно 200» від «було 20, стало 200».
DepthFirst *int
DepthLast *int
DepthPeak *int
DepthPoints int
}
// ChecksFacts — планувальник зондів очима сервера.
//
// Пропущений такт зонд доповідає явно (STATE_SKIPPED), і код причини
// осідає в core.checks.last_error як «still_running: …» або
// «concurrency_limit: …». Це єдина втрата вимірювань, яку взагалі видно
// в базі, тому вона й рахується саме звідси.
type ChecksFacts struct {
Enabled int
Skipped int
// Stale — чек, від якого немає звіту довше трьох його інтервалів.
// Три, а не один: один пропуск буває від сплеску, три поспіль —
// це вже не сплеск.
Stale int
StaleSample []string
}
// NotifyFacts — доставка сповіщень (alr.notifications).
type NotifyFacts struct {
Sent int
Failed int
Throttled int
}
// PoolFacts — пул з'єднань до Postgres у ЦЬОМУ процесі.
//
// Тільки в цьому: у колектора власний пул у власному контейнері, і
// зазирнути в нього звідси нічим. Порожня черга запитів тут не означає,
// що з колектором усе гаразд.
type PoolFacts struct {
Acquired int32
Total int32
Max int32
EmptyAcquires int64
CanceledAcquires int64
AcquireWaitMs int64
}
// CollectQueueFacts знімає всі черги одним походом у базу.
//
// Одна транзакція на всі запити: інакше «надійшло» й «оброблено»
// рахувались би на різних станах бази, і на завантаженій системі
// сторінка сама собі показувала б від'ємний приріст.
func (s *Store) CollectQueueFacts(ctx context.Context, tenantID string,
window time.Duration) (*QueueFacts, error) {
if window <= 0 {
window = time.Hour
}
now := time.Now()
since := now.Add(-window)
f := &QueueFacts{At: now, WindowSec: int(window.Seconds())}
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
if err := collectJobs(ctx, tx, tenantID, since, &f.Jobs); err != nil {
return err
}
if err := collectCommands(ctx, tx, tenantID, since, &f.Commands); err != nil {
return err
}
if err := collectOutbox(ctx, tx, tenantID, since, &f.Outbox); err != nil {
return err
}
if err := collectChecks(ctx, tx, tenantID, &f.Checks); err != nil {
return err
}
if err := collectNotify(ctx, tx, tenantID, since, &f.Notify); err != nil {
return err
}
agents, err := collectAgents(ctx, tx, tenantID, now)
if err != nil {
return err
}
f.Agents = agents
return nil
})
if err != nil {
return nil, err
}
st := s.pool.Stat()
f.Pool = PoolFacts{
Acquired: st.AcquiredConns(),
Total: st.TotalConns(),
Max: st.MaxConns(),
EmptyAcquires: st.EmptyAcquireCount(),
CanceledAcquires: st.CanceledAcquireCount(),
AcquireWaitMs: st.AcquireDuration().Milliseconds(),
}
return f, nil
}
func collectJobs(ctx context.Context, tx pgx.Tx, tenantID string,
since time.Time, out *JobsQueueFacts) error {
// Хвіст, а не вся таблиця: умова збігається з ncm_jobs_pending_idx
// (частковий індекс по 'queued'/'running'), тож читається рівно те,
// що зараз у роботі, скільки б історії за ним не лежало.
rows, err := tx.Query(ctx, `
SELECT status::text,
count(*)::int,
COALESCE(extract(epoch FROM now() - min(created_at)), 0)::int,
count(*) FILTER (
WHERE status = 'running' AND started_at < now() - interval '10 minutes'
)::int
FROM ncm.jobs
WHERE tenant_id = $1 AND status IN ('queued','running')
GROUP BY status
`, tenantID)
if err != nil {
return err
}
for rows.Next() {
var status string
var n, oldest, stuck int
if err := rows.Scan(&status, &n, &oldest, &stuck); err != nil {
rows.Close()
return err
}
switch status {
case "queued":
out.Queued, out.OldestQueuedSec = n, oldest
case "running":
out.Running, out.OldestRunningSec, out.Stuck = n, oldest, stuck
}
}
rows.Close()
if err := rows.Err(); err != nil {
return err
}
// Прибиральника відрізняємо за тривалістю, а не за текстом помилки.
// Текст живе в ReapStuckJobs і може змінитись при першому ж
// редагуванні повідомлення — і лічильник втрат мовчки став би
// нулем. Тривалість же однозначна: власний таймаут завдання — п'ять
// хвилин, тож завдання, яке пробуло в 'running' довше десяти,
// закрив саме прибиральник.
if err := tx.QueryRow(ctx, `
SELECT count(*) FILTER (WHERE created_at > $2)::int,
count(*) FILTER (WHERE finished_at > $2)::int,
count(*) FILTER (WHERE finished_at > $2 AND status = 'failed')::int,
count(*) FILTER (
WHERE finished_at > $2 AND status = 'failed'
AND started_at IS NOT NULL
AND finished_at - started_at >= interval '10 minutes'
)::int
FROM ncm.jobs
WHERE tenant_id = $1 AND (created_at > $2 OR finished_at > $2)
`, tenantID, since).Scan(&out.Arrived, &out.Finished, &out.Failed, &out.Reaped); err != nil {
return err
}
brows, err := tx.Query(ctx, `
SELECT COALESCE(a.name, 'без зонда'),
count(*)::int,
COALESCE(extract(epoch FROM now() - min(j.created_at)), 0)::int
FROM ncm.jobs j
LEFT JOIN core.agents a ON a.id = j.agent_id
WHERE j.tenant_id = $1 AND j.status = 'queued'
GROUP BY a.name
ORDER BY 2 DESC
LIMIT 5
`, tenantID)
if err != nil {
return err
}
defer brows.Close()
for brows.Next() {
var b AgentBacklog
if err := brows.Scan(&b.AgentName, &b.Queued, &b.OldestSec); err != nil {
return err
}
out.ByAgent = append(out.ByAgent, b)
}
return brows.Err()
}
func collectCommands(ctx context.Context, tx pgx.Tx, tenantID string,
since time.Time, out *CommandsQueueFacts) error {
// Знову хвіст за частковим індексом (command_targets_pending_idx):
// прогін на п'ятсот хостів лишає по собі п'ятсот рядків історії, і
// рахувати їх щоразу, щоб дізнатись про п'ять активних, — рівно та
// помилка, від якої ця сторінка мала б стерегти.
rows, err := tx.Query(ctx, `
SELECT status::text,
count(*)::int,
COALESCE(extract(epoch FROM now() - min(created_at)), 0)::int
FROM ncm.command_targets
WHERE tenant_id = $1 AND status IN ('pending','queued','running')
GROUP BY status
`, tenantID)
if err != nil {
return err
}
for rows.Next() {
var status string
var n, oldest int
if err := rows.Scan(&status, &n, &oldest); err != nil {
rows.Close()
return err
}
switch status {
case "pending":
out.Pending, out.OldestPendingSec = n, oldest
case "queued":
out.Queued, out.OldestQueuedSec = n, oldest
case "running":
out.Running = n
}
}
rows.Close()
if err := rows.Err(); err != nil {
return err
}
if err := tx.QueryRow(ctx, `
SELECT count(*)::int, COALESCE(sum(concurrency), 0)::int
FROM ncm.command_runs
WHERE tenant_id = $1 AND status = 'running'
`, tenantID).Scan(&out.ActiveRuns, &out.Concurrency); err != nil {
return err
}
return tx.QueryRow(ctx, `
SELECT count(*)::int,
count(*) FILTER (
WHERE status = 'failed' AND error LIKE 'зонд не відповів%'
)::int
FROM ncm.command_targets
WHERE tenant_id = $1 AND finished_at > $2
`, tenantID, since).Scan(&out.Finished, &out.TimedOut)
}
func collectOutbox(ctx context.Context, tx pgx.Tx, tenantID string,
since time.Time, out *OutboxFacts) error {
// Стеля навмисна. Точне число недоставлених подій потрібне, лише
// поки воно мале; за кількадесят тисяч відповідь однаково одна —
// «публікація стоїть», і платити за неї повним підрахунком таблиці
// на кожне оновлення сторінки безглуздо.
const maxCount = 5000
if err := tx.QueryRow(ctx, `
SELECT count(*)::int FROM (
SELECT 1 FROM core.event_outbox
WHERE tenant_id = $1 AND published_at IS NULL
LIMIT $2
) x
`, tenantID, maxCount+1).Scan(&out.Unpublished); err != nil {
return err
}
if out.Unpublished > maxCount {
out.Unpublished, out.Capped = maxCount, true
}
// Найстаріша неопублікована — один рядок за тим самим частковим
// індексом, у порядку id. Саме її вік і є відповіддю: подія, що
// лежить тут хвилину, вже спізнилась на мапу.
var oldest *time.Time
if err := tx.QueryRow(ctx, `
SELECT created_at FROM core.event_outbox
WHERE tenant_id = $1 AND published_at IS NULL
ORDER BY id
LIMIT 1
`, tenantID).Scan(&oldest); err != nil && !errors.Is(err, pgx.ErrNoRows) {
return err
}
if oldest != nil {
out.OldestSec = int(time.Since(*oldest).Seconds())
}
return tx.QueryRow(ctx, `
SELECT count(*)::int FROM core.event_outbox
WHERE tenant_id = $1 AND created_at > $2
`, tenantID, since).Scan(&out.Created)
}
func collectChecks(ctx context.Context, tx pgx.Tx, tenantID string, out *ChecksFacts) error {
if err := tx.QueryRow(ctx, `
SELECT count(*)::int,
count(*) FILTER (
WHERE last_error LIKE 'still_running%'
OR last_error LIKE 'concurrency_limit%'
)::int,
count(*) FILTER (
WHERE last_run_at IS NULL
OR last_run_at < now() - make_interval(secs => interval_sec * 3)
)::int
FROM core.checks c
WHERE c.tenant_id = $1 AND c.enabled
-- Чеки видаленого хоста нікуди не діваються, а план їх уже
-- не бере — тобто вони «мовчать» назавжди. Рахувати їх як
-- пропущений такт означає тримати на сторінці вічну червону
-- картку, після якої люди перестають на неї дивитись.
AND EXISTS (
SELECT 1 FROM inv.devices d
WHERE d.id = c.device_id AND d.deleted_at IS NULL
)
`, tenantID).Scan(&out.Enabled, &out.Skipped, &out.Stale); err != nil {
return err
}
if out.Stale == 0 {
return nil
}
// Кілька прикладів, а не весь перелік: сторінка має назвати місце,
// куди йти дивитись, а не замінити собою розділ «Хости».
rows, err := tx.Query(ctx, `
SELECT d.name || ' · ' || c.check_type
FROM core.checks c
JOIN inv.devices d ON d.id = c.device_id
WHERE c.tenant_id = $1 AND c.enabled
AND d.deleted_at IS NULL
AND (c.last_run_at IS NULL
OR c.last_run_at < now() - make_interval(secs => c.interval_sec * 3))
ORDER BY c.last_run_at NULLS FIRST
LIMIT 5
`, tenantID)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var s string
if err := rows.Scan(&s); err != nil {
return err
}
out.StaleSample = append(out.StaleSample, s)
}
return rows.Err()
}
func collectNotify(ctx context.Context, tx pgx.Tx, tenantID string,
since time.Time, out *NotifyFacts) error {
// alr.notifications — гіпертаблиця, RLS на ній немає навмисно
// (несумісний зі стисненням, див. 0011), тому tenant_id тут не
// перестраховка, а єдиний бар'єр.
return tx.QueryRow(ctx, `
SELECT count(*) FILTER (WHERE status = 'sent')::int,
count(*) FILTER (WHERE status = 'failed')::int,
count(*) FILTER (WHERE status = 'throttled')::int
FROM alr.notifications
WHERE tenant_id = $1 AND ts > $2
`, tenantID, since).Scan(&out.Sent, &out.Failed, &out.Throttled)
}
func collectAgents(ctx context.Context, tx pgx.Tx, tenantID string,
now time.Time) ([]AgentQueueFacts, error) {
rows, err := tx.Query(ctx, `
SELECT a.id::text, a.name, a.status::text, a.last_heartbeat_at,
(a.health->>'queue_depth')::int,
(a.health->>'dropped_samples')::bigint,
(a.health->>'tasks_running')::int,
(a.health->>'tasks_queued')::int,
(SELECT count(*)::int FROM inv.devices d
WHERE d.agent_id = a.id AND d.deleted_at IS NULL)
FROM core.agents a
WHERE a.tenant_id = $1
ORDER BY a.name
`, tenantID)
if err != nil {
return nil, err
}
var out []AgentQueueFacts
byID := map[string]int{}
for rows.Next() {
var a AgentQueueFacts
var hb *time.Time
if err := rows.Scan(&a.ID, &a.Name, &a.Status, &hb,
&a.QueueDepth, &a.Dropped, &a.TasksRunning, &a.TasksQueued,
&a.Devices); err != nil {
rows.Close()
return nil, err
}
a.HeartbeatSec = -1
if hb != nil {
a.HeartbeatSec = int(now.Sub(*hb).Seconds())
}
byID[a.ID] = len(out)
out = append(out, a)
}
rows.Close()
if err := rows.Err(); err != nil {
return nil, err
}
if len(out) == 0 {
return out, nil
}
// Тренд глибини буфера — з ts.agent_health. Вікно коротке навмисно:
// heartbeat ходить раз на півхвилини, і півгодини дають шістдесят
// точок на зонда — досить, щоб побачити напрямок, і замало, щоб
// прочитати зайвий чанк гіпертаблиці.
trows, err := tx.Query(ctx, `
SELECT agent_id::text,
(array_agg(COALESCE(queue_depth,0) ORDER BY ts))[1],
(array_agg(COALESCE(queue_depth,0) ORDER BY ts DESC))[1],
max(COALESCE(queue_depth,0)),
count(*)::int
FROM ts.agent_health
WHERE tenant_id = $1 AND ts > now() - interval '30 minutes'
GROUP BY agent_id
`, tenantID)
if err != nil {
return nil, err
}
defer trows.Close()
for trows.Next() {
var id string
var first, last, peak, points int
if err := trows.Scan(&id, &first, &last, &peak, &points); err != nil {
return nil, err
}
i, ok := byID[id]
if !ok {
continue
}
out[i].DepthFirst, out[i].DepthLast = &first, &last
out[i].DepthPeak, out[i].DepthPoints = &peak, points
}
return out, trows.Err()
}