Netpulse_SasS/agent/internal/session/session.go
byrsapty 3816e2cfcf Етап 7: збір конфігів запрацював наскрізно
Профілі були даними без виконавця. Тепер ланцюг замкнено: черга в БД →
диспетчер у AgentService → зонд по SSH/Telnet → вивантаження стрімом →
звірка sha256 і дедуплікація.

Агент (internal/ncmx):
- усе побудовано навколо пошуку промпту: у консолі немає ані коду
  завершення, ані довжини відповіді
- промпт шукається лише в хвості 512 байтів, інакше "banner motd #"
  обривав би збір на середині конфігу
- дедлайн на паузу між байтами, а не на всю операцію: збір із шасі
  триває хвилини, а тиша означає завислий пристрій або хибний regex
- ключі SSH обладнання не звіряються свідомо: залізо міняє їх з кожною
  прошивкою, і known_hosts на сотні пристроїв означав би не збирати
  конфіги зовсім

Сервер: черга ncm.jobs із FOR UPDATE SKIP LOCKED (два екземпляри не
надішлють одне завдання двічі), диспетчер, прибирання завислих.

Знайдено живим прогоном:
- сервер ігнорував поле encoding: агент стискав gzip і рахував sha256
  від оригіналу, сервер рахував від стиснених байтів і відхиляв кожен
  бекап. Поле в контракті з Етапу 2, реалізації не було — інтеграційний
  тест користувався "none" і повз дірку проходив
