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

1057 lines
44 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"
"errors"
"fmt"
"strings"
"time"
"github.com/jackc/pgx/v5"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"github.com/netpulse/netpulse/server/internal/crypto"
"google.golang.org/protobuf/types/known/durationpb"
)
// Масове виконання команд по фільтру хостів.
//
// Шлях до пристрою тут той самий, що й у збору конфігів: сервер кладе
// завдання в базу, диспетчер колектора віддає його живій сесії зонда,
// зонд відкриває CLI й повертає вивід потоком чанків. Нового транспорту
// не з'явилось — з'явилась інша одиниця роботи (прогін, а не бекап) і
// власна таблиця під неї (див. міграцію 0036).
//
// Ознака, за якою зонд і приймач розрізняють два види завдань, —
// config_type у ConfigJob. Поле вже було й уже возило вид зрізу
// ('running', 'startup'); CommandConfigType просто стає ще одним його
// значенням. Так вдалося не чіпати .proto взагалі.
// CommandConfigType — значення config_type, яким позначено завдання
// «виконати команди», а не «зняти конфіг».
const CommandConfigType = "command"
// Стелі. Не налаштовуються навмисно: це не параметри роботи, а межі,
// за якими помилка оператора стає аварією мережі.
const (
// Двісті пристроїв — типова дільниця. Тисяча — це вже не «прогнати
// команду», а міграція, і робити її треба свідомо, кількома
// прогонами, дивлячись на результат кожного.
MaxRunDevices = 500
// Довший список команд — це вже скрипт, і йому місце в шаблоні
// профілю, а не в одноразовому вікні.
MaxRunCommands = 20
MaxCommandLength = 512
)
// CommandOutcome — вивід однієї команди на одному хості.
//
// Prep позначає підготовчі команди профілю (screen-length, terminal
// length): їх не набирала людина, і в переліку результатів вони мають
// бути видні окремо — інакше оператор шукає свою команду серед чужих.
type CommandOutcome struct {
Command string `json:"command"`
Output string `json:"output"`
Error string `json:"error,omitempty"`
Prep bool `json:"prep,omitempty"`
}
// DeviceFilter — той самий набір понять, яким оператор фільтрує список
// хостів у розділі «Хости».
//
// Другої мови фільтрів у продукті бути не повинно: людина, яка щойно
// відібрала в інвентарі «усі Huawei на Миронівці», має відібрати те
// саме тут тими самими полями. Тому майданчик тут — це його НАЗВА, а не
// ідентифікатор: саме назву видно в переліку хостів, і саме її людина
// впізнає в підтвердженні прогону.
type DeviceFilter struct {
Query string `json:"query,omitempty"`
GroupIDs []string `json:"group_ids,omitempty"`
Vendors []string `json:"vendors,omitempty"`
Kinds []string `json:"kinds,omitempty"`
Sites []string `json:"sites,omitempty"`
Statuses []string `json:"statuses,omitempty"`
// Models — точним збігом зі списку, як виробник.
//
// Значення приїздять із sysDescr при розпізнаванні й виглядають як
// «DGS-3420-28SC» або «ex4600-40f»: людина обирає їх із переліку
// наявних, а не набирає, тож підрядок тут не потрібен.
Models []string `json:"models,omitempty"`
// Версія ПЗ — ОДНА умова, а не перелік значень.
//
// Питання до версії ставлять інакше, ніж до виробника. «Huawei або
// D-Link» — звичайне питання; «версія 5.70 або V1.2.5P3» — ні. Зате
// «усе, що НЕ 14.1X53-D27.3» — це буквально те, з чого починається
// планування оновлення прошивки, і переліком значень воно не
// виражається взагалі. Тому операція, а не набір.
VersionOp string `json:"version_op,omitempty"`
VersionValue string `json:"version_value,omitempty"`
// OnlyEnabled типово true: вимкнений хост вимкнули свідомо, і
// масова команда не привід про це забути.
OnlyEnabled bool `json:"only_enabled"`
}
// Операції над версією ПЗ.
//
// Рівно ті, на які є питання перед оновленням прошивки, і жодної зайвої:
//
// eq — «покажи всі на 14.1X53-D27.3» (перевірка після оновлення);
// ne — «покажи все, що не на 14.1X53-D27.3» (що лишилось оновити);
// contains — «покажи все сімейство 14.1X53» (версії всередині релізу
// відрізняються суфіксом, і точний збіг тут не працює);
// ncontains — те саме питання «що лишилось» на рівні сімейства;
// empty — «покажи, у кого версія невідома».
//
// Останнє не косметика й не симетрія заради симетрії. `ne` навмисно
// ЗАХОПЛЮЄ хости з порожньою версією — вони справді не на цільовій
// прошивці, і мовчки викидати їх зі списку «що лишилось оновити» було б
// найгіршим із можливих варіантів. Але тоді потрібен спосіб подивитись
// саме на них: D-Link версії в sysDescr не повідомляє взагалі, і пів
// дільниці може виявитись «невідомо».
const (
VersionOpEq = "eq"
VersionOpNe = "ne"
VersionOpContains = "contains"
VersionOpNContains = "ncontains"
VersionOpEmpty = "empty"
)
// versionCond зводить умову до пари, яку розуміє запит.
//
// Невідома операція й порожнє значення однаково означають «умови немає».
// Не помилка: фільтр складають на льоту, і поле версії, у якому людина
// стерла текст, має повернути повний перелік, а не відмову.
func versionCond(f DeviceFilter) (op, val string) {
val = strings.TrimSpace(f.VersionValue)
switch f.VersionOp {
case VersionOpEmpty:
// Єдина операція, якій значення не потрібне.
return VersionOpEmpty, ""
case VersionOpEq, VersionOpNe, VersionOpContains, VersionOpNContains:
if val == "" {
return "", ""
}
return f.VersionOp, val
}
return "", ""
}
// CommandCandidate — хост, що потрапив під фільтр.
//
// Разом із причиною, чому він до команди не готовий: перелік показується
// ДО запуску, і дізнатись про «немає доступу ssh» через хвилину після
// натискання — це дізнатись запізно.
type CommandCandidate struct {
DeviceID string `json:"device_id"`
Name string `json:"name"`
Address string `json:"address,omitempty"`
Kind string `json:"kind"`
Vendor string `json:"vendor,omitempty"`
// Модель і версія їдуть у перелік завжди, а не лише коли за ними
// фільтрують: питання «а точно всі вони на старій прошивці» ставлять
// саме над цим переліком, і відповідь має бути в тому ж рядку.
Model string `json:"model,omitempty"`
OSVersion string `json:"os_version,omitempty"`
SiteName string `json:"site_name,omitempty"`
Status string `json:"status"`
Enabled bool `json:"enabled"`
AgentOnline bool `json:"agent_online"`
Ready bool `json:"ready"`
Reason string `json:"reason,omitempty"`
agentID string
}
// ResolveCommandTargets віддає хости, що збігаються з фільтром.
//
// Фільтрація в самому запиті, а не після вибірки: перелік стає текстом
// підтвердження, і зайвий хост у ньому — це зайвий хост у прогоні.
//
// Права на запис перевіряються тут же. Виконати команду на хості, який
// людині видно лише на читання, не можна: це зміна стану заліза, і
// «бачу» тут не означає «можу».
func (s *Store) ResolveCommandTargets(ctx context.Context, tenantID string,
sc Scope, f DeviceFilter) ([]CommandCandidate, error) {
var out []CommandCandidate
// Умова фільтра — спільна вставка (device_filter.go), а не власна
// копія: поле, що з'явиться у фільтрі, має відібрати те саме й тут,
// і в масовій правці хостів, і на сторінці конфігів.
cond, condArgs := deviceFilterSQL(f, 4)
args := append([]any{tenantID, sc.Unrestricted, nonNilIDs(sc.Writable)}, condArgs...)
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
rows, err := tx.Query(ctx, `
SELECT d.id::text, d.name, COALESCE(host(d.address),''), d.kind::text,
COALESCE(d.vendor,''), COALESCE(d.model,''), COALESCE(d.os_version,''),
COALESCE(st.name,''),
d.status::text, d.enabled,
COALESCE(d.agent_id::text,''),
COALESCE(a.status::text = 'online', false),
EXISTS (
SELECT 1 FROM inv.device_credentials dc
JOIN inv.credentials c ON c.id = dc.credential_id
WHERE dc.device_id = d.id AND c.tenant_id = d.tenant_id
AND c.proto IN ('ssh','telnet')
)
FROM inv.devices d
LEFT JOIN inv.sites st ON st.id = d.site_id
LEFT JOIN core.agents a ON a.id = d.agent_id
WHERE d.tenant_id = $1 AND d.deleted_at IS NULL
AND ($2::boolean OR d.id = ANY($3::uuid[]))`+cond+`
ORDER BY d.name
`, args...)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var c CommandCandidate
var hasCLI bool
if err := rows.Scan(&c.DeviceID, &c.Name, &c.Address, &c.Kind, &c.Vendor,
&c.Model, &c.OSVersion,
&c.SiteName, &c.Status, &c.Enabled, &c.agentID, &c.AgentOnline,
&hasCLI); err != nil {
return err
}
c.Ready, c.Reason = commandReadiness(c.agentID, c.AgentOnline, hasCLI)
out = append(out, c)
}
return rows.Err()
})
return out, err
}
// commandReadiness пояснює, чому хост не візьме команду.
//
// Текстом, а не кодом: причин небагато, кожна лікується по-своєму, і
// тримати той самий перелік ще й у вебі означало б забути оновити одне
// з двох місць. Так само зроблено для «Розпізнати зараз».
func commandReadiness(agentID string, online, hasCLI bool) (bool, string) {
switch {
case agentID == "":
return false, "не прив'язаний до зонда — нікому виконати"
case !hasCLI:
return false, "немає доступу ssh або telnet"
case !online:
// Не відмова: завдання почекає в черзі. Але людина має бачити,
// що результату зараз не буде.
return true, "зонд не на зв'язку — виконається, коли повернеться"
}
return true, ""
}
// ---------------------------------------------------------------------
// Створення прогону
// ---------------------------------------------------------------------
// CommandRunInput — те, що приходить із форми.
type CommandRunInput struct {
Commands []string `json:"commands"`
Filter DeviceFilter `json:"filter"`
// DeviceIDs — перелік, який людина БАЧИЛА в підтвердженні.
//
// Не надлишковий: між переглядом і натисканням хтось міг завести
// новий хост, що теж підпадає під фільтр. Прогін іде по перетину
// фільтра й цього переліку — тобто рівно по тому, що було на екрані.
DeviceIDs []string `json:"device_ids"`
Concurrency int `json:"concurrency"`
TimeoutSec int `json:"timeout_sec"`
}
// CommandRun — прогін у переліку.
type CommandRun struct {
ID string `json:"id"`
Status string `json:"status"`
Commands []string `json:"commands"`
Filter any `json:"filter,omitempty"`
Concurrency int `json:"concurrency"`
TimeoutSec int `json:"timeout_sec"`
Total int `json:"total"`
CreatedBy string `json:"created_by,omitempty"`
CreatedAt time.Time `json:"created_at"`
FinishedAt *time.Time `json:"finished_at,omitempty"`
// Зведення по станах: скільки успішних, скільки в роботі.
Counts map[string]int `json:"counts"`
}
// CommandTarget — хост у прогоні разом із результатом.
type CommandTarget struct {
ID string `json:"id"`
DeviceID string `json:"device_id"`
DeviceName string `json:"device_name"`
Address string `json:"address,omitempty"`
Status string `json:"status"`
Outcomes []CommandOutcome `json:"outcomes"`
Error string `json:"error,omitempty"`
Transcript string `json:"transcript,omitempty"`
StartedAt *time.Time `json:"started_at,omitempty"`
FinishedAt *time.Time `json:"finished_at,omitempty"`
DurationMs int `json:"duration_ms"`
}
// CommandRunDetail — прогін із хостами.
type CommandRunDetail struct {
CommandRun
Targets []CommandTarget `json:"targets"`
}
// CreateCommandRun заводить прогін і його хости.
//
// Одна транзакція на все: прогін без хостів і хости без прогону однаково
// безглузді, а диспетчер починає забирати рядки з наступного такту —
// тобто вже за кілька секунд.
func (s *Store) CreateCommandRun(ctx context.Context, tenantID, userID string,
in CommandRunInput) (CommandRunDetail, error) {
var out CommandRunDetail
cmds, err := sanitizeCommands(in.Commands)
if err != nil {
return out, err
}
if len(in.DeviceIDs) == 0 {
return out, fmt.Errorf("%w: не обрано жодного хоста", ErrInvalid)
}
if len(in.DeviceIDs) > MaxRunDevices {
return out, fmt.Errorf("%w: за раз можна взяти не більше %d хостів",
ErrInvalid, MaxRunDevices)
}
concurrency := in.Concurrency
if concurrency <= 0 {
concurrency = 5
}
if concurrency > 50 {
concurrency = 50
}
timeout := in.TimeoutSec
if timeout <= 0 {
timeout = 120
}
if timeout > 600 {
timeout = 600
}
filterJSON, err := json.Marshal(in.Filter)
if err != nil {
return out, err
}
cmdJSON, err := json.Marshal(cmds)
if err != nil {
return out, err
}
err = s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
var runID string
if err := tx.QueryRow(ctx, `
INSERT INTO ncm.command_runs
(tenant_id, created_by, commands, filter, concurrency, timeout_sec, total)
VALUES ($1, $2, $3::jsonb, $4::jsonb, $5, $6, $7)
RETURNING id::text
`, tenantID, nullUUID(userID), string(cmdJSON), string(filterJSON),
concurrency, timeout, len(in.DeviceIDs)).Scan(&runID); err != nil {
return err
}
// Хости беремо тим самим запитом, що й перелік: інакше між
// підтвердженням і вставкою міг би пролізти рядок, якого людина
// не бачила. ANY($2) звужує до підтвердженого переліку.
//
// Хост без зонда лягає одразу невдалим, а не пропускається:
// прогін на 40 хостів, у якому мовчки виконалось 38, — це
// прогін, після якого ніхто не знає, що сталося з двома.
rows, err := tx.Query(ctx, `
INSERT INTO ncm.command_targets
(run_id, tenant_id, device_id, agent_id, status, error)
SELECT $1, $2, d.id, d.agent_id,
CASE WHEN d.agent_id IS NULL THEN 'failed'::ncm.command_target_status
ELSE 'pending'::ncm.command_target_status END,
CASE WHEN d.agent_id IS NULL
THEN 'хост не прив''язаний до зонда — нікому виконати'
ELSE NULL END
FROM inv.devices d
WHERE d.tenant_id = $2 AND d.deleted_at IS NULL
AND d.id = ANY($3::uuid[])
ORDER BY d.name
RETURNING id::text, device_id::text
`, runID, tenantID, nonNilIDs(in.DeviceIDs))
if err != nil {
return err
}
var inserted int
for rows.Next() {
var id, dev string
if err := rows.Scan(&id, &dev); err != nil {
rows.Close()
return err
}
inserted++
}
rows.Close()
if err := rows.Err(); err != nil {
return err
}
if inserted == 0 {
return fmt.Errorf("%w: жоден із обраних хостів більше не існує", ErrInvalid)
}
if _, err := tx.Exec(ctx,
`UPDATE ncm.command_runs SET total = $2 WHERE id = $1`, runID, inserted); err != nil {
return err
}
out.ID = runID
return nil
})
if err != nil {
return out, err
}
return s.GetCommandRun(ctx, tenantID, out.ID, true)
}
// sanitizeCommands звіряє команди з межами й прибирає порожні рядки.
//
// Переносу рядка тут бути не може: одне поле форми — одна команда, а
// склеєний рядок поїхав би на пристрій як дві, з яких другу ніхто не
// бачив у підтвердженні.
func sanitizeCommands(in []string) ([]string, error) {
out := make([]string, 0, len(in))
for _, c := range in {
c = strings.TrimSpace(c)
if c == "" {
continue
}
if strings.ContainsAny(c, "\r\n") {
return nil, fmt.Errorf("%w: команда не може містити перенесення рядка", ErrInvalid)
}
if len(c) > MaxCommandLength {
return nil, fmt.Errorf("%w: команда довша за %d символів",
ErrInvalid, MaxCommandLength)
}
out = append(out, c)
}
if len(out) == 0 {
return nil, fmt.Errorf("%w: не задано жодної команди", ErrInvalid)
}
if len(out) > MaxRunCommands {
return nil, fmt.Errorf("%w: за раз можна виконати не більше %d команд",
ErrInvalid, MaxRunCommands)
}
return out, nil
}
// ---------------------------------------------------------------------
// Читання
// ---------------------------------------------------------------------
// ListCommandRuns — історія прогонів.
func (s *Store) ListCommandRuns(ctx context.Context, tenantID string, limit int) ([]CommandRun, error) {
if limit <= 0 || limit > 200 {
limit = 50
}
var out []CommandRun
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
rows, err := tx.Query(ctx, `
SELECT r.id::text, r.status::text, r.commands::text,
r.concurrency, r.timeout_sec, r.total,
COALESCE(u.username, ''), r.created_at, r.finished_at,
COALESCE((
SELECT jsonb_object_agg(x.status, x.n)::text
FROM (
SELECT t.status::text AS status, count(*) AS n
FROM ncm.command_targets t
WHERE t.run_id = r.id
GROUP BY t.status
) x
), '{}')
FROM ncm.command_runs r
LEFT JOIN core.users u ON u.id = r.created_by
WHERE r.tenant_id = $1
ORDER BY r.created_at DESC
LIMIT $2
`, tenantID, limit)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var r CommandRun
var cmdRaw, countsRaw string
if err := rows.Scan(&r.ID, &r.Status, &cmdRaw, &r.Concurrency, &r.TimeoutSec,
&r.Total, &r.CreatedBy, &r.CreatedAt, &r.FinishedAt, &countsRaw); err != nil {
return err
}
_ = jsonUnmarshalStrings(cmdRaw, &r.Commands)
r.Counts = map[string]int{}
_ = json.Unmarshal([]byte(countsRaw), &r.Counts)
out = append(out, r)
}
return rows.Err()
})
return out, err
}
// GetCommandRun — прогін із хостами.
//
// withTranscript вимикає стенограму в живому опитуванні: сторінка
// перечитує прогін раз на дві секунди, а стенограма сесії до великого
// шасі — це сотні кілобайт на хост.
func (s *Store) GetCommandRun(ctx context.Context, tenantID, runID string,
withTranscript bool) (CommandRunDetail, error) {
var out CommandRunDetail
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
var cmdRaw, filterRaw string
err := tx.QueryRow(ctx, `
SELECT r.id::text, r.status::text, r.commands::text, r.filter::text,
r.concurrency, r.timeout_sec, r.total,
COALESCE(u.username, ''), r.created_at, r.finished_at
FROM ncm.command_runs r
LEFT JOIN core.users u ON u.id = r.created_by
WHERE r.id = $1 AND r.tenant_id = $2
`, runID, tenantID).Scan(&out.ID, &out.Status, &cmdRaw, &filterRaw,
&out.Concurrency, &out.TimeoutSec, &out.Total,
&out.CreatedBy, &out.CreatedAt, &out.FinishedAt)
if errors.Is(err, pgx.ErrNoRows) {
return ErrNotFound
}
if err != nil {
return err
}
_ = jsonUnmarshalStrings(cmdRaw, &out.Commands)
var filter DeviceFilter
if json.Unmarshal([]byte(filterRaw), &filter) == nil {
out.Filter = filter
}
rows, err := tx.Query(ctx, `
SELECT t.id::text, t.device_id::text, d.name, COALESCE(host(d.address),''),
t.status::text, t.outcomes::text, COALESCE(t.error,''),
CASE WHEN $3::boolean THEN COALESCE(t.transcript,'') ELSE '' END,
t.started_at, t.finished_at, COALESCE(t.duration_ms, 0)
FROM ncm.command_targets t
JOIN inv.devices d ON d.id = t.device_id
WHERE t.run_id = $1 AND t.tenant_id = $2
ORDER BY d.name
`, runID, tenantID, withTranscript)
if err != nil {
return err
}
defer rows.Close()
out.Counts = map[string]int{}
for rows.Next() {
var t CommandTarget
var outcomesRaw string
if err := rows.Scan(&t.ID, &t.DeviceID, &t.DeviceName, &t.Address,
&t.Status, &outcomesRaw, &t.Error, &t.Transcript,
&t.StartedAt, &t.FinishedAt, &t.DurationMs); err != nil {
return err
}
t.Outcomes = []CommandOutcome{}
_ = json.Unmarshal([]byte(outcomesRaw), &t.Outcomes)
out.Counts[t.Status]++
out.Targets = append(out.Targets, t)
}
return rows.Err()
})
if err != nil {
return out, err
}
if out.Targets == nil {
out.Targets = []CommandTarget{}
}
return out, nil
}
// CancelCommandRun зупиняє прогін.
//
// Зупиняє те, до чого ще не дійшло. Хости, на яких сесія вже відкрита,
// доводяться до кінця: обірвана посеред команди сесія лишає пристрій у
// стані, якого ніхто не бачив, і це гірше за зайвий рядок виводу. Тому
// це «більше не починати», а не «скасувати все».
func (s *Store) CancelCommandRun(ctx context.Context, tenantID, runID string, userID string) (int, error) {
var stopped int
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
tag, err := tx.Exec(ctx, `
UPDATE ncm.command_runs
SET status = 'canceled', canceled_at = now(), canceled_by = $3
WHERE id = $1 AND tenant_id = $2 AND status = 'running'
`, runID, tenantID, nullUUID(userID))
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
return ErrNotFound
}
tag, err = tx.Exec(ctx, `
UPDATE ncm.command_targets
SET status = 'canceled', finished_at = now(),
error = 'прогін зупинено оператором'
WHERE run_id = $1 AND tenant_id = $2 AND status = 'pending'
`, runID, tenantID)
if err != nil {
return err
}
stopped = int(tag.RowsAffected())
return nil
})
return stopped, err
}
// ErrRunActive — прогін ще не закритий, видаляти нічого не можна.
//
// Окрема помилка, а не ErrInvalid: запит правильний, стан системи —
// тимчасовий, і людині треба сказати не «неправильно», а «зачекайте».
var ErrRunActive = errors.New("прогін ще виконується")
// DeleteCommandRun прибирає прогін разом із результатами.
//
// Повертає те, що видалив, — і саме тому не однорядковим DELETE. Запис
// в аудит має нести команди, кількість і перелік хостів: після видалення
// відповісти на питання «а що там було» більше нічим, і рядок журналу
// лишається єдиним слідом. Читання й видалення в одній транзакції, щоб
// між ними ніхто не встиг дописати ще один хост.
//
// Живий прогін не видаляється. Його рядки — це водночас завдання для
// диспетчера й адреса, за якою приймач кладе вивід із зонда: прибрати їх
// з-під сесії, що вже відкрита на пристрої, означає загубити результат
// команди, яка на залізі все одно виконається. Тому умова — finished_at,
// а не status: зупинений прогін теж доробляє хости, до яких уже
// під'єднались.
func (s *Store) DeleteCommandRun(ctx context.Context, tenantID, runID string) (CommandRun, error) {
var out CommandRun
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
var cmdRaw string
var finished *time.Time
err := tx.QueryRow(ctx, `
SELECT r.id::text, r.status::text, r.commands::text, r.total,
COALESCE(u.username,''), r.created_at, r.finished_at
FROM ncm.command_runs r
LEFT JOIN core.users u ON u.id = r.created_by
WHERE r.id = $1 AND r.tenant_id = $2
FOR UPDATE OF r
`, runID, tenantID).Scan(&out.ID, &out.Status, &cmdRaw, &out.Total,
&out.CreatedBy, &out.CreatedAt, &finished)
if errors.Is(err, pgx.ErrNoRows) {
return ErrNotFound
}
if err != nil {
return err
}
_ = jsonUnmarshalStrings(cmdRaw, &out.Commands)
out.FinishedAt = finished
if finished == nil {
if out.Status == "canceled" {
return fmt.Errorf("%w: зупинений, але кілька хостів ще доробляють сесію — "+
"прогін закриється сам і стане видимим для видалення", ErrRunActive)
}
return fmt.Errorf("%w: спершу зупиніть його", ErrRunActive)
}
// Зведення по станах — теж у журнал: «видалили прогін, у якому
// було 12 помилок» і «видалили прогін, де все пройшло» — різні
// події, і розрізняти їх постфактум буде нічим.
out.Counts = map[string]int{}
rows, err := tx.Query(ctx, `
SELECT status::text, count(*) FROM ncm.command_targets
WHERE run_id = $1 AND tenant_id = $2 GROUP BY status
`, runID, tenantID)
if err != nil {
return err
}
for rows.Next() {
var st string
var n int
if err := rows.Scan(&st, &n); err != nil {
rows.Close()
return err
}
out.Counts[st] = n
}
rows.Close()
if err := rows.Err(); err != nil {
return err
}
// Хости зникають каскадом (див. міграцію 0036): зв'язок
// один-до-багатьох, і рядок результату без прогону не означає
// нічого.
_, err = tx.Exec(ctx, `
DELETE FROM ncm.command_runs WHERE id = $1 AND tenant_id = $2
`, runID, tenantID)
return err
})
return out, err
}
// ---------------------------------------------------------------------
// Диспетчер
// ---------------------------------------------------------------------
// PendingCommandJob — завдання, готове до відправки зонду.
type PendingCommandJob struct {
TargetID string
TenantID string
AgentID string
DeviceName string
Job *npv1.ConfigJob
}
// ClaimCommandJobs забирає хости, до яких дійшла черга.
//
// Дві межі одночасно. Глобальна (limit) — щоб один такт не вивалив на
// колектор усе, що є. Порогова (concurrency прогону) — щоб та сама
// команда не пішла на всю дільницю одночасно: це стільки ж паралельних
// сесій із зонда, і зонд стоїть на тому самому каналі, що й моніторинг.
//
// Обидві рахуються в одному UPDATE. Повторна перевірка status='pending'
// у WHERE — це і є захист від подвійної відправки: другий екземпляр
// колектора дочекається блокування рядка, перечитає його вже в стані
// 'queued' і не візьме. Ліміт паралельності при перегонах двох
// колекторів може разово перевищитись на одиницю — це дешевше, ніж
// серіалізувати диспетчер на глобальному блокуванні.
func (s *Store) ClaimCommandJobs(ctx context.Context, onlineAgents []string,
limit int, ring *crypto.Keyring) ([]PendingCommandJob, error) {
if len(onlineAgents) == 0 {
return nil, nil
}
if limit <= 0 {
limit = 16
}
rows, err := s.bg.Query(ctx, `
UPDATE ncm.command_targets t
SET status = 'queued'
FROM (
SELECT c.id
FROM (
SELECT p.id, p.run_id,
row_number() OVER (PARTITION BY p.run_id
ORDER BY p.created_at, p.id) AS rn
FROM ncm.command_targets p
JOIN ncm.command_runs r ON r.id = p.run_id
WHERE p.status = 'pending'
AND r.status = 'running'
AND p.agent_id = ANY($1::uuid[])
) c
JOIN ncm.command_runs r ON r.id = c.run_id
WHERE c.rn <= r.concurrency - (
SELECT count(*) FROM ncm.command_targets x
WHERE x.run_id = c.run_id AND x.status IN ('queued','running')
)
ORDER BY c.rn
LIMIT $2
) pick
WHERE t.id = pick.id AND t.status = 'pending'
RETURNING t.id::text, t.tenant_id::text, t.agent_id::text,
t.device_id::text, t.run_id::text
`, onlineAgents, limit)
if err != nil {
return nil, err
}
type claimed struct{ targetID, tenantID, agentID, deviceID, runID string }
var list []claimed
for rows.Next() {
var c claimed
if err := rows.Scan(&c.targetID, &c.tenantID, &c.agentID, &c.deviceID, &c.runID); err != nil {
rows.Close()
return nil, err
}
list = append(list, c)
}
rows.Close()
if err := rows.Err(); err != nil {
return nil, err
}
out := make([]PendingCommandJob, 0, len(list))
for _, c := range list {
job, name, err := s.buildCommandJob(ctx, c.tenantID, c.runID, c.deviceID, c.targetID, ring)
if err != nil {
// Завдання, яке неможливо зібрати, падає зараз і з
// поясненням. Решта прогону при цьому йде далі: провал
// одного хоста — це один рядок, а не зупинка роботи.
_ = s.FinishCommandTarget(ctx, c.targetID, "failed", err.Error(), nil, "")
continue
}
out = append(out, PendingCommandJob{
TargetID: c.targetID, TenantID: c.tenantID, AgentID: c.agentID,
DeviceName: name, Job: job,
})
}
return out, nil
}
// buildCommandJob збирає завдання з прогону, профілю й доступу.
//
// job_id тут — це id рядка ncm.command_targets. Так приймач чанків
// знаходить, куди класти вивід, не заводячи ще однієї відповідності:
// одне завдання — один хост — один рядок результату.
func (s *Store) buildCommandJob(ctx context.Context, tenantID, runID, deviceID, targetID string,
ring *crypto.Keyring) (*npv1.ConfigJob, string, error) {
var (
devName, address string
cmdRaw string
timeoutSec int
profileID *string
credID *string
)
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
return tx.QueryRow(ctx, `
SELECT d.name, COALESCE(host(d.address),''),
r.commands::text, r.timeout_sec,
p.profile_id::text, p.credential_id::text
FROM ncm.command_runs r
JOIN inv.devices d ON d.id = $2 AND d.tenant_id = r.tenant_id
AND d.deleted_at IS NULL
LEFT JOIN ncm.device_policies p ON p.device_id = d.id
WHERE r.id = $1 AND r.tenant_id = $3
`, runID, deviceID, tenantID).
Scan(&devName, &address, &cmdRaw, &timeoutSec, &profileID, &credID)
})
if errors.Is(err, pgx.ErrNoRows) {
return nil, "", ErrNotFound
}
if err != nil {
return nil, "", err
}
if address == "" {
return nil, devName, fmt.Errorf("у хоста %s немає адреси", devName)
}
var userCmds []string
if err := jsonUnmarshalStrings(cmdRaw, &userCmds); err != nil {
return nil, devName, fmt.Errorf("команди прогону: %w", err)
}
if len(userCmds) == 0 {
return nil, devName, fmt.Errorf("прогін без жодної команди")
}
// Профіль дає підготовчі команди й зразок запрошення. Його
// відсутність тут — не привід відмовляти, на відміну від збору
// конфігу: там профіль каже, ЩО саме знімати, і без нього завдання
// не існує, а тут команди задала людина. Лишається типовий зразок
// запрошення — і якщо він не підійде, помилка буде зрозуміла
// («не дочекались запрошення»), а не мовчазна.
prof, perr := s.resolveProfile(ctx, tenantID, deviceID, profileID)
var prep []string
promptRegex := ""
enableRequired := false
profileTransport := ""
if perr == nil {
// Остання команда профілю — це знімання конфігу; тут вона
// зайва. Решта готує консоль (screen-length, terminal length) і
// потрібна рівно так само, як при зборі.
if len(prof.Commands) > 1 {
prep = prof.Commands[:len(prof.Commands)-1]
}
promptRegex = prof.PromptRegex
enableRequired = prof.EnableRequired
profileTransport = prof.Transport
} else if !errors.Is(perr, ErrNoProfile) {
return nil, devName, perr
}
cred, err := s.resolveNcmCredential(ctx, tenantID, deviceID, credID, ring)
if err != nil {
return nil, devName, err
}
transport, port := jobTransport(profileTransport, cred)
if timeoutSec <= 0 {
timeoutSec = 120
}
return &npv1.ConfigJob{
JobId: targetID,
Device: &npv1.DeviceTarget{
DeviceId: deviceID,
Name: devName,
Address: address,
},
Credential: cred,
Transport: transport,
Port: port,
Commands: append(append([]string{}, prep...), userCmds...),
PromptRegex: promptRegex,
EnableRequired: enableRequired,
// Ознака, за якою зонд виконує команди, а не знімає конфіг, і
// за якою приймач чанків кладе вивід у ncm.command_targets.
ConfigType: CommandConfigType,
Timeout: durationpb.New(time.Duration(timeoutSec) * time.Second),
MaxBytes: 8 << 20,
// Стенограма потрібна рівно тоді, коли вивід виглядає не так,
// як очікували: у ній видно, що саме пристрій відповів на вхід
// і на кожну команду. Пароль до неї не потрапляє.
CaptureTranscript: true,
}, devName, nil
}
// MarkCommandTargetSent позначає, що завдання пішло живій сесії.
func (s *Store) MarkCommandTargetSent(ctx context.Context, targetID string) error {
_, err := s.bg.Exec(ctx, `
UPDATE ncm.command_targets
SET status = 'running', started_at = now()
WHERE id = $1 AND status = 'queued'
`, targetID)
return err
}
// FinishCommandTarget закриває хост у прогоні.
//
// Вивід зберігається навіть при невдачі: команда могла впасти на
// третій із п'яти, і те, що встигли побачити перші дві, — половина
// відповіді на питання, чому впала третя.
func (s *Store) FinishCommandTarget(ctx context.Context, targetID, status, errMsg string,
outcomes []CommandOutcome, transcript string) error {
raw := "[]"
if len(outcomes) > 0 {
b, err := json.Marshal(outcomes)
if err != nil {
return err
}
raw = string(b)
}
_, err := s.bg.Exec(ctx, `
UPDATE ncm.command_targets
SET status = $2::ncm.command_target_status,
finished_at = now(),
duration_ms = GREATEST(0,
(EXTRACT(EPOCH FROM (now() - COALESCE(started_at, now()))) * 1000)::int),
outcomes = $3::jsonb,
error = NULLIF($4,''),
transcript = NULLIF($5,'')
WHERE id = $1
`, targetID, status, raw, errMsg, transcript)
return err
}
// CommandRunCommands — команди прогону, до якого належить хост.
//
// Потрібні приймачу, щоб відрізнити підготовчі команди профілю від тих,
// що набрала людина: у переліку результатів вони мають бути видні
// окремо. Окремої колонки під це немає навмисно — перелік уже є в
// прогоні, і друга його копія рано чи пізно розійшлася б із першою.
func (s *Store) CommandRunCommands(ctx context.Context, targetID string) ([]string, error) {
var raw string
err := s.bg.QueryRow(ctx, `
SELECT r.commands::text
FROM ncm.command_targets t
JOIN ncm.command_runs r ON r.id = t.run_id
WHERE t.id = $1
`, targetID).Scan(&raw)
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrNotFound
}
if err != nil {
return nil, err
}
var out []string
return out, jsonUnmarshalStrings(raw, &out)
}
// IsCommandTarget каже, чи належить job_id прогону команд.
//
// Приймач чанків спирається на config_type із заголовка, але заголовок
// приходить від зонда, а зонд може бути старішої версії. Звірка з базою
// коштує один запит за первинним ключем і не дає виводу команд лягти в
// архів конфігів.
func (s *Store) IsCommandTarget(ctx context.Context, jobID string) (bool, error) {
var exists bool
err := s.bg.QueryRow(ctx, `
SELECT EXISTS (SELECT 1 FROM ncm.command_targets WHERE id = $1)
`, jobID).Scan(&exists)
return exists, err
}
// ReapStuckCommandTargets повертає до тями зависле.
//
// Два різні випадки з різним лікуванням. 'queued' без відправки —
// колектор забрав рядок і не встиг надіслати (перезапуск, обрив): такий
// повертається в чергу, бо на пристрої ще нічого не відбувалось.
// 'running' понад бюджет — зонд зник разом із сесією: такий закривається
// невдачею, бо повторний запуск означав би другу сесію до пристрою, на
// якому команда, можливо, вже виконалась.
func (s *Store) ReapStuckCommandTargets(ctx context.Context, grace time.Duration) (int64, error) {
var total int64
tag, err := s.bg.Exec(ctx, `
UPDATE ncm.command_targets
SET status = 'pending'
WHERE status = 'queued' AND created_at < now() - $1::interval
`, (2 * time.Minute).String())
if err != nil {
return 0, err
}
total += tag.RowsAffected()
// Бюджет хоста — таймаут прогону плюс запас на з'єднання й
// вивантаження. Спільної константи тут бути не може: прогін із
// таймаутом у десять хвилин і прогін на десять секунд зависають
// по-різному.
tag, err = s.bg.Exec(ctx, `
UPDATE ncm.command_targets t
SET status = 'failed', finished_at = now(),
error = COALESCE(t.error, 'зонд не відповів у відведений час')
FROM ncm.command_runs r
WHERE r.id = t.run_id
AND t.status = 'running'
AND t.started_at < now() - make_interval(secs => r.timeout_sec) - $1::interval
`, grace.String())
if err != nil {
return total, err
}
return total + tag.RowsAffected(), nil
}
// SettleCommandRuns закриває прогони, у яких не лишилось роботи.
//
// Окремим кроком, а не в FinishCommandTarget: останній хост може
// закритися не через нього, а через прибиральника вище, і тримати
// однакову умову в двох місцях означало б рано чи пізно лишити прогін
// назавжди «у роботі».
func (s *Store) SettleCommandRuns(ctx context.Context) (int64, error) {
tag, err := s.bg.Exec(ctx, `
UPDATE ncm.command_runs r
SET status = CASE WHEN r.status = 'canceled' THEN 'canceled'::ncm.command_run_status
ELSE 'done'::ncm.command_run_status END,
finished_at = now()
WHERE r.finished_at IS NULL
AND NOT EXISTS (
SELECT 1 FROM ncm.command_targets t
WHERE t.run_id = r.id
AND t.status IN ('pending','queued','running')
)
`)
if err != nil {
return 0, err
}
return tag.RowsAffected(), nil
}
// lowerAll — регістронезалежне порівняння виробників.
//
// «Huawei», «huawei» і «HUAWEI» в інвентарі трапляються всі три: поле
// заповнюють і люди руками, і автовизначення за sysDescr.
func lowerAll(in []string) []string {
out := make([]string, 0, len(in))
for _, s := range in {
out = append(out, strings.ToLower(strings.TrimSpace(s)))
}
return out
}