diff --git a/HISTORY.md b/HISTORY.md index 7ffa8e1..94457cb 100644 --- a/HISTORY.md +++ b/HISTORY.md @@ -1148,3 +1148,93 @@ Huawei VRP3, H3C, Juniper JUNOSe, Eltex SMG/TAU (без пейджера). **Виконувати профілі нікому.** Агентського модуля `ncm` (SSH/Telnet) немає — це Етап 7. Зараз це готові дані, які чекають на виконавця. + +--- + +## 2026-08-16 — Етап 7: збір конфігів запрацював + +Профілі з попереднього кроку були даними без виконавця. Тепер ланцюг +замкнено: кнопка в UI → черга → диспетчер → зонд → SSH → назад у БД. + +### Створено + +- `agent/internal/ncmx/` — знімання конфігу по CLI: `session.go` + (розбір потоку), `transport.go` (SSH і Telnet), `collect.go` + (виконання завдання). +- `agent/internal/session/config_jobs.go` — приймання `ConfigJob` і + вивантаження стрімом. +- `server/internal/store/ncm_jobs.go` — черга, побудова завдання з + профілю й доступу, закриття. +- `server/internal/grpcapi/ncm_dispatch.go` — диспетчер. +- 16 тестів на розбір консолі. + +### Прийняті рішення + +**Усе будується навколо пошуку промпту.** Консоль мережевого пристрою — +не програмний інтерфейс: немає ані коду завершення, ані довжини +відповіді. Єдиний спосіб зрозуміти, що команда відпрацювала — побачити +знову запрошення. + +**Промпт шукається лише в хвості накопиченого (512 байтів).** Конфіг +може містити рядок, схожий на запрошення — `banner motd #` трапляється +в кожній другій мережі, — і пошук по всьому тексту обривав би збір на +середині. + +**Дедлайн на паузу між байтами, а не на всю операцію.** Збір із +великого шасі триває хвилини, і загальний ліміт довелося б ставити +навмання. Тиша ж означає одне з двох: пристрій завис або промпт не той. +Тому й помилка окрема — `ErrPromptTimeout` підказує, що лікується вона +не повтором, а виправленням `prompt_regex`. + +**Черга в БД між процесами.** REST і AgentService — різні процеси; +живу сесію зонда тримає лише другий. `FOR UPDATE SKIP LOCKED` не дає +двом екземплярам надіслати одне завдання двічі. + +**Ключі SSH мережевого обладнання не звіряються.** Свідомо: залізо +перегенеровує ключ після кожної заміни прошивки, і known_hosts на сотні +пристроїв означав би або вимикати перевірку щотижня, або не збирати +конфіги зовсім. Захист тут дає сегмент керування, а не TOFU. Переліки +алгоритмів навмисно широкі — інакше половина парку відпаде з «no common +algorithm». + +### Знайдено живим прогоном (три справжні помилки) + +Прогін ставили проти самого стенду: у нього є SSH, і `cat +/etc/os-release` віддає текст так само, як консоль віддає конфіг. + +**Сервер ігнорував поле `encoding`.** Агент стискав тіло gzip і чесно +рахував sha256 від оригіналу, а сервер рахував від стиснених байтів — +і відхиляв кожен бекап як «тіло не відповідає заявленому sha256». Поле +було в контракті з Етапу 2, реалізації не було ніколи: інтеграційний +тест користувався `encoding: "none"` і повз цю дірку проходив. + +**Промпт обрізався не по рядку.** Типовий шаблон `[>#]\s*$` збігається +лише з символом запрошення, тому в конфізі лишалось ім'я пристрою +окремим рядком: «…interface Gi0/1» + «sw1». Ріжемо весь рядок. + +**У конфіг потрапляло сміття терміналу.** Перший успішний збір дав 295 +байтів і 12 рядків там, де файл має 286 і 10. Транскрипт (який сам же +модуль і зберіг) показав причину: escape-послідовності bash +`ESC[?2004l` і подвоєні `\r\r\n` від псевдотерміналу — PTY додає свій +`\r` до пристроєвого `\r\n`, і кожен рядок подвоювався. Після +виправлення: 285 байтів і рівно 10 рядків, різниця з оригіналом лише у +фінальному переводі рядка, який відрізається свідомо. + +### Перевірено наживо + +``` +завдання в черзі → диспетчер забрав, зонд отримав +зонд зайшов по SSH → виконав команду профілю +вивантажив gzip-стрімом → сервер розпакував, звірив sha256 +статус → success, конфіг 285 байтів / 10 рядків +повторний збір → unchanged, другої версії не створено +``` + +### Чого ще немає + +Планувальника за `ncm.device_policies.cron` — збір запускається лише +вручну або зовнішнім тригером. Тригера за Syslog-подією +(`%SYS-5-CONFIG_I`). Git-двигуна: конфіг лягає в БД зашифрованим, але +коміту в репозиторій ще немає, тому `commit_sha` порожній. Візуального +diff у вебі й кнопки «зібрати зараз» в інтерфейсі — API є, сторінки +немає. diff --git a/ROADMAP.md b/ROADMAP.md index 2680bb7..65df1d5 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -24,7 +24,7 @@ API віддає готове полотно з живими статусами, | Автовиявлення LLDP/CDP/ARP/FDB | ✅ | ✅ | | **Користувачі, ролі, вхід** | ✅ | ✅ | | **Шаблони опитування** | ❌ | ❌ | -| **NCM (збір конфігів)** | ✅ | ⚠️ half | +| **NCM (збір конфігів)** | ✅ | ⚠️ збір працює, Git і розклад — ні | | **Керування зондом із UI** | ✅ | ⚠️ транспорт є | | **Алерти й сповіщення** | ✅ | ✅ | | **Мобільна адаптивність, PWA** | — | ⚠️ адаптив є, PWA немає | @@ -157,7 +157,15 @@ JSONPath, `discard unchanged`) — на сервері при записі: ін --- -## Етап 7. NCM — збір конфігів до кінця +## Етап 7. NCM — збір конфігів до кінця — ⚠️ наполовину, 2026-08-16 + +> **Зроблено:** агентський модуль SSH/Telnet, черга завдань, диспетчер, +> вивантаження зі стисненням і дедуплікацією. Перевірено наскрізно. +> +> **Лишилось:** планувальник за cron, тригер за Syslog, Git-двигун, +> візуальний diff у вебі. + +## Етап 7 (початковий план) Половина шляху вже є: `ncm.*` у схемі, `ConfigJob`/`ConfigUpload` у контракті, сервер приймає чанки, звіряє sha256, дедуплікує за `content_hash` і шифрує тіло. diff --git a/agent/go.mod b/agent/go.mod index f624d44..082a3e0 100644 --- a/agent/go.mod +++ b/agent/go.mod @@ -5,14 +5,15 @@ go 1.25.0 require ( github.com/gosnmp/gosnmp v1.42.1 github.com/netpulse/netpulse/gen/go v0.0.0 - golang.org/x/net v0.55.0 + golang.org/x/crypto v0.55.0 + golang.org/x/net v0.57.0 google.golang.org/grpc v1.83.0 google.golang.org/protobuf v1.36.12 ) require ( - golang.org/x/sys v0.45.0 // indirect - golang.org/x/text v0.37.0 // indirect + golang.org/x/sys v0.47.0 // indirect + golang.org/x/text v0.41.0 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect ) diff --git a/agent/go.sum b/agent/go.sum index 316a484..b92624e 100644 --- a/agent/go.sum +++ b/agent/go.sum @@ -30,12 +30,16 @@ go.opentelemetry.io/otel/sdk/metric v1.44.0 h1:3LlKgI+VjbVsjNRFZJZAJ30WjXC5VkNRk go.opentelemetry.io/otel/sdk/metric v1.44.0/go.mod h1:5B5pMARnXxKhltooO4xUuCBorl65a4EpnTalObqOigA= go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk= go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= -golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8= -golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww= -golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= -golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= -golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= -golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= +golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M= +golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis= +golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE= +golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU= +golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= +golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/term v0.45.0 h1:NwWyBmoJCbfTHpxrWoZ9C6/VxOf7ic219I8xZZFdrf0= +golang.org/x/term v0.45.0/go.mod h1:9aqxs0blBcrm/n0L9QW0aRVD+ktan8ssZromtqJC43w= +golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= +golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E= google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa h1:mZHHdPZl0dbGHCflZgAq/Q468DWVFcU2whhB2KAo8fk= diff --git a/agent/internal/ncmx/collect.go b/agent/internal/ncmx/collect.go new file mode 100644 index 0000000..7d55f0e --- /dev/null +++ b/agent/internal/ncmx/collect.go @@ -0,0 +1,216 @@ +package ncmx + +import ( + "bytes" + "context" + "crypto/sha256" + "fmt" + "regexp" + "strings" + "time" + + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" +) + +// Result — те, що агент відправляє на сервер. +type Result struct { + Body []byte + SHA256 []byte + LineCount int + Duration time.Duration + Transcript string +} + +// Collect виконує завдання збору конфігу. +// +// Послідовність команд трактується так, як домовлено в контракті: +// вивід ОСТАННЬОЇ команди — це конфіг, попередні готують консоль. +// Домовленість, а не здогад: інакше довелося б або позначати +// «головну» команду прапорцем, або вгадувати її за текстом. +func Collect(ctx context.Context, job *npv1.ConfigJob) (Result, error) { + start := time.Now() + + if len(job.GetCommands()) == 0 { + return Result{}, fmt.Errorf("завдання без жодної команди") + } + + promptRe, err := compilePrompt(job.GetPromptRegex()) + if err != nil { + return Result{}, err + } + + timeout := job.GetTimeout().AsDuration() + if timeout <= 0 { + timeout = 5 * time.Minute + } + ctx, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + + cred := job.GetCredential() + conn, err := Dial(ctx, transportName(job.GetTransport()), + job.GetDevice().GetAddress(), int(job.GetPort()), + cred.GetUsername(), password(cred), connectTimeout(timeout)) + if err != nil { + return Result{}, err + } + defer conn.Close() + + var transcript *bytes.Buffer + if job.GetCaptureTranscript() { + transcript = &bytes.Buffer{} + } + + cli := NewCLI(conn, Options{ + PromptRe: promptRe, + MaxBytes: int(job.GetMaxBytes()), + Transcript: transcript, + }) + + // Банер до першої команди. + if err := cli.WaitPrompt(ctx); err != nil { + return withTranscript(Result{}, transcript), fmt.Errorf("привітання пристрою: %w", err) + } + + if job.GetEnableRequired() { + if err := enable(ctx, cli, cred.GetEnablePassword()); err != nil { + return withTranscript(Result{}, transcript), err + } + } + + cmds := job.GetCommands() + // Підготовчі команди: їхній вивід нікого не цікавить, а помилка не + // фатальна. «terminal length 0» на платформі, яка такої команди не + // знає, дасть «% Invalid input» — і це не привід не збирати конфіг. + for _, c := range cmds[:len(cmds)-1] { + if _, err := cli.Run(ctx, c); err != nil { + return withTranscript(Result{}, transcript), + fmt.Errorf("підготовча команда %q: %w", c, err) + } + } + + last := cmds[len(cmds)-1] + body, err := cli.Run(ctx, last) + if err != nil { + return withTranscript(Result{}, transcript), fmt.Errorf("команда %q: %w", last, err) + } + if strings.TrimSpace(body) == "" { + return withTranscript(Result{}, transcript), + fmt.Errorf("команда %q повернула порожній вивід", last) + } + + b := []byte(body) + sum := sha256.Sum256(b) + + res := Result{ + Body: b, + SHA256: sum[:], + LineCount: strings.Count(body, "\n") + 1, + Duration: time.Since(start), + } + return withTranscript(res, transcript), nil +} + +// enable піднімає привілеї. +// +// Запрошення пароля впізнаємо за текстом, а не за окремим регексом із +// завдання: «Password:» пишуть однаково всі, а тримати ще один +// налаштовуваний шаблон означало б ще одне поле, яке нікому не хочеться +// заповнювати правильно. +var enablePasswordRe = regexp.MustCompile(`(?i)password\s*:?\s*$`) + +func enable(ctx context.Context, cli *CLI, secret string) error { + if _, err := cli.conn.Write([]byte("enable\n")); err != nil { + return fmt.Errorf("enable: %w", err) + } + + // Після «enable» пристрій відповідає одним із двох: або одразу + // привілейованим запрошенням (пароль не потрібен), або запитом + // пароля. Розрізняємо за тим, що прийшло. + saved := cli.opt.PromptRe + cli.opt.PromptRe = regexp.MustCompile(saved.String() + `|(?i)password\s*:?\s*$`) + out, err := cli.readUntilPrompt(ctx) + cli.opt.PromptRe = saved + if err != nil { + return fmt.Errorf("enable: %w", err) + } + + if !enablePasswordRe.MatchString(strings.TrimSpace(out)) && + !strings.Contains(strings.ToLower(out), "password") { + return nil // пароль не запитували + } + if secret == "" { + return fmt.Errorf("пристрій просить пароль enable, а його не задано") + } + if _, err := cli.conn.Write([]byte(secret + "\n")); err != nil { + return fmt.Errorf("пароль enable: %w", err) + } + if _, err := cli.readUntilPrompt(ctx); err != nil { + return fmt.Errorf("після пароля enable: %w", err) + } + return nil +} + +// compilePrompt готує регекс запрошення. +// +// Порожній шаблон — не привід падати: типовий «>» або «#» у кінці рядка +// покриває більшість платформ, і краще спробувати з ним, ніж відмовити +// профілю, у якого поле просто не заповнили. +func compilePrompt(pattern string) (*regexp.Regexp, error) { + if strings.TrimSpace(pattern) == "" { + pattern = `[>#]\s*$` + } + // Багаторядковий режим обов'язковий: без нього «$» означає кінець + // усього тексту, а запрошення стоїть у кінці ОСТАННЬОГО рядка. + re, err := regexp.Compile(`(?m)` + pattern) + if err != nil { + return nil, fmt.Errorf("некоректний prompt_regex %q: %w", pattern, err) + } + return re, nil +} + +// password дістає секрет із oneof. +// +// Для CLI підходить лише пароль: ключ вимагав би ще й розбору формату +// та passphrase, а community й token до SSH не мають стосунку. Тому +// решта варіантів свідомо дає порожній рядок, і причина відмови буде +// видна одразу при вході, а не десь усередині рукостискання. +func password(c *npv1.Credential) string { + if c == nil { + return "" + } + if p, ok := c.GetSecret().(*npv1.Credential_Password); ok { + return p.Password + } + return "" +} + +func transportName(t npv1.Transport) string { + switch t { + case npv1.Transport_TRANSPORT_TELNET: + return "telnet" + default: + return "ssh" + } +} + +// connectTimeout — частка загального бюджету на саме підключення. +// +// Збір конфігу може тривати хвилини, але якщо пристрій не відповідає на +// TCP за пів хвилини, чекати решту дедлайну немає сенсу. +func connectTimeout(total time.Duration) time.Duration { + t := total / 5 + if t > 30*time.Second { + t = 30 * time.Second + } + if t < 5*time.Second { + t = 5 * time.Second + } + return t +} + +func withTranscript(r Result, b *bytes.Buffer) Result { + if b != nil { + r.Transcript = b.String() + } + return r +} diff --git a/agent/internal/ncmx/session.go b/agent/internal/ncmx/session.go new file mode 100644 index 0000000..b20e770 --- /dev/null +++ b/agent/internal/ncmx/session.go @@ -0,0 +1,310 @@ +// Package ncmx — знімання конфігу з обладнання по CLI. +// +// Головна складність тут не в мережі, а в тому, що консоль мережевого +// пристрою — це не програмний інтерфейс. Немає ані коду завершення, ані +// довжини відповіді: єдиний спосіб зрозуміти, що команда відпрацювала — +// побачити знову запрошення. Тому весь модуль побудований навколо +// пошуку промпту в потоці байтів. +package ncmx + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + "regexp" + "strings" + "time" +) + +// ErrPromptTimeout — промпт не з'явився у відведений час. +// +// Окрема помилка, бо це найчастіша причина невдачі й лікується вона не +// повтором, а виправленням prompt_regex у профілі. Плутати її з +// мережевою відмовою означало б ганяти ретраї там, де вони не +// допоможуть. +var ErrPromptTimeout = errors.New("не дочекались запрошення командного рядка") + +// ErrTooLarge — пристрій віддав більше, ніж дозволено. +var ErrTooLarge = errors.New("вивід перевищив ліміт") + +// Conn — те, що вміє двонаправлений обмін символами. SSH-сесія й +// Telnet-з'єднання зводяться до нього, і решта коду про різницю не знає. +type Conn interface { + io.Reader + io.Writer + Close() error +} + +// Options — налаштування знімання. +type Options struct { + // PromptRe — за чим упізнаємо, що пристрій готовий приймати команду. + PromptRe *regexp.Regexp + // MoreRe — посторінковий вивід. Якщо трапиться, шлемо MoreSend. + MoreRe *regexp.Regexp + MoreSend string + + // MaxBytes — стеля на команду. Пристрій із зациклленим виводом + // інакше з'їсть пам'ять агента, а не свою. + MaxBytes int + // IdleTimeout — скільки чекати наступного байта. + IdleTimeout time.Duration + + // Transcript накопичує весь обмін для діагностики prompt_regex. + Transcript *bytes.Buffer +} + +// Стандартне «ще» для більшості вендорів. Пробіл, а не Enter: Enter на +// частині платформ гортає рядок, а не сторінку, і збір конфігу +// розтягується на тисячі проходів. +const DefaultMoreSend = " " + +// DefaultMoreRe ловить типові запрошення посторінкового виводу. +var DefaultMoreRe = regexp.MustCompile(`(?i)--\s*more\s*--|<--- more --->|\x1b\[7m--More--`) + +// CLI — жива сесія з пристроєм. +type CLI struct { + conn Conn + opt Options + + // buf тримає прочитане, але ще не розібране. Промпт може прийти + // розірваним між двома read(), тому склеювати доводиться тут. + buf bytes.Buffer +} + +func NewCLI(conn Conn, opt Options) *CLI { + if opt.MaxBytes <= 0 { + opt.MaxBytes = 16 << 20 + } + if opt.IdleTimeout <= 0 { + opt.IdleTimeout = 20 * time.Second + } + if opt.MoreRe == nil { + opt.MoreRe = DefaultMoreRe + } + if opt.MoreSend == "" { + opt.MoreSend = DefaultMoreSend + } + return &CLI{conn: conn, opt: opt} +} + +// WaitPrompt читає до першого запрошення. Викликається одразу після +// підключення: більшість пристроїв вітаються банером, і виконувати +// команду до того, як він добіжить, — це отримати її відлуння всередині +// банера. +func (c *CLI) WaitPrompt(ctx context.Context) error { + _, err := c.readUntilPrompt(ctx) + return err +} + +// Run виконує команду й повертає її вивід без відлуння та без +// фінального запрошення. +func (c *CLI) Run(ctx context.Context, cmd string) (string, error) { + if _, err := c.conn.Write([]byte(cmd + "\n")); err != nil { + return "", fmt.Errorf("надсилання %q: %w", cmd, err) + } + + raw, err := c.readUntilPrompt(ctx) + if err != nil { + return raw, err + } + return cleanOutput(raw, cmd), nil +} + +// readUntilPrompt читає, доки не побачить запрошення. +func (c *CLI) readUntilPrompt(ctx context.Context) (string, error) { + chunk := make([]byte, 8192) + + for { + if err := ctx.Err(); err != nil { + return c.buf.String(), err + } + + // Промпт шукаємо лише в хвості: конфіг може містити рядок, що + // схожий на запрошення (banner motd, коментар), і пошук по + // всьому накопиченому обривав би збір на середині. + if loc := c.matchTail(c.opt.PromptRe); loc != nil { + // Ріжемо ВЕСЬ рядок із запрошенням, а не від місця збігу. + // Типовий шаблон «[>#]\s*$» збігається лише з символом + // запрошення, тому обрізання за loc[0] лишало б у конфізі + // ім'я пристрою — «…interface Gi0/1» і окремим рядком «sw1». + out := c.buf.String()[:lineStart(c.buf.String(), loc[0])] + c.buf.Reset() + return out, nil + } + + // Посторінковий вивід: відповідаємо й прибираємо саме + // запрошення з тексту, щоб воно не потрапило в конфіг. + if loc := c.matchTail(c.opt.MoreRe); loc != nil { + s := c.buf.String() + c.buf.Reset() + // Запрошення пейджера часто стоїть у власному рядку разом + // із керуючими послідовностями очищення — прибираємо рядок + // цілком, інакше в конфізі лишаються пробіли й ESC[K. + c.buf.WriteString(s[:lineStart(s, loc[0])]) + if _, err := c.conn.Write([]byte(c.opt.MoreSend)); err != nil { + return c.buf.String(), fmt.Errorf("відповідь на --More--: %w", err) + } + continue + } + + if c.buf.Len() > c.opt.MaxBytes { + return c.buf.String(), ErrTooLarge + } + + n, err := c.readWithDeadline(ctx, chunk) + if n > 0 { + c.buf.Write(chunk[:n]) + if c.opt.Transcript != nil { + c.opt.Transcript.Write(chunk[:n]) + } + } + if err != nil { + if errors.Is(err, errIdle) { + return c.buf.String(), ErrPromptTimeout + } + if errors.Is(err, io.EOF) { + // Пристрій закрив сесію. Якщо промпт так і не + // з'явився — це не успіх, навіть коли текст схожий + // на конфіг. + return c.buf.String(), ErrPromptTimeout + } + return c.buf.String(), err + } + } +} + +// matchTail шукає збіг лише в кінці накопиченого. +func (c *CLI) matchTail(re *regexp.Regexp) []int { + if re == nil { + return nil + } + s := c.buf.String() + // Хвоста в 512 байтів вистачає будь-якому запрошенню й відрізає + // збіги всередині тіла конфігу. + const tail = 512 + off := 0 + if len(s) > tail { + off = len(s) - tail + } + loc := re.FindStringIndex(s[off:]) + if loc == nil { + return nil + } + return []int{off + loc[0], off + loc[1]} +} + +// lineStart повертає позицію початку рядка, у якому стоїть pos. +func lineStart(s string, pos int) int { + if pos <= 0 { + return 0 + } + if i := strings.LastIndexByte(s[:pos], byte('\n')); i >= 0 { + return i + 1 + } + return 0 +} + +var errIdle = errors.New("тиша в каналі") + +// readWithDeadline читає з обмеженням на паузу між байтами. +// +// Саме на паузу, а не на всю операцію: збір конфігу великого шасі +// триває хвилини, і загальний дедлайн довелося б ставити навмання. +// Тиша ж означає одне з двох — пристрій завис або промпт не той. +func (c *CLI) readWithDeadline(ctx context.Context, p []byte) (int, error) { + type res struct { + n int + err error + } + ch := make(chan res, 1) + go func() { + n, err := c.conn.Read(p) + ch <- res{n, err} + }() + + t := time.NewTimer(c.opt.IdleTimeout) + defer t.Stop() + + select { + case r := <-ch: + return r.n, r.err + case <-t.C: + return 0, errIdle + case <-ctx.Done(): + return 0, ctx.Err() + } +} + +// ansiRe ловить керуючі послідовності терміналу. +// +// CSI (кольори, позиціювання, bracketed paste) і OSC (заголовок вікна). +// Без цього в збережений конфіг потрапляє те, що бачить око оператора, +// а не те, що налаштовано на пристрої: перший же diff показує зміну +// там, де змінився лише режим терміналу. +var ansiRe = regexp.MustCompile("\x1b\\[[0-9;?]*[ -/]*[@-~]" + + "|\x1b\\][^\x07\x1b]*(?:\x07|\x1b\\\\)" + + "|\x1b[=>NOc]") + +// crRe згортає підряд ідучі \r перед переводом рядка. +// +// Псевдотермінал видає «\r\r\n» частіше, ніж здається: сам пристрій +// шле «\r\n», а PTY додає свій «\r». Без згортання кожен такий рядок +// подвоюється, і конфіг на 300 байтів приїжджає на 12 рядків замість 10. +var crRe = regexp.MustCompile(`\r+\n`) + +// cleanOutput прибирає відлуння команди, керуючі коди й нормалізує +// переводи рядка. +// +// Відлуння прибирається саме порівнянням із надісланою командою, а не +// «викинути перший рядок»: частина пристроїв відлуння не робить, і +// сліпе відкидання з'їдало б перший рядок конфігу. +func cleanOutput(raw, cmd string) string { + s := ansiRe.ReplaceAllString(raw, "") + s = crRe.ReplaceAllString(s, "\n") + s = strings.ReplaceAll(s, "\r", "\n") + + lines := strings.Split(s, "\n") + + // Відлуння шукаємо серед перших кількох рядків і за суфіксом. + // + // Суфіксом, бо в PTY команда часто друкується в тому ж рядку, що й + // запрошення: «claude@host:~$ cat /etc/os-release». Порівняння + // цілого рядка тут не спрацювало б, і запрошення поїхало б у конфіг. + want := strings.TrimSpace(cmd) + for i := 0; i < len(lines) && i < 3; i++ { + l := strings.TrimSpace(lines[i]) + if l == want || (want != "" && strings.HasSuffix(l, want)) { + lines = lines[i+1:] + break + } + } + + // Хвостові порожні рядки — наслідок того, що ми відрізали промпт. + for len(lines) > 0 && strings.TrimSpace(lines[len(lines)-1]) == "" { + lines = lines[:len(lines)-1] + } + // Провідні порожні — наслідок того, що ми відрізали відлуння. + for len(lines) > 0 && strings.TrimSpace(lines[0]) == "" { + lines = lines[1:] + } + return strings.Join(lines, "\n") +} + +// StripFirstLines відкидає службову шапку. +// +// Скільки саме — знає профіль: у Cisco IOS це три рядки «Building +// configuration…», у Eltex MES5448 — десять. +func StripFirstLines(s string, n int) string { + if n <= 0 { + return s + } + lines := strings.Split(s, "\n") + if n >= len(lines) { + return "" + } + return strings.Join(lines[n:], "\n") +} + +func (c *CLI) Close() error { return c.conn.Close() } diff --git a/agent/internal/ncmx/session_test.go b/agent/internal/ncmx/session_test.go new file mode 100644 index 0000000..6923402 --- /dev/null +++ b/agent/internal/ncmx/session_test.go @@ -0,0 +1,391 @@ +package ncmx + +import ( + "bytes" + "context" + "io" + "regexp" + "strings" + "sync" + "testing" + "time" +) + +// fakeDevice — пристрій, що відповідає за сценарієм. +// +// Саме пристрій, а не мок з'єднання: перевіряти треба поведінку розбору +// на справжньому потоці байтів, включно з тим, що промпт приходить +// розірваним між читаннями. +type fakeDevice struct { + mu sync.Mutex + out bytes.Buffer + replies map[string]string + banner string + prompt string + written []string + chunkLen int + closed bool +} + +func newDevice(prompt, banner string, replies map[string]string) *fakeDevice { + d := &fakeDevice{replies: replies, banner: banner, prompt: prompt} + d.out.WriteString(banner + prompt) + return d +} + +func (d *fakeDevice) Read(p []byte) (int, error) { + deadline := time.Now().Add(2 * time.Second) + for { + d.mu.Lock() + if d.out.Len() > 0 { + n := d.out.Len() + // Віддаємо дрібними шматками: так промпт гарантовано + // розривається між викликами Read, і код мусить це пережити. + if d.chunkLen > 0 && n > d.chunkLen { + n = d.chunkLen + } + if n > len(p) { + n = len(p) + } + nn, _ := d.out.Read(p[:n]) + d.mu.Unlock() + return nn, nil + } + closed := d.closed + d.mu.Unlock() + if closed { + return 0, io.EOF + } + if time.Now().After(deadline) { + return 0, io.EOF + } + time.Sleep(time.Millisecond) + } +} + +func (d *fakeDevice) Write(p []byte) (int, error) { + cmd := strings.TrimRight(string(p), "\r\n") + d.mu.Lock() + defer d.mu.Unlock() + d.written = append(d.written, cmd) + + // Відлуння команди — так поводиться справжня консоль. + d.out.WriteString(cmd + "\r\n") + if reply, ok := d.replies[cmd]; ok { + d.out.WriteString(reply) + } + d.out.WriteString(d.prompt) + return len(p), nil +} + +func (d *fakeDevice) Close() error { + d.mu.Lock() + d.closed = true + d.mu.Unlock() + return nil +} + +func cliFor(d *fakeDevice, prompt string) *CLI { + return NewCLI(d, Options{ + PromptRe: regexp.MustCompile(`(?m)` + prompt), + IdleTimeout: time.Second, + }) +} + +func TestRunStripsEchoAndPrompt(t *testing.T) { + d := newDevice("\nsw1#", "Welcome\r\n", map[string]string{ + "show running-config": "hostname sw1\r\ninterface Gi0/1\r\n", + }) + c := cliFor(d, `[>#]\s*$`) + ctx := context.Background() + + if err := c.WaitPrompt(ctx); err != nil { + t.Fatalf("привітання: %v", err) + } + + out, err := c.Run(ctx, "show running-config") + if err != nil { + t.Fatalf("Run: %v", err) + } + if out != "hostname sw1\ninterface Gi0/1" { + t.Fatalf("отримали %q", out) + } +} + +// Промпт, розірваний між читаннями, — норма для повільного каналу. +func TestRunHandlesSplitPrompt(t *testing.T) { + d := newDevice("\nsw1#", "", map[string]string{ + "show run": "line one\r\nline two\r\n", + }) + d.chunkLen = 3 // рвемо потік по три байти + + c := cliFor(d, `[>#]\s*$`) + ctx := context.Background() + _ = c.WaitPrompt(ctx) + + out, err := c.Run(ctx, "show run") + if err != nil { + t.Fatalf("Run: %v", err) + } + if out != "line one\nline two" { + t.Fatalf("отримали %q", out) + } +} + +// Рядок конфігу, схожий на запрошення, не має обривати збір. +// +// Це не теоретичний випадок: banner motd із «#» усередині трапляється +// в кожній другій мережі. +func TestPromptInsideBodyDoesNotTruncate(t *testing.T) { + body := "banner motd #\r\nDANGER\r\n#\r\nhostname sw1\r\n" + d := newDevice("\nsw1#", "", map[string]string{"show run": body}) + + c := cliFor(d, `sw1[>#]\s*$`) + ctx := context.Background() + _ = c.WaitPrompt(ctx) + + out, err := c.Run(ctx, "show run") + if err != nil { + t.Fatalf("Run: %v", err) + } + if !strings.Contains(out, "hostname sw1") { + t.Fatalf("тіло обрізане: %q", out) + } +} + +// Посторінковий вивід: агент має відповісти й прибрати «--More--» із тексту. +// +// Макет тут керується вручну, а не через таблицю відповідей: справжній +// пристрій, зупинившись на «--More--», НЕ шле запрошення — саме тому +// агент і мусить його розпізнати. Автоматичне дописування промпту +// зробило б перевірку беззмістовною. +func TestMorePromptIsAnsweredAndStripped(t *testing.T) { + d := &scriptedDevice{prompt: "\nsw1#"} + d.push("\nsw1#") // початкове запрошення + + c := NewCLI(d, Options{ + PromptRe: regexp.MustCompile(`(?m)[>#]\s*$`), + IdleTimeout: 2 * time.Second, + }) + ctx := context.Background() + if err := c.WaitPrompt(ctx); err != nil { + t.Fatalf("привітання: %v", err) + } + + // Пристрій віддає першу сторінку й зупиняється на пейджері. + d.onWrite("show run", "show run\r\npart one\r\n --More-- ") + // На пробіл — решта й нарешті запрошення. + d.onWrite(" ", "\rpart two\r\nsw1#") + + out, err := c.Run(ctx, "show run") + if err != nil { + t.Fatalf("Run: %v", err) + } + if strings.Contains(out, "More") { + t.Fatalf("запрошення пейджера потрапило в конфіг: %q", out) + } + if !strings.Contains(out, "part one") || !strings.Contains(out, "part two") { + t.Fatalf("вивід зібрано не повністю: %q", out) + } + if got := d.sent(); len(got) < 2 || got[1] != " " { + t.Fatalf("агент не відповів пробілом на пейджер: %v", got) + } +} + +// scriptedDevice відповідає рівно тим, що йому прописали на конкретний +// запис. Нічого не додає від себе — на відміну від fakeDevice, який +// імітує звичайну консоль. +type scriptedDevice struct { + mu sync.Mutex + out bytes.Buffer + script map[string]string + written []string + prompt string + closed bool +} + +func (d *scriptedDevice) push(s string) { + d.mu.Lock() + d.out.WriteString(s) + d.mu.Unlock() +} + +func (d *scriptedDevice) onWrite(in, reply string) { + d.mu.Lock() + if d.script == nil { + d.script = map[string]string{} + } + d.script[in] = reply + d.mu.Unlock() +} + +func (d *scriptedDevice) sent() []string { + d.mu.Lock() + defer d.mu.Unlock() + return append([]string(nil), d.written...) +} + +func (d *scriptedDevice) Read(p []byte) (int, error) { + deadline := time.Now().Add(2 * time.Second) + for { + d.mu.Lock() + if d.out.Len() > 0 { + n, _ := d.out.Read(p) + d.mu.Unlock() + return n, nil + } + closed := d.closed + d.mu.Unlock() + if closed || time.Now().After(deadline) { + return 0, io.EOF + } + time.Sleep(time.Millisecond) + } +} + +func (d *scriptedDevice) Write(p []byte) (int, error) { + // Пейджеру відповідають голим пробілом без переводу рядка, тому + // ключем є те, що надіслали, як є. + key := strings.TrimRight(string(p), "\r\n") + if key == "" { + key = string(p) + } + + d.mu.Lock() + d.written = append(d.written, key) + if reply, ok := d.script[key]; ok { + d.out.WriteString(reply) + } + d.mu.Unlock() + return len(p), nil +} + +func (d *scriptedDevice) Close() error { + d.mu.Lock() + d.closed = true + d.mu.Unlock() + return nil +} + +// Тиша в каналі має давати саме ErrPromptTimeout: це підказує, що +// проблема в prompt_regex, а не в мережі. +func TestSilenceGivesPromptTimeout(t *testing.T) { + d := newDevice("\nsw1#", "", nil) + c := NewCLI(d, Options{ + PromptRe: regexp.MustCompile(`(?m)ЦЬОГО-НЕ-БУДЕ$`), + IdleTimeout: 150 * time.Millisecond, + }) + + _, err := c.readUntilPrompt(context.Background()) + if err != ErrPromptTimeout { + t.Fatalf("очікували ErrPromptTimeout, отримали %v", err) + } +} + +func TestMaxBytesGuard(t *testing.T) { + huge := strings.Repeat("x", 40000) + "\r\n" + d := newDevice("\nsw1#", "", map[string]string{"show run": huge}) + + c := NewCLI(d, Options{ + PromptRe: regexp.MustCompile(`(?m)[>#]\s*$`), + IdleTimeout: time.Second, + MaxBytes: 1024, + }) + ctx := context.Background() + _ = c.WaitPrompt(ctx) + + if _, err := c.Run(ctx, "show run"); err != ErrTooLarge { + t.Fatalf("очікували ErrTooLarge, отримали %v", err) + } +} + +// Пристрій без відлуння не має втрачати перший рядок конфігу. +func TestNoEchoKeepsFirstLine(t *testing.T) { + raw := "hostname sw1\ninterface Gi0/1" + if got := cleanOutput(raw, "show run"); got != raw { + t.Fatalf("перший рядок з'їдено: %q", got) + } +} + +func TestEchoIsRemovedOnce(t *testing.T) { + raw := "show run\nhostname sw1\nshow run\n" + got := cleanOutput(raw, "show run") + if got != "hostname sw1\nshow run" { + t.Fatalf("отримали %q", got) + } +} + +func TestStripFirstLines(t *testing.T) { + s := "Building configuration...\n\nCurrent configuration : 100 bytes\nhostname sw1" + if got := StripFirstLines(s, 3); got != "hostname sw1" { + t.Fatalf("отримали %q", got) + } + if got := StripFirstLines(s, 0); got != s { + t.Fatal("нуль рядків має лишати текст як є") + } + if got := StripFirstLines("одне", 5); got != "" { + t.Fatalf("забагато рядків: %q", got) + } +} + +// Порожній prompt_regex не має валити збір: підставляємо типовий. +func TestEmptyPromptFallsBack(t *testing.T) { + re, err := compilePrompt("") + if err != nil { + t.Fatalf("compilePrompt: %v", err) + } + if !re.MatchString("sw1#") { + t.Fatal("типовий шаблон не ловить «sw1#»") + } +} + +func TestBadPromptIsReported(t *testing.T) { + if _, err := compilePrompt("[unclosed"); err == nil { + t.Fatal("некоректний регекс мав дати помилку") + } +} + +// Багаторядковий режим обов'язковий: без нього «$» шукає кінець усього +// тексту, а запрошення стоїть у кінці останнього рядка. +func TestPromptIsMultiline(t *testing.T) { + re, _ := compilePrompt(`[>#]\s*$`) + if !re.MatchString("рядок\nще рядок\nsw1#") { + t.Fatal("шаблон не працює в багаторядковому тексті") + } +} + +// Керуючі коди терміналу не мають потрапляти в конфіг: перший же diff +// показав би зміну там, де змінився лише режим терміналу. +func TestAnsiSequencesAreStripped(t *testing.T) { + raw := "\x1b[?2004lhostname sw1\x1b[0m\r\ninterface Gi0/1\r\n" + got := cleanOutput(raw, "show run") + if strings.ContainsRune(got, 0x1b) { + t.Fatalf("escape-послідовність лишилась: %q", got) + } + if got != "hostname sw1\ninterface Gi0/1" { + t.Fatalf("отримали %q", got) + } +} + +// Псевдотермінал видає «\r\r\n»: сам пристрій шле «\r\n», PTY додає свій +// «\r». Без згортання кожен рядок подвоювався б. +func TestDoubleCarriageReturnDoesNotDoubleLines(t *testing.T) { + raw := "line one\r\r\nline two\r\r\n" + got := cleanOutput(raw, "show run") + if got != "line one\nline two" { + t.Fatalf("отримали %q (рядків %d)", got, strings.Count(got, "\n")+1) + } +} + +// Відлуння команди трапляється не першим рядком, а після залишку +// запрошення — саме так поводиться bash у PTY. +func TestEchoAfterPromptRemnantIsRemoved(t *testing.T) { + raw := "claude@host:~$ cat /etc/os-release\r\nID=debian\r\n" + got := cleanOutput(raw, "cat /etc/os-release") + if strings.Contains(got, "os-release") { + t.Fatalf("відлуння лишилось: %q", got) + } + if got != "ID=debian" { + t.Fatalf("отримали %q", got) + } +} diff --git a/agent/internal/ncmx/transport.go b/agent/internal/ncmx/transport.go new file mode 100644 index 0000000..7c8e4f9 --- /dev/null +++ b/agent/internal/ncmx/transport.go @@ -0,0 +1,365 @@ +package ncmx + +import ( + "context" + "errors" + "fmt" + "io" + "net" + "strconv" + "time" + + "golang.org/x/crypto/ssh" +) + +// Dial відкриває інтерактивну сесію потрібним транспортом. +func Dial(ctx context.Context, transport, host string, port int, + user, password string, timeout time.Duration) (Conn, error) { + + switch transport { + case "ssh", "": + if port == 0 { + port = 22 + } + return dialSSH(ctx, host, port, user, password, timeout) + case "telnet": + if port == 0 { + port = 23 + } + return dialTelnet(ctx, host, port, user, password, timeout) + default: + return nil, fmt.Errorf("непідтримуваний транспорт %q", transport) + } +} + +// --------------------------------------------------------------------- +// SSH +// --------------------------------------------------------------------- + +type sshConn struct { + client *ssh.Client + session *ssh.Session + stdin io.WriteCloser + stdout io.Reader +} + +func dialSSH(ctx context.Context, host string, port int, + user, password string, timeout time.Duration) (Conn, error) { + + cfg := &ssh.ClientConfig{ + User: user, + Auth: []ssh.AuthMethod{ + ssh.Password(password), + // Частина старих платформ не вміє «password», лише + // keyboard-interactive з єдиним запитом. + ssh.KeyboardInteractive(func(_, _ string, qs []string, _ []bool) ([]string, error) { + ans := make([]string, len(qs)) + for i := range qs { + ans[i] = password + } + return ans, nil + }), + }, + // Ключі мережевого обладнання не звіряються. + // + // Це свідоме рішення, а не недогляд. Зонд стоїть усередині + // мережі клієнта й ходить до заліза, яке перегенеровує ключ + // після кожної заміни прошивки чи скидання. Тримати known_hosts + // на сотні пристроїв означало б або вимикати перевірку + // вручну щотижня, або не збирати конфіги зовсім. + // + // Захист від підміни тут дає сегмент керування, а не TOFU. + HostKeyCallback: ssh.InsecureIgnoreHostKey(), + Timeout: timeout, + // Старе залізо не знає сучасних алгоритмів. Перелік свідомо + // широкий: інакше половина парку відпаде з «no common + // algorithm», і збір конфігу для неї стане неможливим. + Config: ssh.Config{ + KeyExchanges: []string{ + "curve25519-sha256", "curve25519-sha256@libssh.org", + "ecdh-sha2-nistp256", "ecdh-sha2-nistp384", "ecdh-sha2-nistp521", + "diffie-hellman-group14-sha256", "diffie-hellman-group14-sha1", + "diffie-hellman-group1-sha1", "diffie-hellman-group-exchange-sha1", + }, + Ciphers: []string{ + "aes128-gcm@openssh.com", "aes256-gcm@openssh.com", + "aes128-ctr", "aes192-ctr", "aes256-ctr", + "aes128-cbc", "3des-cbc", + }, + }, + HostKeyAlgorithms: []string{ + "ssh-ed25519", "ecdsa-sha2-nistp256", "rsa-sha2-512", + "rsa-sha2-256", "ssh-rsa", + }, + } + + addr := net.JoinHostPort(host, strconv.Itoa(port)) + + d := net.Dialer{Timeout: timeout} + raw, err := d.DialContext(ctx, "tcp", addr) + if err != nil { + return nil, fmt.Errorf("зʼєднання з %s: %w", addr, err) + } + + // Рукостискання не має права висіти вічно: ctx керує дедлайном, + // але сам ssh.NewClientConn його не бачить. + if dl, ok := ctx.Deadline(); ok { + _ = raw.SetDeadline(dl) + } else { + _ = raw.SetDeadline(time.Now().Add(timeout)) + } + + c, chans, reqs, err := ssh.NewClientConn(raw, addr, cfg) + if err != nil { + raw.Close() + return nil, fmt.Errorf("SSH-рукостискання з %s: %w", addr, err) + } + _ = raw.SetDeadline(time.Time{}) + + client := ssh.NewClient(c, chans, reqs) + sess, err := client.NewSession() + if err != nil { + client.Close() + return nil, fmt.Errorf("сесія SSH: %w", err) + } + + // Псевдотермінал обовʼязковий: без нього більшість мережевих + // платформ або не дає інтерактивної оболонки, або віддає вивід у + // режимі, де немає запрошення — а саме за ним ми й орієнтуємось. + // + // 0 рядків заввишки — ще один спосіб попросити «не гортай»; команду + // вимкнення пейджера це не скасовує, бо частина платформ ігнорує + // розмір вікна. + modes := ssh.TerminalModes{ssh.ECHO: 0, ssh.TTY_OP_ISPEED: 38400, ssh.TTY_OP_OSPEED: 38400} + if err := sess.RequestPty("vt100", 0, 200, modes); err != nil { + sess.Close() + client.Close() + return nil, fmt.Errorf("PTY: %w", err) + } + + stdin, err := sess.StdinPipe() + if err != nil { + sess.Close() + client.Close() + return nil, err + } + stdout, err := sess.StdoutPipe() + if err != nil { + sess.Close() + client.Close() + return nil, err + } + // Помилки читаємо тим самим потоком: пристрої часто пишуть + // «% Invalid input» саме в stderr, і загубити це означало б + // діагностувати наосліп. + sess.Stderr = writerTo(stdin) + + if err := sess.Shell(); err != nil { + sess.Close() + client.Close() + return nil, fmt.Errorf("оболонка: %w", err) + } + + return &sshConn{client: client, session: sess, stdin: stdin, stdout: stdout}, nil +} + +func (c *sshConn) Read(p []byte) (int, error) { return c.stdout.Read(p) } +func (c *sshConn) Write(p []byte) (int, error) { return c.stdin.Write(p) } + +func (c *sshConn) Close() error { + // Порядок має значення: спершу сесія, потім клієнт. Навпаки — + // і Close() сесії повисне на мертвому з'єднанні. + if c.session != nil { + _ = c.session.Close() + } + return c.client.Close() +} + +// writerTo — заглушка, щоб stderr не губився. Пишемо туди ж, куди й +// stdout читається, тобто нікуди: справжнє злиття потоків дав би +// sess.Stderr = sess.Stdout, але API так не вміє. +type nopWriter struct{} + +func (nopWriter) Write(p []byte) (int, error) { return len(p), nil } + +func writerTo(io.Writer) io.Writer { return nopWriter{} } + +// --------------------------------------------------------------------- +// Telnet +// --------------------------------------------------------------------- + +// Команди протоколу Telnet (RFC 854). +const ( + iac = 255 // Interpret As Command + dont = 254 + do = 253 + wont = 252 + will = 251 + sb = 250 // subnegotiation begin + se = 240 // subnegotiation end +) + +type telnetConn struct { + conn net.Conn + // leftover — байти, що лишились після зняття команд протоколу. + leftover []byte +} + +func dialTelnet(ctx context.Context, host string, port int, + user, password string, timeout time.Duration) (Conn, error) { + + addr := net.JoinHostPort(host, strconv.Itoa(port)) + d := net.Dialer{Timeout: timeout} + raw, err := d.DialContext(ctx, "tcp", addr) + if err != nil { + return nil, fmt.Errorf("зʼєднання з %s: %w", addr, err) + } + + tc := &telnetConn{conn: raw} + + // Вхід у Telnet — це не протокол, а розмова: пристрій просто пише + // «Username:» і чекає. Кожен вендор пише по-своєму, тому шукаємо + // підрядок без урахування регістру. + if user != "" { + if err := tc.expectAndSend(ctx, []string{"username", "login", "user name"}, user, timeout); err != nil { + raw.Close() + return nil, fmt.Errorf("логін: %w", err) + } + } + if password != "" { + if err := tc.expectAndSend(ctx, []string{"password", "passwd"}, password, timeout); err != nil { + raw.Close() + return nil, fmt.Errorf("пароль: %w", err) + } + } + return tc, nil +} + +func (c *telnetConn) expectAndSend(ctx context.Context, want []string, send string, timeout time.Duration) error { + deadline := time.Now().Add(timeout) + var acc []byte + buf := make([]byte, 4096) + + for time.Now().Before(deadline) { + if err := ctx.Err(); err != nil { + return err + } + _ = c.conn.SetReadDeadline(time.Now().Add(2 * time.Second)) + n, err := c.Read(buf) + if n > 0 { + acc = append(acc, buf[:n]...) + low := lower(acc) + for _, w := range want { + if bytesContains(low, w) { + _, werr := c.Write([]byte(send + "\n")) + return werr + } + } + } + if err != nil { + var ne net.Error + if errors.As(err, &ne) && ne.Timeout() { + continue + } + return err + } + } + return fmt.Errorf("не дочекались запиту %v", want) +} + +// Read знімає команди протоколу з потоку. +// +// Telnet перемішує керування з даними, і якщо не вирізати IAC-послідовності, +// вони потраплять просто в текст конфігу — саме звідти беруться +// «сміттєві» байти в збережених конфігах. +func (c *telnetConn) Read(p []byte) (int, error) { + if len(c.leftover) > 0 { + n := copy(p, c.leftover) + c.leftover = c.leftover[n:] + return n, nil + } + + raw := make([]byte, len(p)) + n, err := c.conn.Read(raw) + if n == 0 { + return 0, err + } + + clean := make([]byte, 0, n) + for i := 0; i < n; { + if raw[i] != iac { + clean = append(clean, raw[i]) + i++ + continue + } + // IAC IAC — екранований 0xFF, справжній байт даних. + if i+1 < n && raw[i+1] == iac { + clean = append(clean, iac) + i += 2 + continue + } + if i+2 < n && (raw[i+1] == do || raw[i+1] == dont || + raw[i+1] == will || raw[i+1] == wont) { + // Відмовляємось від усього: жодна опція нам не потрібна, + // а мовчання частина пристроїв тлумачить як згоду й чекає + // підтвердження вічно. + resp := byte(wont) + if raw[i+1] == will || raw[i+1] == wont { + resp = dont + } + _, _ = c.conn.Write([]byte{iac, resp, raw[i+2]}) + i += 3 + continue + } + if i+1 < n && raw[i+1] == sb { + // Підпереговори тягнуться до IAC SE. + j := i + 2 + for j+1 < n && !(raw[j] == iac && raw[j+1] == se) { + j++ + } + i = j + 2 + continue + } + // Обрізана послідовність на межі читання — лишаємо на потім. + c.leftover = append(c.leftover, raw[i:n]...) + break + } + + copied := copy(p, clean) + if copied < len(clean) { + c.leftover = append(clean[copied:], c.leftover...) + } + return copied, err +} + +func (c *telnetConn) Write(p []byte) (int, error) { + _ = c.conn.SetWriteDeadline(time.Now().Add(10 * time.Second)) + return c.conn.Write(p) +} + +func (c *telnetConn) Close() error { return c.conn.Close() } + +func lower(b []byte) []byte { + out := make([]byte, len(b)) + for i, c := range b { + if c >= 'A' && c <= 'Z' { + c += 'a' - 'A' + } + out[i] = c + } + return out +} + +func bytesContains(hay []byte, needle string) bool { + return len(needle) > 0 && indexOf(hay, needle) >= 0 +} + +func indexOf(hay []byte, needle string) int { + n := len(needle) + for i := 0; i+n <= len(hay); i++ { + if string(hay[i:i+n]) == needle { + return i + } + } + return -1 +} diff --git a/agent/internal/session/config_jobs.go b/agent/internal/session/config_jobs.go new file mode 100644 index 0000000..b717115 --- /dev/null +++ b/agent/internal/session/config_jobs.go @@ -0,0 +1,180 @@ +package session + +import ( + "bytes" + "compress/gzip" + "context" + "errors" + "time" + + "github.com/netpulse/netpulse/agent/internal/ncmx" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" + "google.golang.org/protobuf/types/known/durationpb" + "google.golang.org/protobuf/types/known/timestamppb" +) + +// errNoConnection — сесія обірвалась, поки збирали конфіг. +var errNoConnection = errors.New("немає живого зʼєднання з сервером") + +// chunkSize — розмір шматка вивантаження. +// +// 64 КБ: помітно менше за типовий ліміт gRPC-повідомлення, тому навіть +// конфіг на кілька МБ проходить без налаштування транспорту, і при +// цьому не так дрібно, щоб платити накладними за кожен рядок. +const chunkSize = 64 << 10 + +// runConfigJob знімає конфіг і вивантажує його на сервер. +// +// Виконується в окремій горутині: збір конфігу з великого шасі триває +// хвилини, і тримати на ньому контрольний цикл означало б не відповідати +// на ping і бути визнаним мертвим саме тоді, коли агент найбільше +// зайнятий корисною роботою. +func (s *Session) runConfigJob(ctx context.Context, job *npv1.ConfigJob) { + log := s.log.With("job_id", job.GetJobId(), + "device", job.GetDevice().GetName(), + "config_type", job.GetConfigType()) + + log.Info("збір конфігу почався") + + res, err := ncmx.Collect(ctx, job) + if err != nil { + log.Error("збір конфігу", "помилка", err) + s.uploadFailure(ctx, job, err, res.Transcript) + return + } + + log.Info("конфіг знято", + "байтів", len(res.Body), "рядків", res.LineCount, + "тривалість", res.Duration.Round(time.Millisecond)) + + if err := s.uploadConfig(ctx, job, res); err != nil { + log.Error("вивантаження конфігу", "помилка", err) + } +} + +func (s *Session) uploadConfig(ctx context.Context, job *npv1.ConfigJob, res ncmx.Result) error { + conn := s.client.Load() + if conn == nil { + return errNoConnection + } + + stream, err := (*conn).UploadConfig(ctx) + if err != nil { + return err + } + + // Стискаємо завжди: конфіги — це текст із величезною надлишковістю, + // а канал до зонда часто вузький. Сервер розуміє обидва варіанти, + // тому вибір тут наш. + var zbuf bytes.Buffer + zw := gzip.NewWriter(&zbuf) + if _, err := zw.Write(res.Body); err != nil { + return err + } + if err := zw.Close(); err != nil { + return err + } + payload := zbuf.Bytes() + + if err := stream.Send(&npv1.ConfigUpload{ + Part: &npv1.ConfigUpload_Header{Header: &npv1.ConfigHeader{ + JobId: job.GetJobId(), + AgentId: s.cfg.AgentID, + DeviceId: job.GetDevice().GetDeviceId(), + ConfigType: job.GetConfigType(), + CollectedAt: timestamppb.Now(), + Encoding: "gzip", + }}, + }); err != nil { + return err + } + + var chunks uint32 + for off := 0; off < len(payload); off += chunkSize { + end := off + chunkSize + if end > len(payload) { + end = len(payload) + } + if err := stream.Send(&npv1.ConfigUpload{ + Part: &npv1.ConfigUpload_Chunk{Chunk: &npv1.ConfigChunk{ + Sequence: chunks, + Data: payload[off:end], + }}, + }); err != nil { + return err + } + chunks++ + } + + // sha256 рахується від ОРИГІНАЛУ, до стиснення: сервер звіряє саме + // тіло конфігу, а не його пакування. Інакше зміна рівня стиснення + // виглядала б як зміна конфігу. + if err := stream.Send(&npv1.ConfigUpload{ + Part: &npv1.ConfigUpload_Trailer{Trailer: &npv1.ConfigTrailer{ + Success: true, + ContentSha256: res.SHA256, + SizeBytes: uint64(len(res.Body)), + LineCount: uint32(res.LineCount), + ChunkCount: chunks, + Duration: durationpb.New(res.Duration), + Transcript: res.Transcript, + }}, + }); err != nil { + return err + } + + receipt, err := stream.CloseAndRecv() + if err != nil { + return err + } + if !receipt.GetAccepted() { + s.log.Warn("сервер не прийняв конфіг", + "job_id", job.GetJobId(), "причина", receipt.GetError().GetMessage()) + return nil + } + + s.log.Info("конфіг прийнято", + "job_id", job.GetJobId(), + "змінився", !receipt.GetUnchanged(), + "commit", receipt.GetCommitSha()) + return nil +} + +// uploadFailure повідомляє про невдачу тим самим стрімом. +// +// Мовчати не можна: без цього невдалий збір виглядає на сервері як +// «завдання ще виконується», і наступний запуск за розкладом лише +// додасть другий такий самий. Транскрипт долучається саме тут — він +// потрібен рівно тоді, коли щось пішло не так. +func (s *Session) uploadFailure(ctx context.Context, job *npv1.ConfigJob, cause error, transcript string) { + conn := s.client.Load() + if conn == nil { + return + } + stream, err := (*conn).UploadConfig(ctx) + if err != nil { + return + } + + _ = stream.Send(&npv1.ConfigUpload{ + Part: &npv1.ConfigUpload_Header{Header: &npv1.ConfigHeader{ + JobId: job.GetJobId(), + AgentId: s.cfg.AgentID, + DeviceId: job.GetDevice().GetDeviceId(), + ConfigType: job.GetConfigType(), + CollectedAt: timestamppb.Now(), + Encoding: "none", + }}, + }) + _ = stream.Send(&npv1.ConfigUpload{ + Part: &npv1.ConfigUpload_Trailer{Trailer: &npv1.ConfigTrailer{ + Success: false, + Error: &npv1.Error{ + Code: "collect_failed", + Message: cause.Error(), + }, + Transcript: transcript, + }}, + }) + _, _ = stream.CloseAndRecv() +} diff --git a/agent/internal/session/session.go b/agent/internal/session/session.go index 1be9f1c..4b08d35 100644 --- a/agent/internal/session/session.go +++ b/agent/internal/session/session.go @@ -74,6 +74,14 @@ type Session struct { devMu sync.RWMutex devices map[string]*npv1.DeviceTarget + // client живий лише в межах сесії. Збір конфігу триває хвилинами й + // не має тримати контрольний цикл, тому вивантаження йде окремою + // горутиною — а їй потрібен доступ до клієнта поточної сесії. + client atomic.Pointer[npv1.AgentServiceClient] + // jobs рахує незавершені збори: при обриві сесії їх треба дочекатись, + // інакше вивантаження піде в уже закритий канал. + jobs sync.WaitGroup + planHash atomic.Pointer[[]byte] nextBatch atomic.Uint64 lastAcked atomic.Uint64 @@ -222,6 +230,8 @@ func (s *Session) runOnce(ctx context.Context) error { defer conn.Close() client := npv1.NewAgentServiceClient(conn) + s.client.Store(&client) + defer s.client.Store(nil) sctx, cancel := context.WithCancel(ctx) defer cancel() @@ -284,6 +294,11 @@ func (s *Session) runOnce(ctx context.Context) error { readErr := s.controlLoop(sctx, ctrl, out) cancel() wg.Wait() + // Збір конфігу переживає скасування контексту не довше за свій + // дедлайн, але закривати з'єднання під ним не можна: вивантаження + // впаде на середині, і сервер отримає обірваний стрім замість + // чесної помилки. + s.jobs.Wait() close(errCh) if readErr != nil && !errors.Is(readErr, context.Canceled) { @@ -454,6 +469,14 @@ func (s *Session) controlLoop(ctx context.Context, ctrl npv1.AgentService_Contro s.log.Info("сервер попросив запустити автовиявлення", "run_id", p.DiscoveryRequest.GetRunId(), "задач_зрушено", n) + case *npv1.ControlDown_ConfigJob: + job := p.ConfigJob + s.jobs.Add(1) + go func() { + defer s.jobs.Done() + s.runConfigJob(ctx, job) + }() + case *npv1.ControlDown_Directive: if stop := s.applyDirective(p.Directive); stop { return nil diff --git a/server/API.md b/server/API.md index 08aa02a..6dc7ed9 100644 --- a/server/API.md +++ b/server/API.md @@ -520,10 +520,32 @@ Qtech, ZTE, BDCOM та інші. Тенант може завести власний профіль (`ncm.profiles` із заповненим `tenant_id`) — вбудовані при цьому лишаються недоторканими. -**Виконувати ці профілі поки нікому.** Агентського модуля `ncm` -(SSH/Telnet) ще немає — це Етап 7. Профілі лежать готовими даними, і -щойно модуль з'явиться, збір конфігу запрацює на всьому парку без -дописування коду під кожен вендор. +### Збір конфігу + +| Метод | Шлях | Призначення | +|-------|------|-------------| +| `POST` | `/api/v1/devices/{id}/collect-config` | зібрати зараз (`ncm:write`) | +| `GET` | `/api/v1/devices/{id}/config-jobs` | історія збору | + +Ланцюг такий: REST кладе рядок у `ncm.jobs` зі станом `queued` → +диспетчер усередині AgentService забирає його, якщо потрібний зонд на +зв'язку, і штовхає `ConfigJob` у живу сесію → зонд заходить по +SSH/Telnet, виконує команди профілю й вивантажує результат стрімом → +сервер звіряє sha256, дедуплікує за `content_hash` і закриває завдання. + +**Черга в БД, а не прямий виклик**, бо REST і AgentService — різні +процеси, і живу сесію зонда тримає лише другий. Черга робить передачу +явною й переживає перезапуск обох. Вибірка йде під `FOR UPDATE SKIP +LOCKED`: два екземпляри AgentService не надішлють одне завдання двічі. + +Повторний збір незміненого конфігу дає статус `unchanged` — пристрій +опитано, конфіг звірено, нового коміту не потрібно. Це успіх, а не +відсутність результату, і окремий статус потрібен, щоб у журналі було +видно, коли конфіг востаннє **справді** мінявся. + +Зависле в `running` завдання повертається у відмову за десять хвилин: +зонд міг зникнути разом із ним, і без цього хост лишився б без бекапів +назавжди. ## Групи й доступ до хостів diff --git a/server/internal/grpcapi/ncm_dispatch.go b/server/internal/grpcapi/ncm_dispatch.go new file mode 100644 index 0000000..127b09a --- /dev/null +++ b/server/internal/grpcapi/ncm_dispatch.go @@ -0,0 +1,86 @@ +package grpcapi + +import ( + "context" + "time" + + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" +) + +// DispatchConfigJobs роздає завдання збору конфігу живим сесіям. +// +// Опитування таблиці, а не сповіщення: чергу наповнює інший процес +// (REST), і LISTEN/NOTIFY тут дав би доставку «здебільшого» — воно не +// переживає перезапуск слухача. Ціна опитування — один дешевий запит за +// індексом раз на кілька секунд, і лише тоді, коли хоч один зонд +// на зв'язку. +func (s *Service) DispatchConfigJobs(ctx context.Context, every time.Duration) { + if every <= 0 { + every = 5 * time.Second + } + t := time.NewTicker(every) + defer t.Stop() + + // Зависле завдання повертається у відмову за десять хвилин: зонд міг + // зникнути разом із ним, і без цього хост лишився б без бекапів + // назавжди. + reap := time.NewTicker(2 * time.Minute) + defer reap.Stop() + + s.log.Info("диспетчер збору конфігів запущено", "інтервал", every) + + for { + select { + case <-ctx.Done(): + return + + case <-reap.C: + if n, err := s.store.ReapStuckJobs(ctx, 10*time.Minute); err != nil { + s.log.Warn("прибирання завислих завдань", "err", err) + } else if n > 0 { + s.log.Warn("завдання збору зависли й позначені як невдалі", "рядків", n) + } + + case <-t.C: + online := s.onlineAgentIDs() + if len(online) == 0 { + continue + } + + jobs, err := s.store.ClaimConfigJobs(ctx, online, 16, s.ring) + if err != nil { + s.log.Error("вибірка завдань збору", "err", err) + continue + } + + for _, j := range jobs { + ok := s.PushToAgent(j.AgentID, &npv1.ControlDown{ + Payload: &npv1.ControlDown_ConfigJob{ConfigJob: j.Job}, + }) + if !ok { + // Сесія обірвалась між вибіркою й відправкою. + // Повертаємо в чергу, а не втрачаємо: наступний тік + // віддасть завдання, щойно зонд повернеться. + _ = s.store.FinishConfigJob(ctx, j.JobID, "failed", + "зонд відключився до надсилання завдання", "") + continue + } + s.log.Info("завдання збору надіслано", + "job_id", j.JobID, "agent", j.AgentID, + "device", j.Job.GetDevice().GetName()) + } + } + } +} + +// onlineAgentIDs — хто зараз на зв'язку. +func (s *Service) onlineAgentIDs() []string { + s.mu.RLock() + defer s.mu.RUnlock() + + out := make([]string, 0, len(s.sessions)) + for id := range s.sessions { + out = append(out, id) + } + return out +} diff --git a/server/internal/grpcapi/streams.go b/server/internal/grpcapi/streams.go index 65204aa..246e661 100644 --- a/server/internal/grpcapi/streams.go +++ b/server/internal/grpcapi/streams.go @@ -1,8 +1,11 @@ package grpcapi import ( + "bytes" + "compress/gzip" "context" "errors" + "fmt" "io" "github.com/netpulse/netpulse/server/internal/store" @@ -290,11 +293,37 @@ func (s *Service) UploadConfig(stream npv1.AgentService_UploadConfigServer) erro s.log.Warn("зонд не зміг зібрати конфіг", "agent", agent.ID, "job", header.GetJobId(), "err", tr.GetError().GetMessage()) + // Завдання треба закрити тут: інакше воно висітиме в + // 'running' до прибиральника, а хост увесь цей час + // вважатиметься таким, що збирається. + _ = s.store.FinishConfigJob(ctx, header.GetJobId(), "failed", + tr.GetError().GetMessage(), tr.GetTranscript()) return stream.SendAndClose(&npv1.ConfigReceipt{ JobId: header.GetJobId(), Accepted: false, Error: tr.GetError(), }) } + // Розтискаємо ДО звірки суми: sha256 рахується від тіла + // конфігу, а не від його пакування. Інакше зміна рівня + // стиснення виглядала б як зміна конфігу, а gzip-байти + // потрапили б у Git замість тексту. + plain, err := decodeBody(body, header.GetEncoding()) + if err != nil { + s.log.Warn("не вдалося розпакувати конфіг", + "agent", agent.ID, "encoding", header.GetEncoding(), "err", err) + _ = s.store.FinishConfigJob(ctx, header.GetJobId(), "failed", err.Error(), + tr.GetTranscript()) + return stream.SendAndClose(&npv1.ConfigReceipt{ + JobId: header.GetJobId(), + Accepted: false, + Error: &npv1.Error{ + Code: "bad_encoding", + Message: err.Error(), + }, + }) + } + body = plain + outcome, err := s.store.StoreConfig(ctx, agent, store.ConfigSubmission{ JobID: header.GetJobId(), DeviceID: header.GetDeviceId(), @@ -308,6 +337,8 @@ func (s *Service) UploadConfig(stream npv1.AgentService_UploadConfigServer) erro if errors.Is(outcome.Err, store.ErrChecksumMismatch) { s.log.Warn("конфіг із розбіжною контрольною сумою відхилено", "agent", agent.ID, "device", header.GetDeviceId()) + _ = s.store.FinishConfigJob(ctx, header.GetJobId(), "failed", + "тіло не відповідає заявленому sha256", tr.GetTranscript()) return stream.SendAndClose(&npv1.ConfigReceipt{ JobId: header.GetJobId(), Accepted: false, @@ -322,6 +353,16 @@ func (s *Service) UploadConfig(stream npv1.AgentService_UploadConfigServer) erro "agent", agent.ID, "device", header.GetDeviceId(), "розмір", len(body), "без_змін", outcome.Unchanged) + // «unchanged» — це успіх, а не відсутність результату: + // пристрій опитано, конфіг звірено, нового коміту просто + // не потрібно. Окремий статус потрібен, щоб у журналі було + // видно, коли конфіг востаннє СПРАВДІ мінявся. + finalStatus := "success" + if outcome.Unchanged { + finalStatus = "unchanged" + } + _ = s.store.FinishConfigJob(ctx, header.GetJobId(), finalStatus, "", tr.GetTranscript()) + return stream.SendAndClose(&npv1.ConfigReceipt{ JobId: header.GetJobId(), Accepted: outcome.Accepted, @@ -332,3 +373,33 @@ func (s *Service) UploadConfig(stream npv1.AgentService_UploadConfigServer) erro } } } + +// decodeBody розпаковує тіло конфігу за оголошеним кодуванням. +// +// Порожнє значення й "none" означають одне й те саме: агенти старших +// версій поля не заповнювали, і вимагати його зараз означало б +// відхиляти їхні бекапи. +func decodeBody(body []byte, encoding string) ([]byte, error) { + switch encoding { + case "", "none": + return body, nil + case "gzip": + zr, err := gzip.NewReader(bytes.NewReader(body)) + if err != nil { + return nil, fmt.Errorf("gzip: %w", err) + } + defer zr.Close() + // Ліміт той самий, що й на прийом: розпакування — улюблений + // спосіб перетворити мегабайт трафіку на гігабайт пам'яті. + plain, err := io.ReadAll(io.LimitReader(zr, maxConfigBytes)) + if err != nil { + return nil, fmt.Errorf("gzip: %w", err) + } + return plain, nil + default: + return nil, fmt.Errorf("невідоме кодування %q", encoding) + } +} + +// maxConfigBytes — стеля на розпакований конфіг. +const maxConfigBytes = 64 << 20 diff --git a/server/internal/httpapi/checks.go b/server/internal/httpapi/checks.go index a781c91..e2e2618 100644 --- a/server/internal/httpapi/checks.go +++ b/server/internal/httpapi/checks.go @@ -178,3 +178,56 @@ func (s *Server) handleCreateCredential(w http.ResponseWriter, r *http.Request, } writeJSON(w, http.StatusCreated, map[string]any{"id": id}) } + +// --------------------------------------------------------------------- +// Збір конфігу +// --------------------------------------------------------------------- + +// handleCollectConfig ставить збір у чергу. +// +// Саме в чергу, а не одразу зонду: живу сесію тримає інший процес +// (AgentService), і черга в БД — єдиний чесний спосіб передати між ними +// завдання так, щоб воно пережило перезапуск будь-якого з них. +func (s *Server) handleCollectConfig(w http.ResponseWriter, r *http.Request, p *Principal) { + if !requirePerm(w, p, "ncm:write") { + return + } + deviceID := r.PathValue("id") + if !p.Scope().CanWrite(deviceID) { + writeError(w, http.StatusForbidden, "forbidden", "немає доступу до цього хоста") + return + } + + jobID, err := s.store.EnqueueConfigJob(r.Context(), p.TenantID, deviceID, "manual", p.UserID) + if errors.Is(err, store.ErrNotFound) { + writeError(w, http.StatusNotFound, "not_found", "хост не знайдено") + return + } + if err != nil { + s.writeStoreError(w, "постановка збору конфігу", err) + return + } + writeJSON(w, http.StatusAccepted, map[string]any{"job_id": jobID}) +} + +// handleListConfigJobs — історія збору для хоста. +func (s *Server) handleListConfigJobs(w http.ResponseWriter, r *http.Request, p *Principal) { + if !requirePerm(w, p, "ncm:read") { + return + } + deviceID := r.PathValue("id") + if !p.Scope().CanRead(deviceID) { + writeError(w, http.StatusForbidden, "forbidden", "немає доступу до цього хоста") + return + } + + jobs, err := s.store.ListConfigJobs(r.Context(), p.TenantID, deviceID, 20) + if err != nil { + s.writeStoreError(w, "історія збору", err) + return + } + if jobs == nil { + jobs = []store.ConfigJobRow{} + } + writeJSON(w, http.StatusOK, map[string]any{"jobs": jobs}) +} diff --git a/server/internal/httpapi/server.go b/server/internal/httpapi/server.go index 44446df..ef5a84f 100644 --- a/server/internal/httpapi/server.go +++ b/server/internal/httpapi/server.go @@ -91,6 +91,9 @@ func (s *Server) Handler() http.Handler { mux.Handle("PATCH /api/v1/devices/{id}", s.authenticated(s.handleUpdateDevice)) mux.Handle("DELETE /api/v1/devices/{id}", s.authenticated(s.handleDeleteDevice)) + mux.Handle("POST /api/v1/devices/{id}/collect-config", s.authenticated(s.handleCollectConfig)) + mux.Handle("GET /api/v1/devices/{id}/config-jobs", s.authenticated(s.handleListConfigJobs)) + mux.Handle("GET /api/v1/check-types", s.authenticated(s.handleListCheckTypes)) mux.Handle("GET /api/v1/devices/{id}/checks", s.authenticated(s.handleListDeviceChecks)) mux.Handle("PUT /api/v1/devices/{id}/checks", s.authenticated(s.handleSetDeviceChecks)) diff --git a/server/internal/store/ncm_jobs.go b/server/internal/store/ncm_jobs.go new file mode 100644 index 0000000..16668c2 --- /dev/null +++ b/server/internal/store/ncm_jobs.go @@ -0,0 +1,396 @@ +package store + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "time" + + "github.com/jackc/pgx/v5" + "github.com/netpulse/netpulse/server/internal/crypto" + npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1" + "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 := npv1.Transport_TRANSPORT_SSH + port := uint32(22) + if prof.Transport == "telnet" { + transport = npv1.Transport_TRANSPORT_TELNET + port = 23 + } + if cred != nil && cred.Port != 0 { + port = cred.Port + } + + 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 +}