/
t3
/
cli
Обзор
Документация
Войти
/
t3
/
cli
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
internal/events/batch.go
91 строка
2 KB
Ivan
cline global refactoring
02 июл 2026, 07:43
02 июл 2026, 07:43
db38e5a
Код
Авторство
О чём код?
package events import ( "sync" "time" "gitverse.ru/t3/cli/pkg/pb" ) type BatchingGrpcStreamSender struct { stream ResponseSender mu sync.Mutex buffer []*pb.LogEvent maxBatch int flushCh chan struct{} closeCh chan struct{} } // NewBatchingGrpcStreamSender initializes the sender and starts the background flushing goroutine. func NewBatchingGrpcStreamSender(stream ResponseSender, maxBatch int, flushInterval time.Duration) *BatchingGrpcStreamSender { s := &BatchingGrpcStreamSender{ stream: stream, maxBatch: maxBatch, flushCh: make(chan struct{}, 1), closeCh: make(chan struct{}), } go s.run(flushInterval) return s } // SendLog adds a LogEvent to the buffer and triggers a flush if the batch size is reached. func (s *BatchingGrpcStreamSender) SendLog(event *pb.LogEvent) error { s.mu.Lock() s.buffer = append(s.buffer, event) shouldFlush := len(s.buffer) >= s.maxBatch s.mu.Unlock() if shouldFlush { select { case s.flushCh <- struct{}{}: default: } } return nil } func (s *BatchingGrpcStreamSender) run(interval time.Duration) { ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ticker.C: s.flush() case <-s.flushCh: s.flush() case <-s.closeCh: s.flush() return } } } func (s *BatchingGrpcStreamSender) flush() { s.mu.Lock() if len(s.buffer) == 0 { s.mu.Unlock() return } batch := s.buffer s.buffer = nil // reset buffer s.mu.Unlock() resp := &pb.AgentResponse{ Payload: &pb.AgentResponse_Logs{ Logs: &pb.LogEvents{ Logs: batch, }, }, } // Note: In production, you should handle the error from Send() appropriately. _ = s.stream.Send(resp) } // Close stops the background goroutine and flushes remaining logs. func (s *BatchingGrpcStreamSender) Close() { close(s.closeCh) }