Netpulse_SasS/server/internal/grpcapi/streams.go
byrsapty ed8fc831bf Дві сесії роботи: 0058–0068, розгортання однією командою, тести
Один коміт, а не десяток тематичних, свідомо: теми переплетені в
спільних файлах (store.go, docker-compose.yml, deploy/README.md), і
розділити їх можна було б лише індексуванням шматків. Коміти, які не
збираються, гірші за один великий — тим паче що це рівно той стан, який
перевірявся разом.

ЩО ПРАЦЮЄ НА СТЕНДІ Й ПЕРЕВІРЕНО ТАМ

  0058  подієві алерти: syslog, ncm, compliance спрацьовують у мить
        події; правило з нереалізованим джерелом більше не зберігається
        мовчки
  0059  snmp.walk і прототипи шаблонів — таблиці з динамічним індексом
        описуються шаблоном, а не Go
  0060  відкат конфігу: план як різниця, маскування паролів із підписом
        плану, обов'язковий контрольний збір, verifying при обриві
  0061  кнопки Telegram: довге опитування, авторизація не з callback_data
  0062  аудит і архів хостів; тест на AST, що падає на ключі без назви
  0063  RLS: три ролі, окремий пул для фонових тактів
  0064  строки зберігання даних і сторінка сховища
  0065  приймач SNMP-трапів; перевірено справжніми пакетами по дроту,
        переклад v1→v2 за RFC 3584 дає правильний OID
  0066  ескалації сповіщень
  0067  алерт про вичерпання диска
  0068  поля заливки конфігу переїхали в каталог профілів

Плюс: 137 тестів вебу з нуля (їх не було взагалі), одинадцять справжніх
вад, знайдених ними й виправлених, і виправлення двох інтеграційних
тестів grpcapi, які мовчки пропускались півтора року.

ЩО ЩЕ НЕ ЗАПУСКАЛОСЬ

  netpulse            установник: одна команда замість 18 змінних і
                      593 рядків інструкції
  RLS з першого запуску  нова інсталяція під політиками одразу;
                      RLS-EXISTING-INSTALL.md лишається тільки для
                      старих інсталяцій
  .forgejo + CI       раннер не зареєстрований

Ці три перевірені компіляцією й міркуванням, але не виконанням.

ГОЛОВНИЙ ВИСНОВОК ДВОХ СЕСІЙ

Зелена перевірка доводить рівно те, що вона перевіряє. Тест ізоляції RLS
був правильний і зелений — і пропустив зламаний вхід, бо перевіряв «чи
не видно чужого», коли зламалось «чи видно своє». Інтеграційні тести
grpcapi були зелені, бо не виконувались. Схема, довідник і протокол
описували те, чого в коді не існувало, і виглядало це як готове.

Тому в кожному завданні цих сесій стояла вимога назвати НЕПОКРИТЕ, а
чотири задачі закінчились не можливістю, а відмовою: правило з
нереалізованим джерелом не зберігається, профіль без команд заливки
каже про це замість мовчазної кнопки, міграція RLS валить сама себе на
таблиці без політики, тест словника аудиту падає на ключі без назви.

Подробиці — HISTORY.md, розділи за 26 і 27 серпня.
2026-08-27 17:32:49 +03:00

