Netpulse_SasS/server/internal/grpcapi/service.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

431 lines
15 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 — серверна реалізація AgentService.
//
// Сервер ніколи не ініціює з'єднання до зонда: у мережі клієнта немає
// ані відкритих портів, ані прокидання NAT. Тому все, що виглядає як
// «команда з сервера», надсилається у зустрічному напрямку стріму
// Control, який відкрив сам агент.
package grpcapi
import (
"context"
"errors"
"io"
"log/slog"
"strings"
"sync"
"time"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"github.com/netpulse/netpulse/server/internal/alerting"
"github.com/netpulse/netpulse/server/internal/crypto"
"github.com/netpulse/netpulse/server/internal/store"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/metadata"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
)
type ctxKey string
const agentCtxKey ctxKey = "netpulse.agent"
// Service реалізує npv1.AgentServiceServer.
type Service struct {
npv1.UnimplementedAgentServiceServer
store *store.Store
ring *crypto.Keyring
log *slog.Logger
// Приймач подій для правил джерел `syslog` і `ncm`.
//
// Живе саме тут, бо саме сюди приходять рядки журналу й зібрані
// конфіги: правило на подію не має де спрацювати, крім тієї миті,
// коли подія надійшла. nil означає інсталяцію з вимкненими
// алертами — тоді приймач просто мовчить.
events *alerting.EventSink
// Живі сесії за agent_id. Потрібні, щоб штовхнути зонду
// TaskDelta або ConfigJob, коли щось змінилось в UI.
mu sync.RWMutex
sessions map[string]*agentSession
}
type agentSession struct {
agent *store.Agent
out chan *npv1.ControlDown
series *store.SeriesTable
opened time.Time
// Хеш плану, який зараз має зонд.
//
// Потрібен, бо зміни в UI роблять інший процес (REST API), а живу
// сесію тримає цей. Без цієї позначки чек, доданий у вебі,
// доїжджав би до зонда лише після обриву зв'язку — тобто, за
// нормальної роботи, ніколи.
planMu sync.Mutex
planHash []byte
}
func (a *agentSession) setPlanHash(h []byte) {
a.planMu.Lock()
a.planHash = h
a.planMu.Unlock()
}
func (a *agentSession) currentPlanHash() []byte {
a.planMu.Lock()
defer a.planMu.Unlock()
return a.planHash
}
func New(st *store.Store, ring *crypto.Keyring, log *slog.Logger) *Service {
if log == nil {
log = slog.Default()
}
return &Service{
store: st,
ring: ring,
log: log,
sessions: make(map[string]*agentSession),
}
}
// WithEventAlerts вмикає подієві алерти на журналі й конфігах.
//
// Окремим методом, а не аргументом New: приймач подій потрібен не
// кожній збірці (тести, читальні інстанси), і вимагати його від них
// означало б тягнути пакет алертів туди, де алертів немає.
func (s *Service) WithEventAlerts(sink *alerting.EventSink) *Service {
s.events = sink
return s
}
// ---------------------------------------------------------------------
// Автентифікація
// ---------------------------------------------------------------------
// StreamInterceptor і UnaryInterceptor автентифікують зонда за токеном
// у метаданих. У проді поверх цього ще mTLS: токен визначає, ЯКИЙ це
// зонд, сертифікат — що він узагалі має право говорити з сервером.
func (s *Service) StreamInterceptor(srv any, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
agent, err := s.authenticate(ss.Context())
if err != nil {
return err
}
return handler(srv, &wrappedStream{ServerStream: ss,
ctx: context.WithValue(ss.Context(), agentCtxKey, agent)})
}
// enrollMethodPrefix — єдиний виклик, який іде без токена зонда.
//
// Інакше й бути не може: саме цим викликом токен і видається. Виняток
// заданий повним префіксом сервісу, а не «містить Enroll»: підрядок у
// назві методу — це спосіб випадково відкрити ще щось, коли сервіс
// назвуть Enrollment-чимось іще.
const enrollMethodPrefix = "/netpulse.v1.EnrollmentService/"
func (s *Service) UnaryInterceptor(ctx context.Context, req any, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (any, error) {
if strings.HasPrefix(info.FullMethod, enrollMethodPrefix) {
return handler(ctx, req)
}
agent, err := s.authenticate(ctx)
if err != nil {
return nil, err
}
return handler(context.WithValue(ctx, agentCtxKey, agent), req)
}
type wrappedStream struct {
grpc.ServerStream
ctx context.Context
}
func (w *wrappedStream) Context() context.Context { return w.ctx }
func (s *Service) authenticate(ctx context.Context) (*store.Agent, error) {
md, ok := metadata.FromIncomingContext(ctx)
if !ok {
return nil, status.Error(codes.Unauthenticated, "немає метаданих")
}
var token string
for _, v := range md.Get("authorization") {
if after, found := strings.CutPrefix(v, "Bearer "); found {
token = after
break
}
}
if token == "" {
return nil, status.Error(codes.Unauthenticated, "немає токена зонда")
}
agent, err := s.store.AuthenticateAgent(ctx, token)
if errors.Is(err, store.ErrAgentNotFound) {
return nil, status.Error(codes.Unauthenticated, "зонда не знайдено")
}
if err != nil {
return nil, status.Error(codes.Internal, "помилка автентифікації")
}
return agent, nil
}
func agentFrom(ctx context.Context) (*store.Agent, error) {
a, ok := ctx.Value(agentCtxKey).(*store.Agent)
if !ok || a == nil {
return nil, status.Error(codes.Unauthenticated, "сесія без зонда")
}
return a, nil
}
// ---------------------------------------------------------------------
// Control
// ---------------------------------------------------------------------
func (s *Service) Control(stream npv1.AgentService_ControlServer) error {
ctx := stream.Context()
agent, err := agentFrom(ctx)
if err != nil {
return err
}
first, err := stream.Recv()
if err != nil {
return err
}
hello := first.GetHello()
if hello == nil {
return status.Error(codes.InvalidArgument, "перше повідомлення має бути Hello")
}
if hello.GetAgentId() != "" && hello.GetAgentId() != agent.ID {
// Токен належить одному зонду, а він назвався іншим.
return status.Error(codes.PermissionDenied, "agent_id не збігається з токеном")
}
if err := s.store.MarkAgentOnline(ctx, agent, hello); err != nil {
s.log.Warn("не вдалося позначити зонд онлайн", "agent", agent.ID, "err", err)
}
sess := &agentSession{
agent: agent,
out: make(chan *npv1.ControlDown, 64),
series: store.NewSeriesTable(),
opened: time.Now(),
}
s.register(sess)
defer func() {
s.unregister(agent.ID)
if err := s.store.MarkAgentOffline(context.WithoutCancel(ctx), agent); err != nil {
s.log.Warn("не вдалося позначити зонд офлайн", "agent", agent.ID, "err", err)
}
}()
// До побудови плану, а не після: інакше щойно заведений чек
// розпізнавання поїхав би до зонда лише наступною звіркою.
if _, err := s.store.EnsureIdentifyChecks(ctx, agent); err != nil {
s.log.Warn("чек розпізнавання", "agent", agent.ID, "err", err)
}
plan, err := s.store.BuildPlan(ctx, agent)
if err != nil {
return status.Errorf(codes.Internal, "побудова плану: %v", err)
}
// Якщо хеш збігся — план у агента вже правильний, і 50 000 задач
// після кожного обриву зв'язку переливати не треба.
planUnchanged := len(hello.GetTaskPlanHash()) > 0 &&
equalBytes(hello.GetTaskPlanHash(), plan.GetPlanHash())
welcome := &npv1.Welcome{
SessionId: agent.ID + "@" + time.Now().UTC().Format(time.RFC3339Nano),
ServerTime: timestamppb.Now(),
HeartbeatInterval: durationpb.New(agent.Limits.HeartbeatEvery),
TelemetryMaxBatchSize: uint32(agent.Limits.BatchSize),
TelemetryMaxBatchInterval: durationpb.New(agent.Limits.BatchInterval),
TelemetryMaxInFlight: uint32(agent.Limits.MaxInFlight),
MaxConcurrentChecks: uint32(agent.Limits.MaxConcurrency),
IcmpRatePps: uint32(agent.Limits.IcmpRatePPS),
TaskPlanFollows: !planUnchanged,
}
if err := stream.Send(&npv1.ControlDown{
Seq: 1, Payload: &npv1.ControlDown_Welcome{Welcome: welcome},
}); err != nil {
return err
}
// Єдиний писар: gRPC не допускає паралельних Send на одному стрімі,
// а слати треба і план, і креденшели, і пуші з UI.
writerDone := make(chan struct{})
go func() {
defer close(writerDone)
var seq uint64 = 1
for {
select {
case <-ctx.Done():
return
case msg, ok := <-sess.out:
if !ok {
return
}
seq++
msg.Seq = seq
if err := stream.Send(msg); err != nil {
return
}
}
}
}()
s.push(sess, &npv1.ControlDown{
Payload: &npv1.ControlDown_ModuleControl{
ModuleControl: store.ModulesForPlan(plan, agent.Modules)},
})
bundle, err := s.store.BuildCredentialBundle(ctx, agent, s.ring)
if err != nil {
s.log.Warn("не вдалося зібрати креденшели", "agent", agent.ID, "err", err)
} else if len(bundle.ByDevice) > 0 {
s.push(sess, &npv1.ControlDown{
Payload: &npv1.ControlDown_Credentials{Credentials: bundle},
})
}
if !planUnchanged {
s.push(sess, &npv1.ControlDown{
Payload: &npv1.ControlDown_TaskPlan{TaskPlan: plan},
})
}
sess.setPlanHash(plan.GetPlanHash())
s.log.Info("зонд підключився",
"agent", agent.ID, "tenant", agent.TenantID,
"tasks", len(plan.GetTasks()), "plan_unchanged", planUnchanged)
err = s.readControl(ctx, stream, agent, sess)
close(sess.out)
<-writerDone
return err
}
func (s *Service) readControl(ctx context.Context, stream npv1.AgentService_ControlServer, agent *store.Agent, sess *agentSession) error {
for {
msg, err := stream.Recv()
if errors.Is(err, io.EOF) {
return nil
}
if err != nil {
return err
}
switch p := msg.Payload.(type) {
case *npv1.ControlUp_Heartbeat:
if err := s.store.RecordHeartbeat(ctx, agent, p.Heartbeat); err != nil {
s.log.Warn("heartbeat не записався", "agent", agent.ID, "err", err)
}
// Ping у відповідь дає обом сторонам свіжий clock skew.
s.push(sess, &npv1.ControlDown{
Payload: &npv1.ControlDown_Ping{Ping: &npv1.Ping{
PingId: uint64(time.Now().UnixNano()), ServerTime: timestamppb.Now(),
}},
})
case *npv1.ControlUp_TaskStatus:
if err := s.store.RecordTaskStatus(ctx, agent, p.TaskStatus); err != nil {
s.log.Warn("статус задачі не записався", "agent", agent.ID, "err", err)
}
case *npv1.ControlUp_CredentialRequest:
bundle, err := s.store.BuildCredentialBundle(ctx, agent, s.ring)
if err != nil {
s.log.Warn("повторна видача креденшелів", "agent", agent.ID, "err", err)
continue
}
s.push(sess, &npv1.ControlDown{
Payload: &npv1.ControlDown_Credentials{Credentials: bundle},
})
case *npv1.ControlUp_ConfigApplyResult:
// Результат заливки конфігу. Обробляється в окремій
// горутині: FinishApply ставить контрольний збір і чекає
// на кілька запитів до бази, а контрольний цикл цієї сесії
// тим часом має відповідати на ping — інакше зонд, який
// щойно зробив найнебезпечнішу роботу, буде визнаний
// мертвим саме через неї.
result := p.ConfigApplyResult
go s.storeApplyResult(context.WithoutCancel(ctx), result)
case *npv1.ControlUp_Event:
s.log.Info("подія зонда",
"agent", agent.ID, "kind", p.Event.GetKind().String(),
"message", p.Event.GetMessage())
case *npv1.ControlUp_Pong:
// Нічого не робимо: сам факт відповіді підтверджує, що
// канал живий, а це вже зафіксував транспорт.
}
}
}
// push кладе повідомлення в чергу писаря. Черга переповнилась —
// зонд не встигає читати; краще втратити пуш, ніж заблокувати
// обробку всіх інших зондів на цьому воркері.
func (s *Service) push(sess *agentSession, msg *npv1.ControlDown) {
select {
case sess.out <- msg:
default:
s.log.Warn("черга контрольного каналу переповнена", "agent", sess.agent.ID)
}
}
func (s *Service) register(sess *agentSession) {
s.mu.Lock()
defer s.mu.Unlock()
if old, ok := s.sessions[sess.agent.ID]; ok {
// Той самий зонд відкрив другу сесію: стара мертва, просто
// ще не помітила. Закривати її чергу тут не можна — це
// зробить її власний defer.
s.log.Warn("друга сесія того самого зонда", "agent", sess.agent.ID,
"попередня_відкрита", old.opened)
}
s.sessions[sess.agent.ID] = sess
}
func (s *Service) unregister(agentID string) {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.sessions, agentID)
}
// SessionCount — скільки зондів онлайн просто зараз.
func (s *Service) SessionCount() int {
s.mu.RLock()
defer s.mu.RUnlock()
return len(s.sessions)
}
// PushToAgent надсилає команду живій сесії зонда. Використовується
// з боку API, коли в UI щось змінили.
func (s *Service) PushToAgent(agentID string, msg *npv1.ControlDown) bool {
s.mu.RLock()
sess, ok := s.sessions[agentID]
s.mu.RUnlock()
if !ok {
return false
}
s.push(sess, msg)
return true
}
func equalBytes(a, b []byte) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i] != b[i] {
return false
}
}
return true
}