From 350426399212ce2d5c2f4b88c5db50dfcad4e205 Mon Sep 17 00:00:00 2001 From: Cedric Date: Sat, 3 Oct 2026 16:26:18 +0200 Subject: [PATCH] Disk buffer for batches VictoriaLogs cannot take, lossless Docker positions, non-blocking reverse DNS - Batches that fail go to /data/spool (SPOOL_MAX_MB, 1 GiB by default) and are sent again oldest first; retries no longer block the store loop and follow the shutdown context. - Docker and host logs wait for room in a full queue instead of being dropped; the Docker position only moves once a line is stored or spooled. - Reverse DNS no longer holds up the syslog listeners, with an LRU cache and a cap on concurrent lookups. Co-Authored-By: Claude Opus 5.5 (1M context) --- README.fr.md | 17 ++-- README.md | 14 ++- docker.go | 17 ++-- histogram_test.go | 4 +- hostlogs.go | 2 + main.go | 34 +++++-- rdns.go | 103 +++++++++++++++++--- rdns_test.go | 77 +++++++++++++++ spool.go | 136 +++++++++++++++++++++++++++ store.go | 232 +++++++++++++++++++++++++++++++++++++--------- store_test.go | 216 ++++++++++++++++++++++++++++++++++++++++++ syslog.go | 6 ++ web/app.js | 3 + web/style.css | 1 + 14 files changed, 778 insertions(+), 84 deletions(-) create mode 100644 rdns_test.go create mode 100644 spool.go create mode 100644 store_test.go diff --git a/README.fr.md b/README.fr.md index fbdfba7..3664538 100644 --- a/README.fr.md +++ b/README.fr.md @@ -18,8 +18,10 @@ devices ──514 udp/tcp──▶ logstream (Go) ──HTTP batches──▶ Vi ![Schéma logique de Logstream](docs/architecture.png) - **Ingestion** : les messages syslog (UDP/TCP) et les logs des conteneurs Docker passent tous - par `sink()` (résolution DNS inverse des hôtes donnés par leur IP), puis par la file du - `Store`, qui les envoie par lots à VictoriaLogs. + par `sink()` (résolution DNS inverse des hôtes donnés par leur IP, sans bloquer la + réception), puis par la file du `Store`, qui les envoie par lots à VictoriaLogs. Quand + VictoriaLogs est injoignable, les lots sont gardés sur disque (`/data/spool`, jusqu'à + `SPOOL_MAX_MB`) et renvoyés, les plus anciens d'abord, dès son retour. - **Direct** : `sink()` publie aussi chaque message dans le `Hub`, qui le diffuse aux navigateurs en SSE. - **Recherche** : l'API HTTP traduit les filtres de l'interface en requêtes LogsQL envoyées à @@ -153,8 +155,10 @@ colorés et exportés comme les messages syslog : quand les conteneurs sont nombreux. Les nouveaux conteneurs sont suivis automatiquement, sauf si cette option est désactivée. Les choix sont enregistrés par service compose (ou nom de conteneur) dans `/data/docker.json` : ils survivent aux recréations. -- Logstream mémorise la position lue dans chaque conteneur (`/data/docker-state.json`) : après - un redémarrage, il reprend sans perdre ni dupliquer de lignes. Un conteneur vu pour la +- Logstream mémorise la position de la dernière ligne stockée pour chaque conteneur + (`/data/docker-state.json`) : après un redémarrage, il reprend sans perdre de lignes. La + position n'avance qu'une fois la ligne dans VictoriaLogs ou dans le tampon disque, et une file + pleine ralentit la lecture au lieu de perdre des lignes. Un conteneur vu pour la première fois est lu à partir de `DOCKER_BACKFILL` en arrière (1 heure par défaut). - Logstream lui-même et le proxy ci-dessous ne sont jamais collectés ; ajoutez l'étiquette `logstream.exclude=true` à tout autre conteneur pour l'exclure définitivement. @@ -361,13 +365,14 @@ résolutions. | `HOST_LOGS_ROOT` | `/host` | emplacement de montage des répertoires de l'hôte | | `TZ` | `Europe/Paris` | fuseau horaire des horodatages RFC 3164 (qui n'en portent pas) | | `BATCH_SIZE`, `FLUSH_MS`, `QUEUE_SIZE` | `1000`, `1000`, `100000` | réglage de l'ingestion | +| `SPOOL_MAX_MB` | `1024` | taille du tampon disque des lots refusés par VictoriaLogs (`0` = pas de tampon : 15 s de tentatives, puis perte) | ## Débogage - `docker compose logs -f logstream` : erreurs de réception et erreurs d'envoi vers VictoriaLogs. -- La barre du bas affiche les compteurs reçus / stockés / perdus et la dernière erreur de - stockage. +- La barre du bas affiche les compteurs reçus / stockés / perdus, les messages en attente dans + le tampon disque et la dernière erreur de stockage. - : l'interface de VictoriaLogs, pour essayer des requêtes LogsQL. - API : diff --git a/README.md b/README.md index e29ae38..93436a5 100644 --- a/README.md +++ b/README.md @@ -17,7 +17,9 @@ devices ──514 udp/tcp──▶ logstream (Go) ──HTTP batches──▶ Vi ![Schéma logique de Logstream](docs/architecture.png) - **Ingestion**: syslog (UDP/TCP) and Docker container logs both go through `sink()` (reverse DNS - on IP hosts), then the `Store` queue, which sends them in batches to VictoriaLogs. + on IP hosts, without holding up the listeners), then the `Store` queue, which sends them in + batches to VictoriaLogs. When VictoriaLogs is unreachable, batches are kept on disk + (`/data/spool`, up to `SPOOL_MAX_MB`) and sent again, oldest first, once it is back. - **Live view**: `sink()` also publishes each message to the `Hub`, which streams it to the browsers over SSE. - **Search**: the HTTP API turns the UI filters into LogsQL queries sent to VictoriaLogs. @@ -137,8 +139,10 @@ colored and exported like syslog messages: to the labels shown) help with many containers. New containers are followed automatically unless that option is turned off. Choices are saved per compose service (or container name) in `/data/docker.json`, so they survive re-creations. -- Logstream remembers the position read in each container (`/data/docker-state.json`): after a - restart it resumes without losing or duplicating lines. A container seen for the first time +- Logstream remembers the position of the last line stored for each container + (`/data/docker-state.json`): after a restart it resumes without losing lines. The position + only moves once a line is in VictoriaLogs or in the disk buffer, and a full queue slows the + reading down instead of dropping lines. A container seen for the first time is read from `DOCKER_BACKFILL` ago (1 hour by default). - Logstream itself and the proxy below are never collected; add the label `logstream.exclude=true` to any other container to exclude it for good. @@ -330,11 +334,13 @@ are only known by your router or a local DNS (Pi-hole, AdGuard, Unbound…), set | `HOST_LOGS_ROOT` | `/host` | where the host directories are mounted | | `TZ` | `Europe/Paris` | time zone for RFC 3164 timestamps (which carry none) | | `BATCH_SIZE`, `FLUSH_MS`, `QUEUE_SIZE` | `1000`, `1000`, `100000` | ingestion tuning | +| `SPOOL_MAX_MB` | `1024` | disk buffer size for batches VictoriaLogs could not take (`0` = no buffer: retried for 15 s, then dropped) | ## Debugging - `docker compose logs -f logstream`: receive errors and errors sending to VictoriaLogs. -- The bottom bar shows received / stored / dropped counters and the last storage error. +- The bottom bar shows received / stored / dropped counters, the messages waiting in the disk + buffer, and the last storage error. - : VictoriaLogs' own UI to try LogsQL queries. - API: ```bash diff --git a/docker.go b/docker.go index ac39112..4aa49aa 100644 --- a/docker.go +++ b/docker.go @@ -61,8 +61,8 @@ type DockerManager struct { cfg dockerConfig hostName string containers []DockerContainer - followers map[string]*follower // container ID -> running follower - checkpoint map[string]time.Time // container ID -> timestamp of the last line read + followers map[string]*follower // container ID -> running follower + checkpoint map[string]time.Time // container ID -> timestamp of the last line stored dirty bool connected bool lastErr string @@ -389,16 +389,19 @@ func (m *DockerManager) follow(ctx context.Context, c DockerContainer) { } } -func sleepCtx(ctx context.Context, d time.Duration) { +// sleepCtx waits for d, or returns false if the context ends first. +func sleepCtx(ctx context.Context, d time.Duration) bool { t := time.NewTimer(d) defer t.Stop() select { case <-ctx.Done(): + return false case <-t.C: + return true } } -// since returns where to resume reading: just after the last line read, or +// since returns where to resume reading: just after the last line stored, or // DOCKER_BACKFILL ago for a container seen for the first time. func (m *DockerManager) since(id string) time.Time { m.mu.Lock() @@ -530,8 +533,6 @@ func (m *DockerManager) emit(c DockerContainer, stream string, line []byte) { } else if t, err := time.Parse(time.RFC3339Nano, s); err == nil { ts, s = t, "" } - defer m.setCheckpoint(c.ID, ts) - s = ansiRe.ReplaceAllString(s, "") if !utf8.ValidString(s) { s = strings.ToValidUTF8(s, string(utf8.RuneError)) @@ -575,6 +576,10 @@ func (m *DockerManager) emit(c DockerContainer, stream string, line []byte) { "compose_service": c.Service, "stream": stream, }, + // Docker keeps the logs: wait for room in the queue rather than drop the + // line, and resume after it only once it is stored. + Wait: true, + Done: func() { m.setCheckpoint(c.ID, ts) }, }) } diff --git a/histogram_test.go b/histogram_test.go index b4e8033..dc16aec 100644 --- a/histogram_test.go +++ b/histogram_test.go @@ -153,7 +153,7 @@ func TestHistogramHandler(t *testing.T) { {"_time":"2000-01-01T00:00:00Z","severity":"info","hits":"9"} `, recent.Format(time.RFC3339), recent.Format(time.RFC3339), older.Format(time.RFC3339)) }} - store := NewStore("http://vl", 10, 10, time.Second) + store := NewStore("http://vl", 10, 10, time.Second, nil) store.streamClient.Transport = vl a := &API{store: store} @@ -207,7 +207,7 @@ func TestHistogramDayInParis(t *testing.T) { {"_time":"2026-09-21T20:00:00Z","severity":"info","hits":"100"} ` }} - store := NewStore("http://vl", 10, 10, time.Second) + store := NewStore("http://vl", 10, 10, time.Second, nil) store.streamClient.Transport = vl a := &API{store: store} from := time.Date(2026, 9, 10, 0, 0, 0, 0, time.UTC) diff --git a/hostlogs.go b/hostlogs.go index 00d1493..2cca358 100644 --- a/hostlogs.go +++ b/hostlogs.go @@ -327,6 +327,7 @@ func (h *HostLogs) readFile(path, name string, off int64) (int64, error) { } e.SourceType = "host" e.Extra = map[string]string{"log_file": "/var/log/" + name} + e.Wait = true h.sink(e) h.count(1, 0) } @@ -400,6 +401,7 @@ func (h *HostLogs) emitJournal(je *journalEntry) { Proto: "journal", SourceType: "host", Extra: map[string]string{"unit": f["_SYSTEMD_UNIT"]}, + Wait: true, }) } diff --git a/main.go b/main.go index ae4d058..32db9e8 100644 --- a/main.go +++ b/main.go @@ -41,6 +41,7 @@ type config struct { batchSize int queueSize int flushEvery time.Duration + spoolMax int64 } func getenv(key, def string) string { @@ -57,6 +58,14 @@ func getenvInt(key string, def int) int { return def } +// getenvIntZero is getenvInt that also accepts 0 (to turn a feature off). +func getenvIntZero(key string, def int) int { + if v, err := strconv.Atoi(os.Getenv(key)); err == nil && v >= 0 { + return v + } + return def +} + func getenvBool(key string, def bool) bool { switch strings.ToLower(os.Getenv(key)) { case "1", "true", "yes", "on": @@ -105,6 +114,7 @@ func main() { batchSize: getenvInt("BATCH_SIZE", 1000), queueSize: getenvInt("QUEUE_SIZE", 100000), flushEvery: time.Duration(getenvInt("FLUSH_MS", 1000)) * time.Millisecond, + spoolMax: int64(getenvIntZero("SPOOL_MAX_MB", 1024)) << 20, } cfg.auth.dataDir = cfg.dataDir @@ -117,7 +127,14 @@ func main() { log.Fatalf("tags: %v", err) } - store := NewStore(cfg.vlogsURL, cfg.batchSize, cfg.queueSize, cfg.flushEvery) + var spool *Spool + if cfg.spoolMax > 0 { + if spool, err = OpenSpool(filepath.Join(cfg.dataDir, "spool"), cfg.spoolMax); err != nil { + log.Printf("disk buffer disabled: %v", err) + spool = nil + } + } + store := NewStore(cfg.vlogsURL, cfg.batchSize, cfg.queueSize, cfg.flushEvery, spool) storeDone := make(chan struct{}) go func() { store.Run(ctx) @@ -128,12 +145,15 @@ func main() { rdns := NewReverseDNS(cfg.rdns, cfg.dnsServer) sink := func(e *Entry) { // Host sent as an IP (or no host in the header): replace it with its DNS name. - // A new IP waits at most 300 ms; slower lookups finish in the background. - if name := rdns.Lookup(e.Host, 300*time.Millisecond); name != "" { - e.HostIP, e.Host = e.Host, name - } - store.Enqueue(e) - hub.Publish(e) + // A new IP waits at most 300 ms, without holding up the listener; slower + // lookups finish in the background. + rdns.Resolve(e.Host, 300*time.Millisecond, func(name string) { + if name != "" { + e.HostIP, e.Host = e.Host, name + } + store.Enqueue(e) + hub.Publish(e) + }) } // Listening errors (port already used…) are shown in Settings > Sources. syslogSrv := NewSyslogServer(ctx, cfg.syslogAddr, getenv("SYSLOG_PUBLIC_PORT", ""), cfg.dataDir, sink) diff --git a/rdns.go b/rdns.go index 4d7157c..958e987 100644 --- a/rdns.go +++ b/rdns.go @@ -1,6 +1,7 @@ package main import ( + "container/list" "context" "net" "strings" @@ -15,17 +16,26 @@ type ReverseDNS struct { posTTL time.Duration // cache duration of a found name negTTL time.Duration // cache duration of "no name" + lookups chan struct{} // limits the lookups running at once + waiters chan struct{} // limits the messages waiting for a lookup (Resolve) + mu sync.Mutex - cache map[string]*rdnsEntry + cache map[string]*list.Element // ip -> element of lru + lru *list.List // *rdnsEntry, most recently used first } type rdnsEntry struct { + ip string name string expires time.Time done chan struct{} // closed once the lookup has finished } -const rdnsMaxEntries = 10000 +const ( + rdnsMaxEntries = 10000 + rdnsMaxLookups = 64 + rdnsMaxWaiters = 1024 +) // NewReverseDNS uses the system resolver, or `server` ("ip" or "ip:port") when set. func NewReverseDNS(enabled bool, server string) *ReverseDNS { @@ -47,7 +57,10 @@ func NewReverseDNS(enabled bool, server string) *ReverseDNS { r: r, posTTL: time.Hour, negTTL: 10 * time.Minute, - cache: make(map[string]*rdnsEntry), + lookups: make(chan struct{}, rdnsMaxLookups), + waiters: make(chan struct{}, rdnsMaxWaiters), + cache: make(map[string]*list.Element), + lru: list.New(), } } @@ -60,6 +73,38 @@ func isClosed(ch chan struct{}) bool { } } +// entry returns the cache entry of ip, starting its lookup when it is missing +// or expired, or nil when too many lookups are already running (a flood of +// unknown addresses). The least recently used entry makes room for a new one. +func (d *ReverseDNS) entry(ip string) *rdnsEntry { + d.mu.Lock() + defer d.mu.Unlock() + if el := d.cache[ip]; el != nil { + e := el.Value.(*rdnsEntry) + if !isClosed(e.done) || time.Now().Before(e.expires) { + d.lru.MoveToFront(el) + return e + } + } + select { + case d.lookups <- struct{}{}: + default: + return nil + } + if el := d.cache[ip]; el != nil { + d.lru.Remove(el) + } + for d.lru.Len() >= rdnsMaxEntries { + old := d.lru.Back() + d.lru.Remove(old) + delete(d.cache, old.Value.(*rdnsEntry).ip) + } + e := &rdnsEntry{ip: ip, done: make(chan struct{})} + d.cache[ip] = d.lru.PushFront(e) + go d.resolve(ip, e) + return e +} + // Lookup returns the name of ip, or "" when ip is not an IP address, has no // PTR record, or is not resolved within `wait`. A lookup that takes longer // keeps running in the background and fills the cache for later calls. @@ -67,18 +112,10 @@ func (d *ReverseDNS) Lookup(ip string, wait time.Duration) string { if !d.enabled || net.ParseIP(ip) == nil { return "" } - d.mu.Lock() - e := d.cache[ip] - if e == nil || (isClosed(e.done) && time.Now().After(e.expires)) { - if len(d.cache) >= rdnsMaxEntries { - d.cache = make(map[string]*rdnsEntry) - } - e = &rdnsEntry{done: make(chan struct{})} - d.cache[ip] = e - go d.resolve(ip, e) + e := d.entry(ip) + if e == nil { + return "" } - d.mu.Unlock() - if !isClosed(e.done) { timer := time.NewTimer(wait) defer timer.Stop() @@ -91,6 +128,43 @@ func (d *ReverseDNS) Lookup(ip string, wait time.Duration) string { return e.name } +// Resolve is Lookup without blocking the caller: fn gets the name (or "") at +// once when it is known, otherwise from a goroutine after at most `wait`. The +// syslog listeners use it so that a slow DNS server never delays the reading +// of the next messages. +func (d *ReverseDNS) Resolve(ip string, wait time.Duration, fn func(name string)) { + if !d.enabled || net.ParseIP(ip) == nil { + fn("") + return + } + e := d.entry(ip) + switch { + case e == nil: + fn("") + return + case isClosed(e.done): + fn(e.name) + return + } + select { + case d.waiters <- struct{}{}: + default: + fn("") // too many messages waiting already + return + } + go func() { + defer func() { <-d.waiters }() + timer := time.NewTimer(wait) + defer timer.Stop() + select { + case <-e.done: + fn(e.name) + case <-timer.C: + fn("") + } + }() +} + // LookupMany resolves several addresses in parallel; unresolved ones are absent. func (d *ReverseDNS) LookupMany(ips []string, wait time.Duration) map[string]string { out := make(map[string]string) @@ -112,6 +186,7 @@ func (d *ReverseDNS) LookupMany(ips []string, wait time.Duration) map[string]str } func (d *ReverseDNS) resolve(ip string, e *rdnsEntry) { + defer func() { <-d.lookups }() ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() ttl := d.negTTL diff --git a/rdns_test.go b/rdns_test.go new file mode 100644 index 0000000..2e2e94c --- /dev/null +++ b/rdns_test.go @@ -0,0 +1,77 @@ +package main + +import ( + "context" + "errors" + "net" + "testing" + "time" +) + +// slowDNS never answers before the context ends. +func slowDNS() *ReverseDNS { + d := NewReverseDNS(true, "") + d.r = &net.Resolver{PreferGo: true, Dial: func(ctx context.Context, _, _ string) (net.Conn, error) { + <-ctx.Done() + return nil, errors.New("timeout") + }} + return d +} + +func TestResolveDoesNotBlock(t *testing.T) { + d := slowDNS() + got := make(chan string, 1) + start := time.Now() + d.Resolve("192.0.2.1", 50*time.Millisecond, func(name string) { got <- name }) + if time.Since(start) > 20*time.Millisecond { + t.Fatal("Resolve waited for the DNS server") + } + if name := <-got; name != "" { + t.Errorf("name %q, want none", name) + } + + // Not an IP address: the callback runs at once. + called := false + d.Resolve("router", time.Second, func(name string) { called = name == "" }) + if !called { + t.Error("callback not called synchronously for a host name") + } +} + +func TestReverseDNSLimits(t *testing.T) { + d := slowDNS() + // Lookups beyond the limit are not started (and not cached). + for i := 0; i < rdnsMaxLookups+10; i++ { + d.Lookup(net.IPv4(10, 0, byte(i>>8), byte(i)).String(), 0) + } + if n := d.lru.Len(); n != rdnsMaxLookups { + t.Errorf("%d entries, want %d", n, rdnsMaxLookups) + } +} + +func TestReverseDNSEvictsLeastRecentlyUsed(t *testing.T) { + d := slowDNS() + done := make(chan struct{}) + close(done) + add := func(ip string) { + e := &rdnsEntry{ip: ip, name: ip + ".lan", expires: time.Now().Add(time.Hour), done: done} + d.cache[ip] = d.lru.PushFront(e) + } + for i := 0; i < rdnsMaxEntries; i++ { + add(net.IPv4(10, 1, byte(i>>8), byte(i)).String()) + } + first := net.IPv4(10, 1, 0, 0).String() + if name := d.Lookup(first, 0); name != first+".lan" { // now the most recent + t.Fatalf("cached name %q", name) + } + d.entry("192.0.2.9") + if d.lru.Len() != rdnsMaxEntries { + t.Errorf("%d entries, want %d", d.lru.Len(), rdnsMaxEntries) + } + if d.cache[first] == nil { + t.Error("recently used entry evicted") + } + if d.cache[net.IPv4(10, 1, 0, 1).String()] != nil { + t.Error("least recently used entry kept") + } +} diff --git a/spool.go b/spool.go new file mode 100644 index 0000000..f288fc4 --- /dev/null +++ b/spool.go @@ -0,0 +1,136 @@ +package main + +import ( + "fmt" + "os" + "path/filepath" + "sort" + "strconv" + "strings" + "sync" + "time" +) + +// Spool keeps on disk the batches VictoriaLogs could not take, so that they +// are sent later instead of being lost. Each batch is one NDJSON file named +// -.ndjson; files are sent back oldest first. +type Spool struct { + dir string + max int64 // maximum total size in bytes + + mu sync.Mutex + size int64 // bytes on disk + lines int64 // messages on disk + seq int64 +} + +// errSpoolFull is returned when a batch does not fit within SPOOL_MAX_MB. +var errSpoolFull = fmt.Errorf("disk buffer full") + +// OpenSpool creates the directory if needed and counts the batches already +// there (left by a previous run). +func OpenSpool(dir string, max int64) (*Spool, error) { + if err := os.MkdirAll(dir, 0o755); err != nil { + return nil, err + } + s := &Spool{dir: dir, max: max} + files, err := s.files() + if err != nil { + return nil, err + } + for _, f := range files { + s.size += f.size + s.lines += f.lines + } + return s, nil +} + +type spoolFile struct { + path string + size int64 + lines int64 +} + +// files lists the batches on disk, oldest first. Unfinished writes (.tmp) +// are removed. +func (s *Spool) files() ([]spoolFile, error) { + entries, err := os.ReadDir(s.dir) + if err != nil { + return nil, err + } + var out []spoolFile + for _, e := range entries { + name := e.Name() + if strings.HasSuffix(name, ".tmp") { + _ = os.Remove(filepath.Join(s.dir, name)) + continue + } + base, ok := strings.CutSuffix(name, ".ndjson") + if !ok { + continue + } + _, n, _ := strings.Cut(base, "-") + lines, _ := strconv.ParseInt(n, 10, 64) + info, err := e.Info() + if err != nil { + continue + } + out = append(out, spoolFile{path: filepath.Join(s.dir, name), size: info.Size(), lines: lines}) + } + // The names start with a fixed-width timestamp: string order is time order. + sort.Slice(out, func(i, j int) bool { return out[i].path < out[j].path }) + return out, nil +} + +// Write saves one batch of `lines` messages. +func (s *Spool) Write(body []byte, lines int) error { + s.mu.Lock() + defer s.mu.Unlock() + if s.size+int64(len(body)) > s.max { + return errSpoolFull + } + s.seq++ + name := fmt.Sprintf("%020d%04d-%d.ndjson", time.Now().UnixNano(), s.seq%10000, lines) + path := filepath.Join(s.dir, name) + tmp := path + ".tmp" + if err := os.WriteFile(tmp, body, 0o644); err != nil { + _ = os.Remove(tmp) + return err + } + if err := os.Rename(tmp, path); err != nil { + _ = os.Remove(tmp) + return err + } + s.size += int64(len(body)) + s.lines += int64(lines) + return nil +} + +// Oldest returns the oldest batch, or ok=false when the spool is empty. +func (s *Spool) Oldest() (f spoolFile, body []byte, ok bool, err error) { + files, err := s.files() + if err != nil || len(files) == 0 { + return f, nil, false, err + } + f = files[0] + body, err = os.ReadFile(f.path) + return f, body, err == nil, err +} + +// Remove deletes a batch once VictoriaLogs has taken it. +func (s *Spool) Remove(f spoolFile) { + if err := os.Remove(f.path); err != nil && !os.IsNotExist(err) { + return + } + s.mu.Lock() + s.size -= f.size + s.lines -= f.lines + s.mu.Unlock() +} + +// Pending returns the number of messages and bytes waiting on disk. +func (s *Spool) Pending() (lines, size int64) { + s.mu.Lock() + defer s.mu.Unlock() + return s.lines, s.size +} diff --git a/store.go b/store.go index b7eafdd..2598a37 100644 --- a/store.go +++ b/store.go @@ -17,13 +17,18 @@ import ( ) // Store sends messages to VictoriaLogs in batches and queries it with LogsQL. +// Batches VictoriaLogs cannot take go to the disk spool (when enabled) and +// are sent again, oldest first, once it answers. type Store struct { base string client *http.Client streamClient *http.Client - in chan *Entry - batchSize int - flushEvery time.Duration + in chan *Entry + quit chan struct{} // closed on shutdown: unblocks waiting producers + batchSize int + flushEvery time.Duration + spool *Spool // nil: no disk buffer + spooled chan struct{} // wakes the replay loop up after a write to the spool received atomic.Int64 ingested atomic.Int64 @@ -31,29 +36,42 @@ type Store struct { lastErr atomic.Value // string } -func NewStore(base string, batchSize, queueSize int, flushEvery time.Duration) *Store { +func NewStore(base string, batchSize, queueSize int, flushEvery time.Duration, spool *Spool) *Store { s := &Store{ - base: strings.TrimRight(base, "/"), - client: &http.Client{Timeout: 60 * time.Second}, + base: strings.TrimRight(base, "/"), + client: &http.Client{Timeout: 60 * time.Second}, // No global timeout: a large export can take longer than a minute. The // request context still cancels it when the browser goes away. streamClient: &http.Client{}, - in: make(chan *Entry, queueSize), - batchSize: batchSize, - flushEvery: flushEvery, + in: make(chan *Entry, queueSize), + quit: make(chan struct{}), + batchSize: batchSize, + flushEvery: flushEvery, + spool: spool, + spooled: make(chan struct{}, 1), } s.lastErr.Store("") return s } -// Enqueue never blocks: when the queue is full, the message is counted as dropped. +// Enqueue adds a message to the queue. When the queue is full, a message +// whose producer can wait (Entry.Wait: Docker, host logs) blocks until there +// is room; any other one (syslog) is counted as dropped. func (s *Store) Enqueue(e *Entry) { s.received.Add(1) select { case s.in <- e: + return default: - s.dropped.Add(1) } + if e.Wait { + select { + case s.in <- e: + return + case <-s.quit: + } + } + s.dropped.Add(1) } // Run drains the queue into VictoriaLogs until the context is cancelled, @@ -62,20 +80,16 @@ func (s *Store) Run(ctx context.Context) { ticker := time.NewTicker(s.flushEvery) defer ticker.Stop() batch := make([]*Entry, 0, s.batchSize) + if s.spool != nil { + go s.replay(ctx) + } - flush := func() { - if len(batch) == 0 { - return + flush := func(ctx context.Context) { + if len(batch) > 0 { + s.store(ctx, batch) + clear(batch) + batch = batch[:0] } - 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 { @@ -83,20 +97,24 @@ func (s *Store) Run(ctx context.Context) { case e := <-s.in: batch = append(batch, e) if len(batch) >= s.batchSize { - flush() + flush(ctx) } case <-ticker.C: - flush() + flush(ctx) case <-ctx.Done(): + close(s.quit) + // The last batches get one short attempt, then go to the spool. + end, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() for { select { case e := <-s.in: batch = append(batch, e) if len(batch) >= s.batchSize { - flush() + flush(end) } default: - flush() + flush(end) return } } @@ -104,38 +122,158 @@ func (s *Store) Run(ctx context.Context) { } } -func (s *Store) insert(batch []*Entry) error { +// store sends one batch to VictoriaLogs, or to the spool when VictoriaLogs +// fails or older batches are still waiting there (to keep their order). +// Entry.Done is called once the batch is stored or spooled. +func (s *Store) store(ctx context.Context, batch []*Entry) { + body, err := encodeBatch(batch) + if err == nil { + err = s.save(ctx, body, len(batch)) + } + if err != nil { + s.dropped.Add(int64(len(batch))) + s.lastErr.Store(err.Error()) + log.Printf("victorialogs: %d messages dropped: %v", len(batch), err) + return + } + for _, e := range batch { + if e.Done != nil { + e.Done() + } + } +} + +func (s *Store) save(ctx context.Context, body []byte, lines int) error { + if s.spool == nil { + // No disk buffer: retry a few times, then give up. + var err error + for attempt := 0; attempt < 5; attempt++ { + if attempt > 0 && !sleepCtx(ctx, time.Duration(1< 0 { + log.Printf("spool: %d messages (%d bytes) waiting from a previous run", n, size) + } + backoff := time.Second + for { + f, body, ok, err := s.spool.Oldest() + if err != nil { + log.Printf("spool: %v", err) + } + if !ok { + select { + case <-ctx.Done(): + return + case <-s.spooled: + case <-time.After(time.Minute): + } + continue + } + err = s.post(ctx, body) + switch { + case err == nil: + s.spool.Remove(f) + s.ingested.Add(f.lines) + s.lastErr.Store("") + backoff = time.Second + continue + case isRejected(err): + // VictoriaLogs refuses the data itself: sending it again would not help. + s.spool.Remove(f) + s.dropped.Add(f.lines) + log.Printf("spool: %d messages refused by victorialogs: %v", f.lines, err) + continue + } + s.lastErr.Store(err.Error()) + if !sleepCtx(ctx, backoff) { + return + } + backoff = min(2*backoff, 30*time.Second) + } +} + +func encodeBatch(batch []*Entry) ([]byte, 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 + return nil, 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<= 400 && ie.status < 500 && ie.status != http.StatusTooManyRequests +} + +func (s *Store) post(ctx context.Context, body []byte) error { + u := s.base + "/insert/jsonline?_stream_fields=host,app&_msg_field=_msg&_time_field=_time" + req, err := http.NewRequestWithContext(ctx, http.MethodPost, u, bytes.NewReader(body)) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/stream+json") + resp, err := s.client.Do(req) if err != nil { return err } defer resp.Body.Close() if resp.StatusCode/100 != 2 { msg, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) - return fmt.Errorf("HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(msg))) + return &insertError{status: resp.StatusCode, msg: strings.TrimSpace(string(msg))} } _, _ = io.Copy(io.Discard, resp.Body) return nil @@ -265,11 +403,15 @@ func queryStatus(err error) int { } func (s *Store) Stats() map[string]any { - return map[string]any{ + st := map[string]any{ "received": s.received.Load(), "ingested": s.ingested.Load(), "dropped": s.dropped.Load(), "queue": len(s.in), "lastError": s.lastErr.Load(), } + if s.spool != nil { + st["spooled"], st["spoolBytes"] = s.spool.Pending() + } + return st } diff --git a/store_test.go b/store_test.go new file mode 100644 index 0000000..0062b39 --- /dev/null +++ b/store_test.go @@ -0,0 +1,216 @@ +package main + +import ( + "bytes" + "context" + "errors" + "io" + "net/http" + "os" + "path/filepath" + "strings" + "sync" + "sync/atomic" + "testing" + "time" +) + +// insertVL plays VictoriaLogs' insert endpoint: it answers `status` (0 = the +// connection fails) and keeps the lines it accepted. +type insertVL struct { + mu sync.Mutex + status int + lines []string +} + +func (v *insertVL) set(status int) { + v.mu.Lock() + v.status = status + v.mu.Unlock() +} + +func (v *insertVL) got() []string { + v.mu.Lock() + defer v.mu.Unlock() + return append([]string(nil), v.lines...) +} + +func (v *insertVL) RoundTrip(r *http.Request) (*http.Response, error) { + body, _ := io.ReadAll(r.Body) + v.mu.Lock() + defer v.mu.Unlock() + if v.status == 0 { + return nil, errors.New("connection refused") + } + if v.status == http.StatusOK { + for _, l := range strings.Split(strings.TrimSpace(string(body)), "\n") { + v.lines = append(v.lines, l) + } + } + return &http.Response{StatusCode: v.status, Body: io.NopCloser(strings.NewReader("")), Header: http.Header{}}, nil +} + +func testEntry(msg string, done *atomic.Int64) *Entry { + e := &Entry{Received: time.Now(), Time: time.Now(), Host: "h", App: "a", Message: msg} + if done != nil { + e.Done = func() { done.Add(1) } + } + return e +} + +func waitFor(t *testing.T, what string, cond func() bool) { + t.Helper() + for deadline := time.Now().Add(5 * time.Second); time.Now().Before(deadline); time.Sleep(10 * time.Millisecond) { + if cond() { + return + } + } + t.Fatalf("timed out waiting for %s", what) +} + +func TestSpoolFiles(t *testing.T) { + dir := t.TempDir() + sp, err := OpenSpool(dir, 100) + if err != nil { + t.Fatal(err) + } + if err := sp.Write([]byte("a\nb\n"), 2); err != nil { + t.Fatal(err) + } + if err := sp.Write([]byte("c\n"), 1); err != nil { + t.Fatal(err) + } + if err := sp.Write(bytes.Repeat([]byte("x"), 100), 1); !errors.Is(err, errSpoolFull) { + t.Fatalf("write over the limit: %v, want errSpoolFull", err) + } + _ = os.WriteFile(filepath.Join(dir, "broken.ndjson.tmp"), []byte("z"), 0o644) + + // A new run finds the batches left on disk and drops unfinished writes. + sp, err = OpenSpool(dir, 100) + if err != nil { + t.Fatal(err) + } + if n, size := sp.Pending(); n != 3 || size != 6 { + t.Fatalf("pending %d lines %d bytes, want 3 and 6", n, size) + } + f, body, ok, err := sp.Oldest() + if !ok || err != nil || string(body) != "a\nb\n" || f.lines != 2 { + t.Fatalf("oldest: %q %+v %v %v", body, f, ok, err) + } + sp.Remove(f) + if _, body, _, _ := sp.Oldest(); string(body) != "c\n" { + t.Fatalf("next oldest %q", body) + } + if _, err := os.Stat(filepath.Join(dir, "broken.ndjson.tmp")); !os.IsNotExist(err) { + t.Error("unfinished write left in the spool") + } +} + +func TestStoreSpoolsWhileVictoriaLogsIsDown(t *testing.T) { + sp, err := OpenSpool(t.TempDir(), 1<<20) + if err != nil { + t.Fatal(err) + } + vl := &insertVL{status: 0} + s := NewStore("http://vl", 2, 100, 20*time.Millisecond, sp) + s.client.Transport = vl + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go s.Run(ctx) + + var done atomic.Int64 + for _, m := range []string{"one", "two", "three"} { + s.Enqueue(testEntry(m, &done)) + } + // VictoriaLogs is down: the messages are kept on disk, and acknowledged. + waitFor(t, "3 spooled messages", func() bool { n, _ := sp.Pending(); return n == 3 }) + if done.Load() != 3 || s.dropped.Load() != 0 { + t.Fatalf("done %d dropped %d, want 3 and 0", done.Load(), s.dropped.Load()) + } + // While batches wait on disk, new ones queue behind them. + vl.set(http.StatusServiceUnavailable) + s.Enqueue(testEntry("four", &done)) + waitFor(t, "4 spooled messages", func() bool { n, _ := sp.Pending(); return n == 4 }) + + vl.set(http.StatusOK) + waitFor(t, "the spool to drain", func() bool { n, _ := sp.Pending(); return n == 0 }) + got := strings.Join(vl.got(), "\n") + for i, m := range []string{"one", "two", "three", "four"} { + if !strings.Contains(got, `"_msg":"`+m+`"`) { + t.Errorf("message %d %q not sent: %s", i, m, got) + } + } + if strings.Index(got, `"one"`) > strings.Index(got, `"four"`) { + t.Error("spooled batches sent out of order") + } + if s.ingested.Load() != 4 || s.lastErr.Load() != "" { + t.Errorf("ingested %d lastErr %q", s.ingested.Load(), s.lastErr.Load()) + } + + // Once the spool is empty, batches go straight to VictoriaLogs again. + s.Enqueue(testEntry("five", &done)) + waitFor(t, "a direct insert", func() bool { return s.ingested.Load() == 5 }) +} + +func TestStoreDropsRefusedSpooledBatch(t *testing.T) { + sp, _ := OpenSpool(t.TempDir(), 1<<20) + _ = sp.Write([]byte("{bad json\n"), 1) + vl := &insertVL{status: http.StatusBadRequest} + s := NewStore("http://vl", 10, 10, time.Second, sp) + s.client.Transport = vl + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go s.Run(ctx) + waitFor(t, "the refused batch to be dropped", func() bool { return s.dropped.Load() == 1 }) + if n, _ := sp.Pending(); n != 0 { + t.Errorf("%d messages left in the spool", n) + } +} + +func TestStoreWithoutSpoolDrops(t *testing.T) { + vl := &insertVL{status: 0} + s := NewStore("http://vl", 10, 10, time.Hour, nil) + s.client.Transport = vl + ctx, cancel := context.WithCancel(context.Background()) + cancel() // shutdown: one attempt, no retry pauses + var done atomic.Int64 + s.store(ctx, []*Entry{testEntry("x", &done)}) + if done.Load() != 0 || s.dropped.Load() != 1 { + t.Errorf("done %d dropped %d, want 0 and 1", done.Load(), s.dropped.Load()) + } +} + +func TestEnqueueWait(t *testing.T) { + s := NewStore("http://vl", 10, 1, time.Hour, nil) + s.Enqueue(testEntry("fills the queue", nil)) + s.Enqueue(testEntry("syslog: dropped", nil)) + if s.dropped.Load() != 1 { + t.Fatalf("dropped %d, want 1", s.dropped.Load()) + } + // A producer that can wait blocks until there is room… + queued := make(chan struct{}) + go func() { + e := testEntry("docker: waits", nil) + e.Wait = true + s.Enqueue(e) + close(queued) + }() + select { + case <-queued: + t.Fatal("did not wait for room in the queue") + case <-time.After(50 * time.Millisecond): + } + <-s.in + <-queued + if s.dropped.Load() != 1 { + t.Errorf("dropped %d, want 1", s.dropped.Load()) + } + // …or until shutdown. + close(s.quit) + e := testEntry("docker: shutdown", nil) + e.Wait = true + s.Enqueue(e) + if s.dropped.Load() != 2 { + t.Errorf("dropped %d after shutdown, want 2", s.dropped.Load()) + } +} diff --git a/syslog.go b/syslog.go index 9239510..17fdb38 100644 --- a/syslog.go +++ b/syslog.go @@ -38,6 +38,12 @@ type Entry struct { SourceType string // "syslog" or "docker" Extra map[string]string // additional fields (docker: container, image, …) + + // Not stored. Wait: the producer can be slowed down when the queue is full + // (Docker, host logs) instead of losing the message. Done: called once the + // message is stored in VictoriaLogs or in the disk spool. + Wait bool + Done func() } // Record returns the entry in the shape stored in VictoriaLogs and returned by the API. diff --git a/web/app.js b/web/app.js index 807ade5..fb0aea6 100644 --- a/web/app.js +++ b/web/app.js @@ -33,6 +33,7 @@ const I18N = { skipped: (n) => `${n} messages not shown (rate too high)`, stats: (s) => `received ${s.received} · stored ${s.ingested} · dropped ${s.dropped} · queue ${s.queue}`, storageErr: 'storage: ', + spooled: (n) => `${n} waiting on disk`, unreachable: 'server unreachable', fTime: 'timestamp', fMsg: 'message', fPid: 'pid', fSd: 'structured data', filterHost: 'Filter on this host', filterApp: 'Filter on this app', @@ -184,6 +185,7 @@ const I18N = { skipped: (n) => `${n} messages non affichés (débit trop élevé)`, stats: (s) => `reçus ${s.received} · stockés ${s.ingested} · perdus ${s.dropped} · file ${s.queue}`, storageErr: 'stockage : ', + spooled: (n) => `${n} en attente sur disque`, unreachable: 'serveur injoignable', fTime: 'horodatage', fMsg: 'message', fPid: 'pid', fSd: 'données structurées', filterHost: 'Filtrer sur cet hôte', filterApp: 'Filtrer sur cette appli', @@ -1664,6 +1666,7 @@ function renderStats() { const n = { received: fmtNum(s.received), ingested: fmtNum(s.ingested), dropped: '%DROPPED%', queue: fmtNum(s.queue) }; const dropped = s.dropped ? `${fmtNum(s.dropped)}` : fmtNum(0); let html = esc(t('stats', n)).replace('%DROPPED%', dropped); + if (s.spooled) html += ` · ${esc(t('spooled', fmtNum(s.spooled)))}`; if (s.lastError) html += ` · ${esc(t('storageErr') + s.lastError)}`; el.innerHTML = html; } diff --git a/web/style.css b/web/style.css index a54b7dd..bcf6cbc 100644 --- a/web/style.css +++ b/web/style.css @@ -452,6 +452,7 @@ mark.hit { background: var(--hit); color: inherit; border-radius: 3px; padding: .conn.ok::before { color: var(--sev-info); } .conn.ko::before { color: var(--sev-err); } #stats .bad { color: var(--sev-err); } +#stats .warn { color: var(--sev-warning); } /* "Back to top", at the right end of the status bar: never over a log row */ .to-top { width: 24px; height: 24px; margin: -5px -6px -5px auto; flex: none; } /* no taller bar */ .to-top svg { width: 16px; height: 16px; }