- промпт обрізався не по рядку: шаблон [>#]\s*$ ловить лише символ, і
  в конфізі лишалось ім'я пристрою окремим рядком
- у конфіг потрапляли escape-послідовності bash і подвоєні \r\r\n від
  PTY: 295 байтів / 12 рядків замість 286 / 10. Після виправлення —
  285 / 10, різниця лише у фінальному переводі рядка

Живий прогін проти стенду: success, конфіг 285 байтів, повторний збір —
unchanged без другої версії.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-24 12:02:50 +03:00

764 lines
23 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 session тримає з'єднання з сервером і все, що з нього випливає:
// рукостискання, застосування плану задач, відправку телеметрії,
// підтвердження, реконект.
//
// Усе з'єднання ініціює агент — у мережі клієнта немає ані відкритих
// портів, ані прокидання NAT. Тому «команда з сервера» технічно є
// повідомленням у зустрічному напрямку вже відкритого стріму Control.
package session
import (
"context"
"errors"
"fmt"
"log/slog"
"math/rand"
"runtime"
"sync"
"sync/atomic"
"time"
"github.com/netpulse/netpulse/agent/internal/module"
"github.com/netpulse/netpulse/agent/internal/scheduler"
"github.com/netpulse/netpulse/agent/internal/telemetry"
npv1 "github.com/netpulse/netpulse/gen/go/netpulse/v1"
"google.golang.org/grpc"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
)
// Conn — те, що вміє *grpc.ClientConn. Винесено в інтерфейс, щоб тести
// підставляли bufconn без справжньої мережі.
type Conn interface {
grpc.ClientConnInterface
Close() error
}
// DialFunc відкриває з'єднання до сервера.
type DialFunc func(ctx context.Context) (Conn, error)
type Config struct {
AgentID string
Hostname string
Build *npv1.AgentBuild
Registry *module.Registry
Buffer *telemetry.Buffer
Scheduler *scheduler.Scheduler
Dial DialFunc
Logger *slog.Logger
MinBackoff time.Duration
MaxBackoff time.Duration
// Модулі, увімкнені до першої відповіді сервера.
DefaultModules []string
}
type Session struct {
cfg Config
log *slog.Logger
// Статуси задач — діагностика. Канал з буфером і скиданням при
// переповненні: краще втратити рядок журналу, ніж загальмувати
// опитування через повільний контрольний канал.
statusCh chan *npv1.TaskStatusUpdate
// Звіти автовиявлення. Окремо від телеметрії: вони рідкі, великі
// й їдуть унарним ReportDiscovery, а не стрімом.
discoCh chan *npv1.DiscoveryReport
credMu sync.RWMutex
creds map[string][]*npv1.Credential
credExpiry time.Time
devMu sync.RWMutex
devices map[string]*npv1.DeviceTarget
// client живий лише в межах сесії. Збір конфігу триває хвилинами й
// не має тримати контрольний цикл, тому вивантаження йде окремою
// горутиною — а їй потрібен доступ до клієнта поточної сесії.
client atomic.Pointer[npv1.AgentServiceClient]
// jobs рахує незавершені збори: при обриві сесії їх треба дочекатись,
// інакше вивантаження піде в уже закритий канал.
jobs sync.WaitGroup
planHash atomic.Pointer[[]byte]
nextBatch atomic.Uint64
lastAcked atomic.Uint64
skewNanos atomic.Int64
connected atomic.Bool
seq atomic.Uint64
startedAt time.Time
}
func New(cfg Config) *Session {
if cfg.Logger == nil {
cfg.Logger = slog.Default()
}
if cfg.MinBackoff <= 0 {
cfg.MinBackoff = time.Second
}
if cfg.MaxBackoff <= 0 {
cfg.MaxBackoff = 2 * time.Minute
}
return &Session{
cfg: cfg,
log: cfg.Logger,
statusCh: make(chan *npv1.TaskStatusUpdate, 512),
discoCh: make(chan *npv1.DiscoveryReport, 16),
creds: make(map[string][]*npv1.Credential),
devices: make(map[string]*npv1.DeviceTarget),
startedAt: time.Now(),
}
}
// SetScheduler замикає взаємну залежність: планувальник потребує
// колбеків сесії (статуси, креденшели), а сесія — планувальника, щоб
// застосовувати план. Конструктором це не виражається, тому окремий сетер.
// Викликати до Run.
func (s *Session) SetScheduler(sc *scheduler.Scheduler) { s.cfg.Scheduler = sc }
// Credentials — провайдер для планувальника.
func (s *Session) Credentials(deviceID string) []*npv1.Credential {
s.credMu.RLock()
defer s.credMu.RUnlock()
// Прострочені креденшели не віддаємо: краще явна помилка
// "немає креденшелів", ніж заблокований обліковий запис на
// половині комутаторів клієнта.
if !s.credExpiry.IsZero() && time.Now().After(s.credExpiry) {
return nil
}
return s.creds[deviceID]
}
// ReportDiscovery — колбек для планувальника.
//
// Черга навмисно коротка: звіт автовиявлення актуальний рівно доти,
// доки описує поточний стан мережі. Накопичувати десяток застарілих
// знімків, поки немає зв'язку, немає сенсу — наступний обхід віддасть
// свіжіші дані.
func (s *Session) ReportDiscovery(res module.Result) {
if len(res.Neighbors) == 0 && len(res.InterfaceRecords) == 0 {
return
}
rep := &npv1.DiscoveryReport{
AgentId: s.cfg.AgentID,
Neighbors: res.Neighbors,
Interfaces: res.InterfaceRecords,
Final: true,
StartedAt: timestamppb.Now(),
FinishedAt: timestamppb.Now(),
}
select {
case s.discoCh <- rep:
default:
// Витісняємо найстаріший, а не відкидаємо новий.
select {
case <-s.discoCh:
default:
}
select {
case s.discoCh <- rep:
default:
}
}
}
// ReportStatus — колбек для планувальника.
func (s *Session) ReportStatus(u *npv1.TaskStatusUpdate) {
select {
case s.statusCh <- u:
default:
}
}
// Connected — чи є жива сесія (для /healthz і самометрик).
func (s *Session) Connected() bool { return s.connected.Load() }
// ClockSkew — поправка годинника відносно сервера.
func (s *Session) ClockSkew() time.Duration {
return time.Duration(s.skewNanos.Load())
}
// ---------------------------------------------------------------------
// Цикл підключення
// ---------------------------------------------------------------------
// Run тримає з'єднання, доки не скасують контекст.
func (s *Session) Run(ctx context.Context) error {
backoff := s.cfg.MinBackoff
for {
if ctx.Err() != nil {
return ctx.Err()
}
err := s.runOnce(ctx)
s.connected.Store(false)
if ctx.Err() != nil {
return ctx.Err()
}
if err != nil {
s.log.Warn("сесію розірвано", "err", err, "retry_in", backoff)
}
// Джитер — щоб сотня зондів не ломанулась перепідключатись
// одночасно після рестарту сервера.
jitter := time.Duration(rand.Int63n(int64(backoff/2 + 1)))
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(backoff + jitter):
}
backoff *= 2
if backoff > s.cfg.MaxBackoff {
backoff = s.cfg.MaxBackoff
}
}
}
func (s *Session) runOnce(ctx context.Context) error {
conn, err := s.cfg.Dial(ctx)
if err != nil {
return fmt.Errorf("dial: %w", err)
}
defer conn.Close()
client := npv1.NewAgentServiceClient(conn)
s.client.Store(&client)
defer s.client.Store(nil)
sctx, cancel := context.WithCancel(ctx)
defer cancel()
ctrl, err := client.Control(sctx)
if err != nil {
return fmt.Errorf("відкрити Control: %w", err)
}
if err := ctrl.Send(s.hello()); err != nil {
return fmt.Errorf("Hello: %w", err)
}
first, err := ctrl.Recv()
if err != nil {
return fmt.Errorf("очікування Welcome: %w", err)
}
welcome := first.GetWelcome()
if welcome == nil {
return fmt.Errorf("замість Welcome прийшло %T", first.Payload)
}
s.applyWelcome(welcome)
// Таблиця series_ref на сервері живе рівно стільки, скільки сесія.
// Перереєстровуємо всі відомі серії, замість обнуляти нумерацію:
// інакше довелося б викинути буфер, накопичений за час обриву, —
// саме ті дані, заради яких він і накопичувався.
s.cfg.Buffer.Interner().MarkAllPending()
s.connected.Store(true)
s.log.Info("сесію встановлено",
"session_id", welcome.SessionId,
"clock_skew", s.ClockSkew())
// Єдиний писар у контрольний стрім: gRPC не допускає паралельних
// Send, а слати треба і heartbeat, і статуси задач.
out := make(chan *npv1.ControlUp, 256)
var wg sync.WaitGroup
// Місткість дорівнює числу горутин: інакше та, що впала останньою,
// заблокується на записі й ніколи не дочекається wg.Wait().
errCh := make(chan error, 8)
spawn := func(name string, fn func() error) {
wg.Add(1)
go func() {
defer wg.Done()
if err := fn(); err != nil && !errors.Is(err, context.Canceled) {
errCh <- fmt.Errorf("%s: %w", name, err)
}
cancel()
}()
}
spawn("writer", func() error { return s.writerLoop(sctx, ctrl, out) })
spawn("heartbeat", func() error { return s.heartbeatLoop(sctx, out, welcome) })
spawn("status", func() error { return s.statusLoop(sctx, out) })
spawn("telemetry", func() error { return s.telemetryLoop(sctx, client, welcome) })
spawn("discovery", func() error { return s.discoveryLoop(sctx, client) })
// Читання команд — у цій же горутині.
readErr := s.controlLoop(sctx, ctrl, out)
cancel()
wg.Wait()
// Збір конфігу переживає скасування контексту не довше за свій
// дедлайн, але закривати з'єднання під ним не можна: вивантаження
// впаде на середині, і сервер отримає обірваний стрім замість
// чесної помилки.
s.jobs.Wait()
close(errCh)
if readErr != nil && !errors.Is(readErr, context.Canceled) {
return readErr
}
for e := range errCh {
if e != nil {
return e
}
}
return nil
}
func (s *Session) hello() *npv1.ControlUp {
var hash []byte
if p := s.planHash.Load(); p != nil {
hash = *p
}
return &npv1.ControlUp{
Seq: s.seq.Add(1),
Payload: &npv1.ControlUp_Hello{Hello: &npv1.Hello{
AgentId: s.cfg.AgentID,
Build: s.cfg.Build,
Hostname: s.cfg.Hostname,
// Хеш плану дозволяє серверу не перезаливати 50 000 задач
// після кожного обриву зв'язку.
TaskPlanHash: hash,
// Продовжуємо з місця розриву, а не з нуля.
LastAckedBatchId: s.lastAcked.Load(),
StartedAt: timestamppb.New(s.startedAt),
}},
}
}
func (s *Session) applyWelcome(w *npv1.Welcome) {
if w.ServerTime != nil {
s.skewNanos.Store(int64(time.Until(w.ServerTime.AsTime())))
}
s.cfg.Registry.EnsureDefaults(s.cfg.DefaultModules...)
}
// ---------------------------------------------------------------------
// Цикли
// ---------------------------------------------------------------------
func (s *Session) writerLoop(ctx context.Context, ctrl npv1.AgentService_ControlClient, out <-chan *npv1.ControlUp) error {
for {
select {
case <-ctx.Done():
return nil
case msg := <-out:
if err := ctrl.Send(msg); err != nil {
return err
}
}
}
}
func (s *Session) enqueue(ctx context.Context, out chan<- *npv1.ControlUp, msg *npv1.ControlUp) {
msg.Seq = s.seq.Add(1)
select {
case out <- msg:
case <-ctx.Done():
default:
// Черга забита — контрольний канал не встигає. Втратити
// heartbeat не страшно: сервер помітить це за таймаутом.
}
}
func (s *Session) heartbeatLoop(ctx context.Context, out chan<- *npv1.ControlUp, w *npv1.Welcome) error {
interval := 30 * time.Second
if w.HeartbeatInterval != nil && w.HeartbeatInterval.AsDuration() > 0 {
interval = w.HeartbeatInterval.AsDuration()
}
t := time.NewTicker(interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return nil
case <-t.C:
hb := &npv1.Heartbeat{Ts: timestamppb.Now(), Health: s.health()}
if s.cfg.Scheduler != nil {
hb.TasksRunning = uint32(s.cfg.Scheduler.Running())
hb.TasksQueued = uint32(s.cfg.Scheduler.Len())
}
s.enqueue(ctx, out, &npv1.ControlUp{
Payload: &npv1.ControlUp_Heartbeat{Heartbeat: hb},
})
}
}
}
func (s *Session) health() *npv1.AgentHealth {
var ms runtime.MemStats
runtime.ReadMemStats(&ms)
stats := s.cfg.Buffer.Stats()
return &npv1.AgentHealth{
RssBytes: ms.Sys,
Goroutines: uint32(runtime.NumGoroutine()),
QueueDepth: uint32(stats.Items),
DroppedSamples: stats.Dropped,
Uptime: durationpb.New(time.Since(s.startedAt)),
ClockSkew: durationpb.New(s.ClockSkew()),
}
}
func (s *Session) statusLoop(ctx context.Context, out chan<- *npv1.ControlUp) error {
for {
select {
case <-ctx.Done():
return nil
case u := <-s.statusCh:
s.enqueue(ctx, out, &npv1.ControlUp{
Payload: &npv1.ControlUp_TaskStatus{TaskStatus: u},
})
}
}
}
// controlLoop читає команди сервера.
func (s *Session) controlLoop(ctx context.Context, ctrl npv1.AgentService_ControlClient, out chan<- *npv1.ControlUp) error {
for {
msg, err := ctrl.Recv()
if err != nil {
if ctx.Err() != nil {
return nil
}
return err
}
switch p := msg.Payload.(type) {
case *npv1.ControlDown_TaskPlan:
s.applyPlan(p.TaskPlan)
case *npv1.ControlDown_TaskDelta:
s.applyDelta(p.TaskDelta)
case *npv1.ControlDown_ModuleControl:
enabled := make(map[string]bool, len(p.ModuleControl.Modules))
for _, m := range p.ModuleControl.Modules {
enabled[m.Key] = m.Enabled
}
s.cfg.Registry.SetActive(enabled, p.ModuleControl.Exclusive)
case *npv1.ControlDown_Credentials:
s.applyCredentials(p.Credentials)
case *npv1.ControlDown_Ping:
s.enqueue(ctx, out, &npv1.ControlUp{
Payload: &npv1.ControlUp_Pong{Pong: &npv1.Pong{
PingId: p.Ping.PingId,
AgentTime: timestamppb.Now(),
}},
})
if p.Ping.ServerTime != nil {
s.skewNanos.Store(int64(time.Until(p.Ping.ServerTime.AsTime())))
}
case *npv1.ControlDown_DiscoveryRequest:
if s.cfg.Scheduler == nil {
continue
}
n := s.cfg.Scheduler.TriggerNow(p.DiscoveryRequest.GetDeviceIds(), "topo.")
s.log.Info("сервер попросив запустити автовиявлення",
"run_id", p.DiscoveryRequest.GetRunId(), адач_зрушено", n)
case *npv1.ControlDown_ConfigJob:
job := p.ConfigJob
s.jobs.Add(1)
go func() {
defer s.jobs.Done()
s.runConfigJob(ctx, job)
}()
case *npv1.ControlDown_Directive:
if stop := s.applyDirective(p.Directive); stop {
return nil
}
case *npv1.ControlDown_Welcome:
// Повторний Welcome у межах сесії — сервер щось наплутав.
s.log.Warn("повторний Welcome у живій сесії — ігноруємо")
}
}
}
func (s *Session) applyDirective(d *npv1.Directive) (stop bool) {
switch d.Action {
case npv1.Directive_ACTION_PAUSE:
s.log.Info("опитування призупинено сервером", "reason", d.Reason)
s.cfg.Scheduler.SetPaused(true)
case npv1.Directive_ACTION_RESUME:
s.log.Info("опитування відновлено")
s.cfg.Scheduler.SetPaused(false)
case npv1.Directive_ACTION_RESET_SERIES_TABLE:
s.cfg.Buffer.ResetSeries()
case npv1.Directive_ACTION_RECONNECT, npv1.Directive_ACTION_DRAIN:
s.log.Info("сервер попросив перепідключитись", "reason", d.Reason)
return true
case npv1.Directive_ACTION_UPDATE:
// Застосування оновлення — окремий крок; без валідного підпису
// Ed25519 воно не має відбутись узагалі, тому тут лише журнал.
s.log.Info("доступне оновлення", "version", d.Update.GetVersion())
}
return false
}
func (s *Session) applyCredentials(b *npv1.CredentialBundle) {
s.credMu.Lock()
defer s.credMu.Unlock()
s.creds = make(map[string][]*npv1.Credential, len(b.ByDevice))
for devID, list := range b.ByDevice {
s.creds[devID] = list.Credentials
}
if b.ExpiresAt != nil {
s.credExpiry = b.ExpiresAt.AsTime()
} else {
s.credExpiry = time.Time{}
}
}
// ---------------------------------------------------------------------
// План задач
// ---------------------------------------------------------------------
func (s *Session) applyPlan(plan *npv1.TaskPlan) {
// Часткові плани не застосовуємо: половина розкладу гірша за старий
// цілий. Сервер має надіслати final = true.
if !plan.Final {
s.stashDevices(plan.Devices)
return
}
s.stashDevices(plan.Devices)
tasks, meta := s.buildTasks(plan.Tasks)
s.cfg.Scheduler.Apply(tasks, meta)
h := plan.PlanHash
s.planHash.Store(&h)
s.log.Info("план задач застосовано", "tasks", len(tasks))
}
func (s *Session) applyDelta(d *npv1.TaskDelta) {
s.stashDevices(d.UpsertDevices)
if len(d.RemoveCheckIds) > 0 {
s.cfg.Scheduler.Remove(d.RemoveCheckIds)
}
if len(d.Upsert) > 0 {
tasks, meta := s.buildTasks(d.Upsert)
s.cfg.Scheduler.Upsert(tasks, meta)
}
if len(d.RemoveDeviceIds) > 0 {
s.devMu.Lock()
for _, id := range d.RemoveDeviceIds {
delete(s.devices, id)
}
s.devMu.Unlock()
}
h := d.PlanHash
s.planHash.Store(&h)
}
func (s *Session) stashDevices(devs []*npv1.DeviceTarget) {
if len(devs) == 0 {
return
}
s.devMu.Lock()
defer s.devMu.Unlock()
for _, d := range devs {
s.devices[d.DeviceId] = d
}
}
func (s *Session) buildTasks(in []*npv1.Task) ([]module.Task, []scheduler.TaskMeta) {
s.devMu.RLock()
defer s.devMu.RUnlock()
tasks := make([]module.Task, 0, len(in))
meta := make([]scheduler.TaskMeta, 0, len(in))
for _, t := range in {
if !t.Enabled {
continue
}
dev := s.devices[t.DeviceId]
if dev == nil {
// Задача без пристрою — сервер надіслав неповний план.
s.log.Warn("задача посилається на невідомий пристрій",
"check_id", t.CheckId, "device_id", t.DeviceId)
continue
}
tasks = append(tasks, module.Task{
CheckID: t.CheckId,
DeviceID: t.DeviceId,
InterfaceID: t.InterfaceId,
CheckType: t.CheckType,
Params: t.ParamsJson,
Target: module.Target{
DeviceID: dev.DeviceId,
Name: dev.Name,
Address: dev.Address,
},
Timeout: t.Timeout.AsDuration(),
})
meta = append(meta, scheduler.TaskMeta{
Interval: t.Interval.AsDuration(),
Offset: t.ScheduleOffset.AsDuration(),
})
}
return tasks, meta
}
// ---------------------------------------------------------------------
// Автовиявлення
// ---------------------------------------------------------------------
func (s *Session) discoveryLoop(ctx context.Context, client npv1.AgentServiceClient) error {
for {
select {
case <-ctx.Done():
return nil
case rep := <-s.discoCh:
ack, err := client.ReportDiscovery(ctx, rep)
if err != nil {
if ctx.Err() != nil {
return nil
}
// Звіт втрачено свідомо: наступний обхід за розкладом
// віддасть свіжіші дані, а ретрай застарілого знімка
// топології нічого не додає.
s.log.Warn("звіт автовиявлення не доставлено", "err", err)
continue
}
s.log.Info("автовиявлення прийнято сервером",
"сусідів", len(rep.Neighbors),
"зіставлено", ack.GetNeighborsResolved(),
"лінків", ack.GetLinksCreated())
}
}
}
// ---------------------------------------------------------------------
// Телеметрія
// ---------------------------------------------------------------------
func (s *Session) telemetryLoop(ctx context.Context, client npv1.AgentServiceClient, w *npv1.Welcome) error {
stream, err := client.StreamTelemetry(ctx)
if err != nil {
return err
}
maxBatch := int(w.TelemetryMaxBatchSize)
if maxBatch <= 0 {
maxBatch = 500
}
flush := 5 * time.Second
if w.TelemetryMaxBatchInterval != nil && w.TelemetryMaxBatchInterval.AsDuration() > 0 {
flush = w.TelemetryMaxBatchInterval.AsDuration()
}
maxInFlight := int(w.TelemetryMaxInFlight)
if maxInFlight <= 0 {
maxInFlight = 4
}
var (
mu sync.Mutex
inFlight = make(map[uint64]*npv1.TelemetryBatch)
)
// Читач підтверджень. Send і Recv на різних горутинах — це те
// єдине, що gRPC дозволяє робити паралельно на одному стрімі.
ackErr := make(chan error, 1)
go func() {
for {
ack, err := stream.Recv()
if err != nil {
ackErr <- err
return
}
if ack.ResetSeriesTable {
s.log.Warn("сервер попросив перереєструвати серії")
s.cfg.Buffer.ResetSeries()
}
mu.Lock()
for id := range inFlight {
if id <= ack.AckedThroughBatchId {
delete(inFlight, id)
}
}
mu.Unlock()
if ack.AckedThroughBatchId > s.lastAcked.Load() {
s.lastAcked.Store(ack.AckedThroughBatchId)
}
if ack.RetryAfter != nil && ack.RetryAfter.AsDuration() > 0 {
select {
case <-ctx.Done():
return
case <-time.After(ack.RetryAfter.AsDuration()):
}
}
}
}()
// Усе, що не встигло підтвердитись, повертаємо в буфер:
// схема БД робить повторний запис безпечним (at-least-once).
defer func() {
mu.Lock()
for _, b := range inFlight {
b.IsRetransmit = true
s.cfg.Buffer.Requeue(b)
}
mu.Unlock()
}()
ticker := time.NewTicker(flush)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return nil
case err := <-ackErr:
return err
case <-ticker.C:
case <-s.cfg.Buffer.Notify():
}
mu.Lock()
busy := len(inFlight) >= maxInFlight
mu.Unlock()
if busy {
// Зворотний тиск: сервер не встигає — не додаємо йому роботи.
continue
}
batchID := s.nextBatch.Add(1)
batch := s.cfg.Buffer.Drain(batchID, s.cfg.AgentID, maxBatch)
if batch == nil {
s.nextBatch.Add(^uint64(0)) // відкотити невикористаний номер
continue
}
mu.Lock()
inFlight[batchID] = batch
mu.Unlock()
if err := stream.Send(batch); err != nil {
mu.Lock()
delete(inFlight, batchID)
mu.Unlock()
s.cfg.Buffer.Requeue(batch)
return err
}
}
}