641 lines
27 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package grpcapi
import (
"bytes"
"compress/gzip"
"context"
"errors"
"fmt"
"io"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"github.com/netpulse/netpulse/server/internal/alerting"
"github.com/netpulse/netpulse/server/internal/store"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)
// ---------------------------------------------------------------------
// Телеметрія
// ---------------------------------------------------------------------
func (s *Service) StreamTelemetry(stream npv1.AgentService_StreamTelemetryServer) error {
ctx := stream.Context()
agent, err := agentFrom(ctx)
if err != nil {
return err
}
// Таблиця series_ref живе рівно стільки, скільки цей стрім.
// Агент знає про це й перереєстровує серії на початку кожної сесії.
table := store.NewSeriesTable()
var acked uint64
for {
batch, err := stream.Recv()
if errors.Is(err, io.EOF) {
return nil
}
if err != nil {
return err
}
st, writeErr := s.store.WriteBatch(ctx, agent, batch, table)
var unknown *store.ErrUnknownSeriesRef
if errors.As(writeErr, &unknown) {
// Ми загубили відповідність — чесно просимо перереєстрацію
// замість тихо викинути дані.
s.log.Warn("невідомий series_ref", "agent", agent.ID, "ref", unknown.Ref)
if err := stream.Send(&npv1.TelemetryAck{
AckedThroughBatchId: acked,
ResetSeriesTable: true,
MaxInFlight: uint32(agent.Limits.MaxInFlight),
Error: &npv1.Error{
Code: "unknown_series_ref",
Message: unknown.Error(),
Retryable: true,
},
}); err != nil {
return err
}
continue
}
if writeErr != nil {
s.log.Error("батч не записався", "agent", agent.ID,
"batch", batch.GetBatchId(), "err", writeErr)
// Батч не підтверджуємо: агент перешле його після
// реконекту, а ON CONFLICT DO NOTHING зробить це безпечним.
if err := stream.Send(&npv1.TelemetryAck{
AckedThroughBatchId: acked,
NackBatchIds: []uint64{batch.GetBatchId()},
MaxInFlight: uint32(agent.Limits.MaxInFlight),
Error: &npv1.Error{
Code: "write_failed", Message: writeErr.Error(), Retryable: true,
},
}); err != nil {
return err
}
continue
}
if batch.GetBatchId() > acked {
acked = batch.GetBatchId()
}
s.log.Debug("батч записано", "agent", agent.ID, "batch", batch.GetBatchId(),
"samples", st.Samples, "icmp", st.Icmp, "interfaces", st.Interfaces)
// Перелік рядків динамічних таблиць їде в payload результатів
// snmp.walk — тим самим батчем, що й метрики, бо збирає його та
// сама задача зонда.
//
// Окремо від WriteBatch навмисно: там усе лягає одним pgx.Batch
// без жодного читання, а тут — читання, порівняння зі станом і
// перебудова чеків. Помилка тут не має нікачити батч: метрики
// вже записані, і просити зонд переслати їх заради рядків
// означало б подвоїти телеметрію через таблицю дисків.
//
// Штовхати зонду новий план звідси не треба: звірка планів
// (SyncPlans) щоп'ять секунд помітить інший хеш і перезаллє
// його сама — тим самим шляхом, яким доїжджають чеки, створені
// у вебі.
if changed, err := s.store.ApplyWalkResults(ctx, agent, batch.GetCheckResults()); err != nil {
s.log.Warn("рядки прототипів не застосувались", "agent", agent.ID, "err", err)
} else if changed > 0 {
s.log.Info("склад рядків прототипів змінився",
"agent", agent.ID, "хостів", changed)
}
if err := stream.Send(&npv1.TelemetryAck{
AckedThroughBatchId: acked,
MaxInFlight: uint32(agent.Limits.MaxInFlight),
}); err != nil {
return err
}
}
}
// ---------------------------------------------------------------------
// Журнали
// ---------------------------------------------------------------------
func (s *Service) StreamLogs(stream npv1.AgentService_StreamLogsServer) error {
ctx := stream.Context()
agent, err := agentFrom(ctx)
if err != nil {
return err
}
var acked uint64
for {
batch, err := stream.Recv()
if errors.Is(err, io.EOF) {
return nil
}
if err != nil {
return err
}
if _, err := s.store.WriteLogs(ctx, agent, batch); err != nil {
s.log.Error("журнали не записались", "agent", agent.ID, "err", err)
if err := stream.Send(&npv1.LogAck{
AckedThroughBatchId: acked,
Error: &npv1.Error{Code: "write_failed", Message: err.Error(), Retryable: true},
}); err != nil {
return err
}
continue
}
if batch.GetDropped() > 0 {
s.log.Warn("зонд відкинув події через ліміт",
"agent", agent.ID, "dropped", batch.GetDropped())
}
s.triggerSyslogBackups(ctx, agent, batch.GetSyslog())
s.raiseSyslogAlerts(ctx, agent, batch.GetSyslog())
s.handleTraps(ctx, agent, batch.GetTraps())
if batch.GetBatchId() > acked {
acked = batch.GetBatchId()
}
if err := stream.Send(&npv1.LogAck{AckedThroughBatchId: acked}); err != nil {
return err
}
}
}
// triggerSyslogBackups ставить позачерговий збір конфігу за подією в
// журналі.
//
// Заради цього приймач журналу й потрібен NCM: нічний cron ловить зміну
// в середньому через дванадцять годин, а «%SYS-5-CONFIG_I» — за секунди.
// Інженер, який щось зламав о десятій ранку, встигає піти додому, і без
// цього тригера лишається вчорашній конфіг і сьогоднішня аварія без
// нічого між ними.
//
// Помилка тут не зриває приймання журналу: події вже записані, і
// втратити їх через проблему з чергою завдань було б обміном гіршим за
// пропущений бекап.
func (s *Service) triggerSyslogBackups(ctx context.Context, agent *store.Agent, entries []*npv1.SyslogEntry) {
if len(entries) == 0 {
return
}
devices, err := s.store.SyslogBackupTriggers(ctx, agent.TenantID, entries)
if err != nil {
s.log.Error("звірка журналу з політиками бекапу", "agent", agent.ID, "err", err)
return
}
for _, deviceID := range devices {
// Порожній користувач — завдання від системи, а не від людини.
// Повторний виклик для хоста з уже поставленим завданням нічого
// не додає: EnqueueConfigJob сам це відсіює, тож сплеск однакових
// рядків не перетворюється на сплеск сесій до пристрою.
if _, err := s.store.EnqueueConfigJob(ctx, agent.TenantID, deviceID, "syslog", ""); err != nil {
s.log.Error("бекап за подією журналу не поставлено",
"device", deviceID, "err", err)
continue
}
s.log.Info("подія журналу призначила збір конфігу", "device", deviceID)
}
}
// raiseSyslogAlerts звіряє щойно прийняті рядки з подієвими правилами.
//
// Робиться тут, одразу після запису, і саме з тієї ж причини, що й
// позачерговий бекап поруч: правило «у журналі зʼявилось %LINK-3-UPDOWN»
// має спрацювати за секунди, а не тоді, коли хтось наступного разу
// відкриє журнал. Опитувати ts.syslog розкладом було б і дорожче
// (гіпертаблиця на мільярд рядків), і брехливіше — вікно опитування
// завжди або пропускає події, або рахує їх двічі.
//
// Помилка тут не зриває приймання: рядки вже записані, і втратити
// журнал через проблему з алертами було б обміном гіршим за пропущений
// алерт. Тому весь розбір мовчить у сам приймач, а гучний він усередині.
func (s *Service) raiseSyslogAlerts(ctx context.Context, agent *store.Agent, entries []*npv1.SyslogEntry) {
if s.events == nil || len(entries) == 0 {
return
}
evs := make([]alerting.SyslogEvent, 0, len(entries))
for _, e := range entries {
evs = append(evs, alerting.SyslogEvent{
DeviceID: e.GetDeviceId(),
Message: e.GetMessage(),
Tag: e.GetTag(),
Severity: int(e.GetSeverity()),
})
}
s.events.OnSyslog(ctx, agent.TenantID, evs)
}
// handleTraps робить із щойно прийнятих трапів дві речі.
//
// Перша — алерти, тим самим шляхом і з тих самих міркувань, що й для
// журналу: правило «linkDown на магістральному порту» має спрацювати за
// секунди, а не тоді, коли хтось наступного разу відкриє журнал.
//
// Друга — облік відправників, яких зонд не зміг зіставити з хостом.
// Це не побічний ефект, а половина сенсу приймача. Трап від адреси,
// якої немає в інвентарі, — найчастіше перший слід нового заліза в
// мережі, і рівно він губиться у всіх системах, де подія без хоста
// просто відкидається. Алертом його не зробиш (алерт без хоста нікуди
// не маршрутизується), тому він потрапляє в окремий перелік, який видно
// на сторінці трапів.
//
// Обидві дії гучні всередині й мовчазні назовні: трапи вже записані, і
// втратити стрім через проблему з алертами було б обміном гіршим за
// пропущений алерт.
func (s *Service) handleTraps(ctx context.Context, agent *store.Agent, traps []*npv1.SnmpTrap) {
if len(traps) == 0 {
return
}
if err := s.store.NoteUnknownTrapSources(ctx, agent.TenantID, agent.ID, traps); err != nil {
s.log.Error("облік невідомих джерел трапів", "agent", agent.ID, "err", err)
}
if s.events == nil {
return
}
evs := make([]alerting.TrapEvent, 0, len(traps))
for _, t := range traps {
ev := alerting.TrapEvent{
DeviceID: t.GetDeviceId(),
SourceIP: t.GetSourceIp(),
TrapOID: t.GetTrapOid(),
}
for _, vb := range t.GetVarbinds() {
ev.Varbinds = append(ev.Varbinds, alerting.TrapVarbind{
OID: vb.GetOid(), Value: vb.GetValue(),
})
}
evs = append(evs, ev)
}
s.events.OnTrap(ctx, agent.TenantID, evs)
}
// raiseConfigAlert доводить долю збору конфігу до подієвих правил.
//
// Дві події, а не одна: «конфіг змінився» і «конфіг не зібрався» —
// різні новини для різних людей. Перша цікавить того, хто відповідає за
// зміни; друга — того, хто відповідає за те, щоб бекапи взагалі були.
// Звести їх в одну означало б, що ввімкнувши потрібну, отримуєш і зайву.
func (s *Service) raiseConfigAlert(ctx context.Context, tenantID, deviceID, configType, kind, detail string) {
if s.events == nil || deviceID == "" {
return
}
s.events.OnConfig(ctx, tenantID, alerting.ConfigEvent{
DeviceID: deviceID,
ConfigType: configType,
Kind: kind,
Detail: detail,
})
}
// ---------------------------------------------------------------------
// Автовиявлення
// ---------------------------------------------------------------------
func (s *Service) ReportDiscovery(ctx context.Context, rep *npv1.DiscoveryReport) (*npv1.DiscoveryAck, error) {
agent, err := agentFrom(ctx)
if err != nil {
return nil, err
}
st, err := s.store.ApplyDiscovery(ctx, agent, rep)
if err != nil {
s.log.Error("звіт автовиявлення не застосувався", "agent", agent.ID, "err", err)
return &npv1.DiscoveryAck{
Accepted: false,
Error: &npv1.Error{Code: "apply_failed", Message: err.Error(), Retryable: true},
}, nil
}
s.log.Info("автовиявлення застосовано",
"agent", agent.ID,
"сусідів", st.NeighborsSeen, "зіставлено", st.NeighborsResolved,
інків_створено", st.LinksCreated, інків_оновлено", st.LinksUpdated)
s.refreshInterfaceChecks(ctx, agent, rep)
// Системна інформація приїжджає тим самим звітом і дає найдешевший
// онбординг з можливих: пристрій сам сказав, що він таке, і шаблон
// причепився без жодного натискання.
info, err := s.store.ApplySystemInfo(ctx, agent.TenantID, rep.GetDevices())
switch {
case err != nil:
s.log.Error("системна інформація не застосувалась", "agent", agent.ID, "err", err)
case info.Failed > 0:
// Не Error: решта хостів у звіті оброблена, і зупиняти на цьому
// онбординг немає підстав. Але й ховати не можна — причина
// лежить у картці кожного з них, а тут видно масштаб.
s.log.Warn("частину хостів не розпізнано",
"agent", agent.ID, "хостів", info.Described,
"невдач", info.Failed, "err", info.LastError)
case info.Assigned > 0 || info.HardwareChanged > 0:
s.log.Info("розпізнавання застосовано",
"agent", agent.ID, "хостів", info.Described,
"шаблонів", info.Assigned, амінааліза", info.HardwareChanged)
}
return &npv1.DiscoveryAck{
Accepted: true,
NeighborsResolved: uint32(st.NeighborsResolved),
LinksCreated: uint32(st.LinksCreated),
}, nil
}
// refreshInterfaceChecks створює або оновлює snmp.if-чеки за щойно
// виявленими інтерфейсами й одразу штовхає зміну живому зонду.
//
// Без цього кроку автовиявлення наповнює inv.interfaces, але ніхто їх
// не опитує: на мапі є лінки й немає трафіку. Чекати наступного
// перепідключення агента, щоб він забрав новий план, — це години
// порожніх графіків після кожного нового комутатора.
func (s *Service) refreshInterfaceChecks(ctx context.Context, agent *store.Agent, rep *npv1.DiscoveryReport) {
seen := make(map[string]bool)
for _, r := range rep.GetInterfaces() {
if id := r.GetDeviceId(); id != "" {
seen[id] = true
}
}
if len(seen) == 0 {
return
}
var (
tasks []*npv1.Task
deviceIDs []string
)
for deviceID := range seen {
task, err := s.store.EnsureInterfaceChecks(ctx, agent, deviceID)
if err != nil {
// Обрізаний список портів — теж помилка, але чек при цьому
// створено, тому це попередження, а не привід зупинитись.
s.log.Warn("snmp.if-чек", "device", deviceID, "err", err)
}
if task != nil {
tasks = append(tasks, task)
deviceIDs = append(deviceIDs, deviceID)
}
}
if len(tasks) == 0 {
return
}
// Новий хеш обов'язковий: інакше після реконекту агент доповість
// старий, сервер вирішить, що план застарів, і перезаллє все.
hash, err := s.store.PlanHash(ctx, agent)
if err != nil {
s.log.Warn("не вдалося перерахувати хеш плану", "agent", agent.ID, "err", err)
return
}
devices, err := s.store.DeviceTargets(ctx, agent, deviceIDs)
if err != nil {
s.log.Warn("не вдалося зібрати описи пристроїв", "agent", agent.ID, "err", err)
return
}
delivered := s.PushToAgent(agent.ID, &npv1.ControlDown{
Payload: &npv1.ControlDown_TaskDelta{TaskDelta: &npv1.TaskDelta{
PlanHash: hash,
Upsert: tasks,
UpsertDevices: devices,
}},
})
if delivered {
s.setAgentPlanHash(agent.ID, hash)
}
s.log.Info("snmp.if-чеки оновлено",
"agent", agent.ID, "задач", len(tasks), адісланоаживо", delivered)
}
// ---------------------------------------------------------------------
// Конфігурації
// ---------------------------------------------------------------------
// UploadConfig збирає конфіг із чанків: header → chunk* → trailer.
func (s *Service) UploadConfig(stream npv1.AgentService_UploadConfigServer) error {
ctx := stream.Context()
agent, err := agentFrom(ctx)
if err != nil {
return err
}
var (
header *npv1.ConfigHeader
body []byte
wantSeq uint32
)
for {
msg, err := stream.Recv()
if errors.Is(err, io.EOF) {
return status.Error(codes.InvalidArgument, "стрім завершився без trailer")
}
if err != nil {
return err
}
switch p := msg.Part.(type) {
case *npv1.ConfigUpload_Header:
header = p.Header
case *npv1.ConfigUpload_Chunk:
if header == nil {
return status.Error(codes.InvalidArgument, "chunk раніше за header")
}
// Порядок важливий: конфіг склеюється байт-у-байт, і
// переставлені чанки дали б тихо зіпсований бекап.
if p.Chunk.GetSequence() != wantSeq {
return status.Errorf(codes.InvalidArgument,
"чанки не по порядку: чекали %d, отримали %d", wantSeq, p.Chunk.GetSequence())
}
wantSeq++
body = append(body, p.Chunk.GetData()...)
case *npv1.ConfigUpload_Trailer:
if header == nil {
return status.Error(codes.InvalidArgument, "trailer без header")
}
tr := p.Trailer
// Масове виконання команд повертається цим самим стрімом:
// шлях сервер→зонд→сервер уже є, і другий такий самий
// заради іншого призначення виводу був би копією з власними
// помилками. Розвилка стоїть саме тут, до розбору тіла, бо
// далі йде логіка бекапу — звірка з попередньою версією,
// коміт у Git, — якої для `show version` не існує.
if s.isCommandUpload(ctx, header) {
return stream.SendAndClose(
s.storeCommandResult(ctx, header, body, tr))
}
if !tr.GetSuccess() {
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())
s.raiseConfigAlert(ctx, agent.TenantID, header.GetDeviceId(),
header.GetConfigType(), "backup_failed", tr.GetError().GetMessage())
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())
s.raiseConfigAlert(ctx, agent.TenantID, header.GetDeviceId(),
header.GetConfigType(), "backup_failed", err.Error())
return stream.SendAndClose(&npv1.ConfigReceipt{
JobId: header.GetJobId(),
Accepted: false,
Error: &npv1.Error{
Code: "bad_encoding",
Message: err.Error(),
},
})
}
body = plain
// Набір конфіг-файлів сервера несе відбиток машини, з якої
// його знято. Звіряємо ДО збереження.
//
// Перевірка стоїть тут, а не в StoreConfig: той нічого не
// знає про зонди й машини й не має починати. А знати про це
// мусить рівно одне місце — те, куди приходять набори.
//
// Ціна помилки саме тут найвища: історія конфігів сервера
// живе в одній гілці Git, і файли іншої машини, дописані в
// неї, виглядають звичайною зміною конфігу. Помітити таке
// можна хіба через півроку — коли з архіву треба
// відновлюватись.
if mid := header.GetMachineId(); mid != "" {
if err := s.store.PinSelfMachine(ctx, agent.TenantID,
header.GetDeviceId(), mid); err != nil {
s.log.Warn("набір конфіг-файлів відхилено: не та машина",
"agent", agent.ID, "device", header.GetDeviceId(), "err", err)
_ = s.store.FinishConfigJob(ctx, header.GetJobId(), "failed",
err.Error(), tr.GetTranscript())
s.raiseConfigAlert(ctx, agent.TenantID, header.GetDeviceId(),
header.GetConfigType(), "backup_failed", err.Error())
return stream.SendAndClose(&npv1.ConfigReceipt{
JobId: header.GetJobId(),
Accepted: false,
Error: &npv1.Error{
Code: "machine_mismatch",
Message: err.Error(),
},
})
}
}
outcome, err := s.store.StoreConfig(ctx, agent, store.ConfigSubmission{
JobID: header.GetJobId(),
DeviceID: header.GetDeviceId(),
ConfigType: header.GetConfigType(),
Body: body,
ClaimedHash: tr.GetContentSha256(),
}, s.ring)
if err != nil {
return status.Errorf(codes.Internal, "збереження конфігу: %v", err)
}
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())
s.raiseConfigAlert(ctx, agent.TenantID, header.GetDeviceId(),
header.GetConfigType(), "backup_failed",
"тіло не відповідає заявленому sha256")
return stream.SendAndClose(&npv1.ConfigReceipt{
JobId: header.GetJobId(),
Accepted: false,
Error: &npv1.Error{
Code: "checksum_mismatch", Retryable: true,
Message: "тіло не відповідає заявленому sha256",
},
})
}
s.log.Info("конфіг прийнято",
"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())
// Подія рівно тоді, коли конфіг СПРАВДІ інший. Збіг хеша —
// не зміна, і алертувати на кожен нічний збір означало б
// щоранку віддавати черговому сорок повідомлень «усе як
// було».
if outcome.Accepted && !outcome.Unchanged {
s.raiseConfigAlert(ctx, agent.TenantID, header.GetDeviceId(),
header.GetConfigType(), "changed", outcome.CommitSHA)
}
return stream.SendAndClose(&npv1.ConfigReceipt{
JobId: header.GetJobId(),
Accepted: outcome.Accepted,
Unchanged: outcome.Unchanged,
ConfigId: outcome.ConfigID,
CommitSha: outcome.CommitSHA,
})
}
}
}
// 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