Найбільше вузьке місце до запуску: агент заводився INSERT-ом у базу, а токен вписувався в командний рядок руками. Поставити зонд у клієнта було неможливо. core.agent_enrollments тримає sha256 одноразового токена; сам токен повертається рівно один раз. Видача під FOR UPDATE в одній транзакції: два агенти з однієї скопійованої команди інакше створили б два зонди з одного запрошення. Відповідь на «немає», «згоріло» і «використано» однакова — розрізняти їх означає підказувати тому, хто підбирає токени. Токен зонда їде окремим полем agent_token, а не в certificate: сертифікат відповідає на інше питання й живе за іншим циклом. Агент зберігає посвідчення в /etc/netpulse/agent.json з правами 0600, через тимчасовий файл і перейменування — обрив живлення посеред запису інакше лишив би половину токена. Знайдено живим прогоном: реєстрація не проходила автентифікацію, бо інтерсептор стоїть на всьому сервері, а не на окремому сервісі — мій же коментар стверджував протилежне. І запуск із самим посвідченням падав: validate() вимагав -agent-id, не знаючи про файл. Сторінка зондів: команда встановлення з токеном, відкликання запрошень, керування модулями й лімітами, видалення. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
396 lines
13 KiB
Go
396 lines
13 KiB
Go
// Package grpcapi — серверна реалізація AgentService.
|
||
//
|
||
// Сервер ніколи не ініціює з'єднання до зонда: у мережі клієнта немає
|
||
// ані відкритих портів, ані прокидання NAT. Тому все, що виглядає як
|
||
// «команда з сервера», надсилається у зустрічному напрямку стріму
|
||
// Control, який відкрив сам агент.
|
||
package grpcapi
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"io"
|
||
"log/slog"
|
||
"strings"
|
||
"sync"
|
||
"time"
|
||
|
||
"github.com/netpulse/netpulse/server/internal/crypto"
|
||
"github.com/netpulse/netpulse/server/internal/store"
|
||
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
|
||
"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
|
||
|
||
// Живі сесії за 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),
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Автентифікація
|
||
// ---------------------------------------------------------------------
|
||
|
||
// 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)
|
||
}
|
||
}()
|
||
|
||
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_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
|
||
}
|