Netpulse_SasS/server/internal/store/ncm_jobs.go
byrsapty fe4ed4e7d9
Some checks failed
CI / web (push) Has been cancelled
CI / server (push) Has been cancelled
CI / agent (push) Has been cancelled
Дедлайн сокета від входу через telnet обривав збір конфігу
Знайдено на живому ZTE C320. Вхід чекає «Username:» з дедлайном читання
у дві секунди, а дедлайн сокета липкий: він лишався на всі наступні
читання.

Наслідок подвійний. Довгий конфіг обривався на середині помилкою
«read tcp: i/o timeout» — мережевою там, де мережа ні до чого. А коли
ще й запрошення не збігалось зі зразком, той самий дедлайн спрацьовував
раніше за власний таймер, і замість «не дочекались запрошення» людина
бачила ту саму мережеву помилку.

Щоб це побачити, довелося спершу полагодити діагностику: стенограма
починалась після входу, тобто після того місця, де все зупинялось.
Тепер розмова входу пишеться теж — без пароля, лише те, що надіслав
пристрій.

Плюс вбудований профіль zte-zxan: наявні профілі ZTE чекають «>» або
«]», а OLT показує «ZXAN#». Окремий профіль, бо в Comware трапляються
рядки з самої решітки, і зразок із «#» обірвав би їхній конфіг.

Результат на живому C320: 26 555 рядків.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-25 19:04:36 +03:00

442 lines
16 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"
)
// ErrNoProfile лишається для порівняння через errors.Is; текст
// уточнюється на місці — «немає профілю» має три різні причини, і
// кожна лікується по-своєму.
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
// Виробник хоста — єдине, за чим профіль підбирається сам. Читаємо
// його заздалегідь, щоб відмова могла сказати, чого саме бракує:
// «виробник не заданий» і «під цього виробника немає профілю» —
// різні проблеми з різними діями.
var vendor string
_ = s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
return tx.QueryRow(ctx, `
SELECT COALESCE(vendor,'') FROM inv.devices WHERE id = $1 AND tenant_id = $2
`, deviceID, tenantID).Scan(&vendor)
})
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) {
switch {
case vendor == "":
return p, fmt.Errorf("%w: у хоста не заданий виробник, і профіль нема за чим "+
"підібрати — оберіть його вручну у вкладці «Збір конфігів»", ErrNoProfile)
default:
return p, fmt.Errorf("%w: під виробника %q немає готового профілю — "+
"оберіть інший у вкладці «Збір конфігів» або створіть свій", ErrNoProfile, vendor)
}
}
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, fmt.Errorf("%w: обраний профіль не містить жодної команди", 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
}