package main import ( "bufio" "bytes" "context" "encoding/json" "fmt" "io" "log" "net/http" "net/url" "strings" "sync/atomic" "time" ) // Store sends messages to VictoriaLogs in batches and queries it with LogsQL. type Store struct { base string client *http.Client in chan *Entry batchSize int flushEvery time.Duration received atomic.Int64 ingested atomic.Int64 dropped atomic.Int64 lastErr atomic.Value // string } func NewStore(base string, batchSize, queueSize int, flushEvery time.Duration) *Store { s := &Store{ base: strings.TrimRight(base, "/"), client: &http.Client{Timeout: 60 * time.Second}, in: make(chan *Entry, queueSize), batchSize: batchSize, flushEvery: flushEvery, } s.lastErr.Store("") return s } // Enqueue never blocks: when the queue is full, the message is counted as dropped. func (s *Store) Enqueue(e *Entry) { s.received.Add(1) select { case s.in <- e: default: s.dropped.Add(1) } } // Run drains the queue into VictoriaLogs until the context is cancelled, // then flushes whatever is left. func (s *Store) Run(ctx context.Context) { ticker := time.NewTicker(s.flushEvery) defer ticker.Stop() batch := make([]*Entry, 0, s.batchSize) flush := func() { if len(batch) == 0 { return } if err := s.insert(batch); err != nil { s.dropped.Add(int64(len(batch))) s.lastErr.Store(err.Error()) log.Printf("victorialogs: %d messages dropped: %v", len(batch), err) } else { s.ingested.Add(int64(len(batch))) s.lastErr.Store("") } batch = batch[:0] } for { select { case e := <-s.in: batch = append(batch, e) if len(batch) >= s.batchSize { flush() } case <-ticker.C: flush() case <-ctx.Done(): for { select { case e := <-s.in: batch = append(batch, e) if len(batch) >= s.batchSize { flush() } default: flush() return } } } } } func (s *Store) insert(batch []*Entry) error { var buf bytes.Buffer enc := json.NewEncoder(&buf) enc.SetEscapeHTML(false) for _, e := range batch { if err := enc.Encode(e.Record()); err != nil { return err } } u := s.base + "/insert/jsonline?_stream_fields=host,app&_msg_field=_msg&_time_field=_time" var err error for attempt := 0; attempt < 5; attempt++ { if attempt > 0 { time.Sleep(time.Duration(1<