Netpulse_SasS/agent/internal/scheduler/scheduler.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

504 lines
14 KiB
Go
Raw 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 scheduler виконує план задач, отриманий від сервера.
//
// Дві вимоги визначають будову:
// 1. 5000 чеків з інтервалом 60 с не мають стартувати одночасно —
// інакше зонд раз на хвилину викидає сплеск на всю мережу.
// Розведення дає schedule_offset, який рахує СЕРВЕР (детерміновано
// від check_id), щоб воно зберігалось між перезапусками агента.
// 2. Повільний пристрій не має гальмувати решту: кожна задача має
// власний дедлайн, а паралельність обмежена семафором.
package scheduler
import (
"container/heap"
"context"
"strings"
"sync"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
)
// Sink приймає результати. Реалізується telemetry.Buffer.
type Sink interface {
Add(deviceID, pluginKey string, res module.Result)
AddCheckResult(cr *npv1.CheckResult)
}
// StatusFunc доповідає серверу про життєвий цикл задачі.
type StatusFunc func(*npv1.TaskStatusUpdate)
// DiscoveryFunc отримує результати автовиявлення.
//
// Вони не йдуть телеметричним стрімом: сусіди й інвентар портів рідкі,
// великі й не прив'язані до моменту часу так, як метрики. Для них
// окремий RPC ReportDiscovery.
type DiscoveryFunc func(res module.Result)
// CredentialFunc віддає облікові дані пристрою на момент виконання.
//
// Навмисно не поле в Task: креденшели мають TTL і оновлюються окремим
// повідомленням від сервера. Якби вони копіювались у план, після
// ротації паролів агент довбав би пристрої простроченими даними,
// доки не приїде новий план.
type CredentialFunc func(deviceID string) []*npv1.Credential
type entry struct {
task module.Task
interval time.Duration
offset time.Duration
nextRun time.Time
running bool
index int
}
// Scheduler — планувальник із min-heap за часом наступного запуску.
type Scheduler struct {
registry *module.Registry
sink Sink
onStatus StatusFunc
creds CredentialFunc
onDisco DiscoveryFunc
mu sync.Mutex
entries map[string]*entry
queue taskHeap
wake chan struct{}
paused bool
sem chan struct{}
nowFn func() time.Time
}
type Config struct {
Registry *module.Registry
Sink Sink
OnStatus StatusFunc
Credentials CredentialFunc
OnDiscovery DiscoveryFunc
MaxConcurrency int
// Підміна годинника в тестах.
Now func() time.Time
}
func New(cfg Config) *Scheduler {
if cfg.MaxConcurrency <= 0 {
cfg.MaxConcurrency = 64
}
if cfg.Now == nil {
cfg.Now = time.Now
}
if cfg.OnStatus == nil {
cfg.OnStatus = func(*npv1.TaskStatusUpdate) {}
}
if cfg.Credentials == nil {
cfg.Credentials = func(string) []*npv1.Credential { return nil }
}
if cfg.OnDiscovery == nil {
cfg.OnDiscovery = func(module.Result) {}
}
return &Scheduler{
registry: cfg.Registry,
sink: cfg.Sink,
onStatus: cfg.OnStatus,
creds: cfg.Credentials,
onDisco: cfg.OnDiscovery,
entries: make(map[string]*entry),
wake: make(chan struct{}, 1),
sem: make(chan struct{}, cfg.MaxConcurrency),
nowFn: cfg.Now,
}
}
// SetPaused зупиняє видачу задач, не втрачаючи розкладу.
// Сервер вмикає це на grace-періоді після несплати: канал живий,
// опитування стоїть.
func (s *Scheduler) SetPaused(v bool) {
s.mu.Lock()
s.paused = v
s.mu.Unlock()
s.wakeup()
}
// NextRun рахує момент наступного запуску так, щоб він був вирівняний
// по сітці інтервалу й зсунутий на offset.
//
// Вирівнювання по сітці (а не "зараз + інтервал") робить розклад
// відтворюваним: після рестарту агента задача повертається у свій
// слот, а не з'їжджає на випадковий момент.
func NextRun(now time.Time, interval, offset time.Duration) time.Time {
if interval <= 0 {
interval = time.Minute
}
if offset < 0 || offset >= interval {
offset = time.Duration(0)
}
next := now.Truncate(interval).Add(offset)
for !next.After(now) {
next = next.Add(interval)
}
return next
}
// ---------------------------------------------------------------------
// Застосування плану
// ---------------------------------------------------------------------
// Apply замінює весь план (ControlDown.TaskPlan).
func (s *Scheduler) Apply(tasks []module.Task, meta []TaskMeta) {
s.mu.Lock()
defer s.mu.Unlock()
// Зберігаємо стан running: задача, що зараз виконується, не має
// зникнути з-під ніг у горутини, яка її виконує.
old := s.entries
s.entries = make(map[string]*entry, len(tasks))
s.queue = nil
now := s.nowFn()
for i, t := range tasks {
m := meta[i]
e := &entry{task: t, interval: m.Interval, offset: m.Offset}
if prev, ok := old[t.CheckID]; ok {
e.running = prev.running
}
e.nextRun = NextRun(now, e.interval, e.offset)
s.entries[t.CheckID] = e
heap.Push(&s.queue, e)
}
s.wakeup()
}
// Upsert застосовує інкрементальну зміну (ControlDown.TaskDelta).
func (s *Scheduler) Upsert(tasks []module.Task, meta []TaskMeta) {
s.mu.Lock()
defer s.mu.Unlock()
now := s.nowFn()
for i, t := range tasks {
m := meta[i]
if e, ok := s.entries[t.CheckID]; ok {
e.task = t
// Розклад перераховуємо лише якщо він реально змінився,
// інакше редагування опису пристрою збивало б фазу.
if e.interval != m.Interval || e.offset != m.Offset {
e.interval, e.offset = m.Interval, m.Offset
e.nextRun = NextRun(now, e.interval, e.offset)
heap.Fix(&s.queue, e.index)
}
continue
}
e := &entry{task: t, interval: m.Interval, offset: m.Offset}
e.nextRun = NextRun(now, e.interval, e.offset)
s.entries[t.CheckID] = e
heap.Push(&s.queue, e)
}
s.wakeup()
}
// Remove знімає задачі з розкладу.
func (s *Scheduler) Remove(checkIDs []string) {
s.mu.Lock()
defer s.mu.Unlock()
for _, id := range checkIDs {
e, ok := s.entries[id]
if !ok {
continue
}
if e.index >= 0 && e.index < len(s.queue) {
heap.Remove(&s.queue, e.index)
}
delete(s.entries, id)
}
s.wakeup()
}
// TaskMeta — розкладова частина задачі.
type TaskMeta struct {
Interval time.Duration
Offset time.Duration
}
// TriggerNow ставить у чергу негайне виконання задач, що підпадають
// під фільтр. Потрібне для кнопки «Запустити виявлення» в UI: сервер
// надсилає DiscoveryRequest, а не чекає наступного циклу розкладу.
//
// Повертає, скільки задач зрушено — сервер має знати, що його запит
// не потрапив у порожнечу.
func (s *Scheduler) TriggerNow(deviceIDs []string, checkTypePrefix string) int {
s.mu.Lock()
defer s.mu.Unlock()
want := make(map[string]bool, len(deviceIDs))
for _, id := range deviceIDs {
want[id] = true
}
now := s.nowFn()
n := 0
for _, e := range s.entries {
if checkTypePrefix != "" && !strings.HasPrefix(e.task.CheckType, checkTypePrefix) {
continue
}
if len(want) > 0 && !want[e.task.DeviceID] {
continue
}
if e.running {
continue
}
e.nextRun = now
if e.index >= 0 && e.index < len(s.queue) {
heap.Fix(&s.queue, e.index)
}
n++
}
if n > 0 {
s.wakeup()
}
return n
}
// Len — скільки задач у розкладі.
func (s *Scheduler) Len() int {
s.mu.Lock()
defer s.mu.Unlock()
return len(s.entries)
}
// Running — скільки задач виконується просто зараз.
func (s *Scheduler) Running() int {
s.mu.Lock()
defer s.mu.Unlock()
n := 0
for _, e := range s.entries {
if e.running {
n++
}
}
return n
}
func (s *Scheduler) wakeup() {
select {
case s.wake <- struct{}{}:
default:
}
}
// ---------------------------------------------------------------------
// Головний цикл
// ---------------------------------------------------------------------
// Run крутиться до скасування контексту.
func (s *Scheduler) Run(ctx context.Context) {
timer := time.NewTimer(time.Hour)
defer timer.Stop()
var wg sync.WaitGroup
defer wg.Wait()
for {
now := s.nowFn()
due, wait := s.popDue(now)
for _, e := range due {
wg.Add(1)
go func(e *entry) {
defer wg.Done()
s.execute(ctx, e)
}(e)
}
if !timer.Stop() {
select {
case <-timer.C:
default:
}
}
timer.Reset(wait)
select {
case <-ctx.Done():
return
case <-timer.C:
case <-s.wake:
}
}
}
// popDue дістає задачі, яким час виконуватись, і повертає час до
// наступної. Задача, попередній запуск якої ще триває, пропускається
// з явним STATE_SKIPPED — тиха втрата циклу гірша за видиму.
func (s *Scheduler) popDue(now time.Time) ([]*entry, time.Duration) {
s.mu.Lock()
defer s.mu.Unlock()
if s.paused {
return nil, time.Second
}
var due []*entry
for s.queue.Len() > 0 {
top := s.queue[0]
if top.nextRun.After(now) {
break
}
heap.Pop(&s.queue)
top.nextRun = NextRun(now, top.interval, top.offset)
heap.Push(&s.queue, top)
if top.running {
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: top.task.CheckID,
State: npv1.TaskStatusUpdate_STATE_SKIPPED,
Ts: timestamppb.New(now),
Error: &npv1.Error{
Code: "still_running",
Message: "попередній запуск не завершився до наступного інтервалу",
Retryable: true,
},
})
continue
}
top.running = true
due = append(due, top)
}
wait := time.Hour
if s.queue.Len() > 0 {
if d := s.queue[0].nextRun.Sub(now); d < wait {
wait = d
}
}
if wait < time.Millisecond {
wait = time.Millisecond
}
return due, wait
}
func (s *Scheduler) execute(ctx context.Context, e *entry) {
defer func() {
s.mu.Lock()
e.running = false
s.mu.Unlock()
}()
task := e.task
// Креденшели беремо на момент виконання, а не з плану: у них TTL.
task.Credentials = s.creds(task.DeviceID)
mod, err := s.registry.Resolve(task.CheckType)
if err != nil {
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: task.CheckID,
State: npv1.TaskStatusUpdate_STATE_REJECTED,
Ts: timestamppb.Now(),
Error: &npv1.Error{Code: "no_module", Message: err.Error()},
})
return
}
// Семафор обмежує паралельність. Чекаємо на слот, але не довше,
// ніж дозволяє власний таймаут задачі.
select {
case s.sem <- struct{}{}:
defer func() { <-s.sem }()
case <-ctx.Done():
return
case <-time.After(task.Timeout):
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: task.CheckID,
State: npv1.TaskStatusUpdate_STATE_SKIPPED,
Ts: timestamppb.Now(),
Error: &npv1.Error{
Code: "concurrency_limit",
Message: "не дочекались вільного слота виконання",
Retryable: true,
},
})
return
}
timeout := task.Timeout
if timeout <= 0 {
timeout = 10 * time.Second
}
runCtx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
started := time.Now()
res, runErr := mod.Run(runCtx, task)
elapsed := time.Since(started)
cr := &npv1.CheckResult{
CheckId: task.CheckID,
DeviceId: task.DeviceID,
CheckType: task.CheckType,
Ts: timestamppb.New(started),
Duration: durationpb.New(elapsed),
Success: runErr == nil,
PayloadJson: res.Payload,
}
if runErr != nil {
cr.Error = &npv1.Error{
Code: classify(runErr, runCtx),
Message: runErr.Error(),
Retryable: true,
}
s.sink.AddCheckResult(cr)
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: task.CheckID,
State: npv1.TaskStatusUpdate_STATE_FAILED,
Ts: timestamppb.Now(),
Error: cr.Error,
})
return
}
s.sink.Add(task.DeviceID, module.ModuleKey(task.CheckType), res)
s.sink.AddCheckResult(cr)
// Devices теж рахуються: у режимі самого лише розпізнавання звіт
// не містить ні сусідів, ні портів — тільки системну групу, заради
// якої чек і заведено. Без цієї умови вона нікуди не їхала.
if len(res.Neighbors) > 0 || len(res.InterfaceRecords) > 0 || len(res.Devices) > 0 {
s.onDisco(res)
}
s.onStatus(&npv1.TaskStatusUpdate{
CheckId: task.CheckID,
State: npv1.TaskStatusUpdate_STATE_SUCCEEDED,
Ts: timestamppb.Now(),
})
}
func classify(err error, ctx context.Context) string {
if ctx.Err() == context.DeadlineExceeded {
return "timeout"
}
return "check_failed"
}
// ---------------------------------------------------------------------
// heap
// ---------------------------------------------------------------------
type taskHeap []*entry
func (h taskHeap) Len() int { return len(h) }
func (h taskHeap) Less(i, j int) bool { return h[i].nextRun.Before(h[j].nextRun) }
func (h taskHeap) Swap(i, j int) { h[i], h[j] = h[j], h[i]; h[i].index = i; h[j].index = j }
func (h *taskHeap) Push(x any) { e := x.(*entry); e.index = len(*h); *h = append(*h, e) }
func (h *taskHeap) Pop() any {
old := *h
n := len(old)
e := old[n-1]
old[n-1] = nil
e.index = -1
*h = old[:n-1]
return e
}