Один коміт, а не десяток тематичних, свідомо: теми переплетені в
спільних файлах (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 серпня.
431 lines
15 KiB
Go
431 lines
15 KiB
Go
// 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
|
||
}
|