Netpulse_SasS/server/internal/store/ncm_jobs.go
byrsapty 452345da32
Some checks are pending
CI / web (push) Waiting to run
CI / server (push) Waiting to run
CI / agent (push) Waiting to run
Транспорт збору конфігу бере доступ хоста, а не профіль моделі
Питання «чому всі профілі ssh» виявило справжню помилку. Список
правильний: ssh — розумний типовий вибір. Але це поле ВИРІШУВАЛО
транспорт, а профіль описує модель, не конкретну коробку.

Хост зі старою прошивкою, де є лише telnet, і з чесно заведеним
telnet-доступом усе одно набирався по SSH — і не збирався ніколи.
inv.credentials.proto вже ніс потрібну відповідь і доїжджав до зонда
невикористаним.

Тепер транспорт бере доступ; профільне поле лишилось підказкою й діє,
поки доступу немає. Логіка винесена в jobTransport() — щоб її можна
було перевірити без бази й без пристрою.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-25 16:45:48 +03:00

421 lines
15 KiB
Go
Raw 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"
"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"
)
var ErrNoProfile = errors.New("для хоста не задано профіль збору конфігу")
// EnqueueConfigJob ставить збір конфігу в чергу.
//
// Саме черга в БД, а не прямий виклик: HTTP-процес і AgentService —
// різні процеси, і живу сесію зонда тримає лише другий. Черга робить
// передачу між ними явною й переживає перезапуск обох.
func (s *Store) EnqueueConfigJob(ctx context.Context, tenantID, deviceID, trigger, userID string) (string, error) {
var id string
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
// Другий queued для того самого хоста нічого не додає: перший
// збере той самий конфіг. Тому натискання кнопки двічі не
// множить роботу.
err := tx.QueryRow(ctx, `
SELECT id::text FROM ncm.jobs
WHERE tenant_id = $1 AND device_id = $2 AND status IN ('queued','running')
LIMIT 1
`, tenantID, deviceID).Scan(&id)
if err == nil {
return nil
}
if !errors.Is(err, pgx.ErrNoRows) {
return err
}
return tx.QueryRow(ctx, `
INSERT INTO ncm.jobs (tenant_id, device_id, agent_id, trigger, requested_by)
SELECT $1, d.id, d.agent_id, $3::ncm.job_trigger, $4
FROM inv.devices d
WHERE d.id = $2 AND d.tenant_id = $1 AND d.deleted_at IS NULL
RETURNING id::text
`, tenantID, deviceID, trigger, nullUUID(userID)).Scan(&id)
})
if errors.Is(err, pgx.ErrNoRows) {
return "", ErrNotFound
}
return id, err
}
// PendingConfigJob — завдання, готове до відправки зонду.
type PendingConfigJob struct {
JobID string
TenantID string
AgentID string
Job *npv1.ConfigJob
}
// ClaimConfigJobs забирає завдання для зондів, які зараз на зв'язку.
//
// Забирає, а не читає: рядок одразу переходить у 'running'. Два
// екземпляри AgentService за балансувальником інакше надіслали б одне
// завдання двічі, і пристрій отримав би дві паралельні сесії.
func (s *Store) ClaimConfigJobs(ctx context.Context, onlineAgents []string, limit int, ring *crypto.Keyring) ([]PendingConfigJob, error) {
if len(onlineAgents) == 0 {
return nil, nil
}
if limit <= 0 {
limit = 16
}
rows, err := s.pool.Query(ctx, `
UPDATE ncm.jobs j
SET status = 'running', started_at = now()
WHERE j.id IN (
SELECT id FROM ncm.jobs
WHERE status = 'queued' AND agent_id = ANY($1::uuid[])
ORDER BY created_at
-- SKIP LOCKED: другий екземпляр не чекає на нас, а бере
-- наступні завдання. Без цього паралельні диспетчери
-- вишикувались би в чергу за одним рядком.
FOR UPDATE SKIP LOCKED
LIMIT $2
)
RETURNING j.id::text, j.tenant_id::text, j.agent_id::text, j.device_id::text
`, onlineAgents, limit)
if err != nil {
return nil, err
}
type claimed struct{ jobID, tenantID, agentID, deviceID string }
var list []claimed
for rows.Next() {
var c claimed
if err := rows.Scan(&c.jobID, &c.tenantID, &c.agentID, &c.deviceID); err != nil {
rows.Close()
return nil, err
}
list = append(list, c)
}
rows.Close()
if err := rows.Err(); err != nil {
return nil, err
}
out := make([]PendingConfigJob, 0, len(list))
for _, c := range list {
job, err := s.buildConfigJob(ctx, c.tenantID, c.deviceID, c.jobID, ring)
if err != nil {
// Завдання, яке неможливо зібрати, має впасти зараз і з
// поясненням, а не висіти в 'running' до перезапуску.
_ = s.FinishConfigJob(ctx, c.jobID, "failed", err.Error(), "")
continue
}
out = append(out, PendingConfigJob{
JobID: c.jobID, TenantID: c.tenantID, AgentID: c.agentID, Job: job,
})
}
return out, nil
}
// buildConfigJob збирає завдання з профілю, хоста й доступу.
func (s *Store) buildConfigJob(ctx context.Context, tenantID, deviceID, jobID string, ring *crypto.Keyring) (*npv1.ConfigJob, error) {
var (
devName, address string
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),''),
p.profile_id::text, p.credential_id::text
FROM inv.devices d
LEFT JOIN ncm.device_policies p ON p.device_id = d.id
WHERE d.id = $1 AND d.tenant_id = $2 AND d.deleted_at IS NULL
`, deviceID, tenantID).Scan(&devName, &address, &profileID, &credID)
})
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrNotFound
}
if err != nil {
return nil, err
}
if address == "" {
return nil, fmt.Errorf("у хоста %s немає адреси", devName)
}
prof, err := s.resolveProfile(ctx, tenantID, deviceID, profileID)
if err != nil {
return nil, err
}
cred, err := s.resolveNcmCredential(ctx, tenantID, deviceID, credID, ring)
if err != nil {
return nil, err
}
transport, port := jobTransport(prof.Transport, cred)
return &npv1.ConfigJob{
JobId: jobID,
Device: &npv1.DeviceTarget{
DeviceId: deviceID,
Name: devName,
Address: address,
},
Credential: cred,
Transport: transport,
Port: port,
Commands: prof.Commands,
PromptRegex: prof.PromptRegex,
EnableRequired: prof.EnableRequired,
ConfigType: "running",
Timeout: durationpb.New(5 * time.Minute),
MaxBytes: 16 << 20,
// Транскрипт пишемо завжди: він потрібен рівно тоді, коли збір
// не вдався, а вдруге відтворити ту саму сесію не вийде.
CaptureTranscript: true,
}, nil
}
type ncmProfile struct {
Commands []string
PromptRegex string
EnableRequired bool
Transport string
// rawCommands — jsonb із БД до розбору.
rawCommands string
}
// jsonUnmarshalStrings розбирає масив рядків із jsonb.
func jsonUnmarshalStrings(raw string, out *[]string) error {
if raw == "" {
return nil
}
return json.Unmarshal([]byte(raw), out)
}
// resolveProfile шукає профіль: спершу заданий явно, потім за вендором.
//
// Автопідбір за вендором — не здогад, а єдиний спосіб не змушувати
// адміністратора вручну призначати профіль кожному з тисячі хостів.
// Явно заданий при цьому завжди виграє.
func (s *Store) resolveProfile(ctx context.Context, tenantID, deviceID string, explicit *string) (ncmProfile, error) {
var p ncmProfile
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
if explicit != nil && *explicit != "" {
return tx.QueryRow(ctx, `
SELECT commands::text, COALESCE(prompt_regex,''), enable_required, transport::text
FROM ncm.profiles WHERE id = $1
`, *explicit).Scan(&p.rawCommands, &p.PromptRegex, &p.EnableRequired, &p.Transport)
}
return tx.QueryRow(ctx, `
SELECT pr.commands::text, COALESCE(pr.prompt_regex,''),
pr.enable_required, pr.transport::text
FROM inv.devices d
JOIN ncm.profiles pr
ON lower(pr.vendor) = lower(d.vendor)
AND (pr.tenant_id IS NULL OR pr.tenant_id = d.tenant_id)
WHERE d.id = $1 AND d.vendor IS NOT NULL
-- Профіль тенанта важить більше за вбудований: клієнт міг
-- підправити команди під свою прошивку.
ORDER BY pr.tenant_id NULLS LAST, pr.key
LIMIT 1
`, deviceID).Scan(&p.rawCommands, &p.PromptRegex, &p.EnableRequired, &p.Transport)
})
if errors.Is(err, pgx.ErrNoRows) {
return p, ErrNoProfile
}
if err != nil {
return p, err
}
if err := jsonUnmarshalStrings(p.rawCommands, &p.Commands); err != nil {
return p, fmt.Errorf("команди профілю: %w", err)
}
if len(p.Commands) == 0 {
return p, ErrNoProfile
}
return p, nil
}
// resolveNcmCredential бере доступ для CLI й розшифровує його.
func (s *Store) resolveNcmCredential(ctx context.Context, tenantID, deviceID string,
explicit *string, ring *crypto.Keyring) (*npv1.Credential, error) {
var (
id, proto, username string
port int
keyID, aad *string
nonce, ct, tag []byte
)
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
q := `
SELECT c.id::text, c.proto::text, COALESCE(c.username,''), COALESCE(c.port,0),
s.key_id, s.nonce, s.ciphertext, s.auth_tag, COALESCE(s.aad,'')
FROM inv.credentials c
LEFT JOIN core.secrets s ON s.id = c.secret_id
WHERE c.tenant_id = $1 AND c.id = $2
`
if explicit != nil && *explicit != "" {
return tx.QueryRow(ctx, q, tenantID, *explicit).
Scan(&id, &proto, &username, &port, &keyID, &nonce, &ct, &tag, &aad)
}
// Не задано явно — беремо прив'язаний до хоста доступ, придатний
// для CLI. SNMP-community сюди не годиться.
return tx.QueryRow(ctx, `
SELECT c.id::text, c.proto::text, COALESCE(c.username,''), COALESCE(c.port,0),
s.key_id, s.nonce, s.ciphertext, s.auth_tag, COALESCE(s.aad,'')
FROM inv.device_credentials dc
JOIN inv.credentials c ON c.id = dc.credential_id
LEFT JOIN core.secrets s ON s.id = c.secret_id
WHERE dc.device_id = $1 AND c.tenant_id = $2
AND c.proto IN ('ssh','telnet')
ORDER BY dc.priority
LIMIT 1
`, deviceID, tenantID).
Scan(&id, &proto, &username, &port, &keyID, &nonce, &ct, &tag, &aad)
})
if errors.Is(err, pgx.ErrNoRows) {
return nil, fmt.Errorf("для хоста не прив'язано доступу ssh або telnet")
}
if err != nil {
return nil, err
}
cred := &npv1.Credential{
CredentialId: id,
Username: username,
Port: uint32(port),
}
if proto == "telnet" {
cred.Transport = npv1.Transport_TRANSPORT_TELNET
} else {
cred.Transport = npv1.Transport_TRANSPORT_SSH
}
if keyID != nil && ring != nil {
plain, err := ring.Decrypt(&crypto.Secret{
KeyID: *keyID, Nonce: nonce, Ciphertext: ct, AuthTag: tag,
}, derefStr(aad))
if err != nil {
return nil, fmt.Errorf("розшифровка доступу: %w", err)
}
cred.Secret = &npv1.Credential_Password{Password: string(plain)}
}
return cred, nil
}
// FinishConfigJob закриває завдання.
func (s *Store) FinishConfigJob(ctx context.Context, jobID, status, errMsg, transcript string) error {
_, err := s.pool.Exec(ctx, `
UPDATE ncm.jobs
SET status = $2::ncm.job_status,
finished_at = now(),
duration_ms = GREATEST(0, EXTRACT(EPOCH FROM (now() - COALESCE(started_at, now())))::int * 1000),
error = NULLIF($3,''),
log = NULLIF($4,'')
WHERE id = $1
`, jobID, status, errMsg, transcript)
return err
}
// ReapStuckJobs повертає в чергу завдання, що зависли в 'running'.
//
// Зонд міг зникнути разом із завданням: без цього хост залишився б без
// бекапів назавжди, а причина була б видна лише в таблиці.
func (s *Store) ReapStuckJobs(ctx context.Context, olderThan time.Duration) (int64, error) {
tag, err := s.pool.Exec(ctx, `
UPDATE ncm.jobs
SET status = 'failed', finished_at = now(),
error = COALESCE(error, 'зонд не відповів у відведений час')
WHERE status = 'running' AND started_at < now() - $1::interval
`, olderThan.String())
if err != nil {
return 0, err
}
return tag.RowsAffected(), nil
}
// ConfigJobRow — рядок історії збору для UI.
type ConfigJobRow struct {
ID string `json:"id"`
Status string `json:"status"`
Trigger string `json:"trigger"`
StartedAt *time.Time `json:"started_at,omitempty"`
FinishedAt *time.Time `json:"finished_at,omitempty"`
DurationMs int `json:"duration_ms"`
Error string `json:"error,omitempty"`
CreatedAt time.Time `json:"created_at"`
}
func (s *Store) ListConfigJobs(ctx context.Context, tenantID, deviceID string, limit int) ([]ConfigJobRow, error) {
if limit <= 0 || limit > 200 {
limit = 20
}
var out []ConfigJobRow
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
rows, err := tx.Query(ctx, `
SELECT id::text, status::text, trigger::text,
started_at, finished_at, COALESCE(duration_ms,0),
COALESCE(error,''), created_at
FROM ncm.jobs
WHERE tenant_id = $1 AND device_id = $2
ORDER BY created_at DESC
LIMIT $3
`, tenantID, deviceID, limit)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var j ConfigJobRow
if err := rows.Scan(&j.ID, &j.Status, &j.Trigger, &j.StartedAt,
&j.FinishedAt, &j.DurationMs, &j.Error, &j.CreatedAt); err != nil {
return err
}
out = append(out, j)
}
return rows.Err()
})
return out, err
}
// jobTransport каже, чим набирати цей хост і на який порт.
//
// Транспорт вирішує ДОСТУП, а не профіль. Профіль описує модель заліза:
// які команди віддати й як упізнати запрошення. А чим до конкретної
// коробки достукатись — властивість самої коробки: та сама модель у
// клієнта може стояти з увімкненим SSH, а на сусідньому вузлі — зі
// старою прошивкою, де є лише telnet.
//
// Раніше вирішував профіль, і оскільки всі 147 вбудованих кажуть «ssh»,
// хост із telnet-доступом набирався по SSH і не збирався ніколи.
// Профільне поле лишилось підказкою «чим це залізо зазвичай беруть» і
// діє, тільки поки доступу немає.
func jobTransport(profileTransport string, cred *npv1.Credential) (npv1.Transport, uint32) {
transport := npv1.Transport_TRANSPORT_SSH
switch {
case cred != nil:
transport = cred.Transport
case profileTransport == "telnet":
transport = npv1.Transport_TRANSPORT_TELNET
}
port := uint32(22)
if transport == npv1.Transport_TRANSPORT_TELNET {
port = 23
}
// Порт із доступу перекриває типовий: залізо за NAT цілком може
// слухати SSH на 2222.
if cred != nil && cred.Port != 0 {
port = cred.Port
}
return transport, port
}