Files
cedricandClaude Opus 5.5 3504263992 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) <noreply@anthropic.com>
2026-10-03 16:26:18 +02:00

455 lines
11 KiB
Go

package main
import (
"bytes"
"context"
"encoding/json"
"errors"
"io"
"io/fs"
"log"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"sync"
"syscall"
"time"
"unicode/utf8"
)
// hostLogsConfig is saved in /data/hostlogs.json and edited in Settings > Sources.
type hostLogsConfig struct {
Enabled bool `json:"enabled"`
}
// hostPos is where reading resumes in a journal file (by file ID) or a text
// file (by name, with its inode to detect rotations). Off -1: archived file
// read to the end.
type hostPos struct {
Off int64 `json:"off"`
Ino uint64 `json:"ino,omitempty"`
}
// HostLogs collects the system logs of the machine hosting the stack: the
// systemd journal when there is one, else the text files of /var/log. The
// host directories are mounted read-only under root (/host by default).
type HostLogs struct {
root string
sink func(*Entry)
backfill time.Duration
cfgPath string
statePath string
wake chan struct{}
// Owned by the Run goroutine.
pos map[string]hostPos
dirty bool
started bool // a first scan was done since enabling: new files are read from their start
mu sync.Mutex
cfg hostLogsConfig
reset bool
mode string // "journal", "files" or "" (nothing found)
files int
read uint64
compressed uint64
errCode string
errDetail string
}
func NewHostLogs(root, dataDir string, backfill time.Duration, sink func(*Entry)) *HostLogs {
h := &HostLogs{
root: root,
sink: sink,
backfill: backfill,
cfgPath: filepath.Join(dataDir, "hostlogs.json"),
statePath: filepath.Join(dataDir, "hostlogs-state.json"),
wake: make(chan struct{}, 1),
pos: map[string]hostPos{},
}
if b, err := os.ReadFile(h.cfgPath); err == nil {
_ = json.Unmarshal(b, &h.cfg)
}
if b, err := os.ReadFile(h.statePath); err == nil {
_ = json.Unmarshal(b, &h.pos)
h.started = len(h.pos) > 0
}
return h
}
// Run polls the host logs every second while the source is enabled.
func (h *HostLogs) Run(ctx context.Context) {
tick := time.NewTicker(time.Second)
defer tick.Stop()
n := 0
for {
select {
case <-ctx.Done():
h.saveState()
return
case <-tick.C:
case <-h.wake:
}
h.mu.Lock()
enabled, reset := h.cfg.Enabled, h.reset
h.reset = false
h.mu.Unlock()
if reset {
// Disabled then enabled again: start over from HOST_LOGS_BACKFILL ago
// rather than reading everything written in between.
h.pos, h.started, h.dirty = map[string]hostPos{}, false, true
}
if enabled {
h.scan(time.Now())
}
if n++; n%5 == 0 || reset {
h.saveState()
}
}
}
func (h *HostLogs) saveState() {
if !h.dirty {
return
}
h.dirty = false
if err := writeJSONFile(h.statePath, h.pos); err != nil {
log.Printf("host logs: saving positions: %v", err)
}
}
func (h *HostLogs) setStatus(mode string, files int, code, detail string) {
h.mu.Lock()
h.mode, h.files, h.errCode, h.errDetail = mode, files, code, detail
h.mu.Unlock()
}
// scan reads what was added since the previous poll.
func (h *HostLogs) scan(now time.Time) {
notBefore := time.Time{}
if !h.started {
notBefore = now.Add(-h.backfill)
}
defer func() { h.started = true }()
var journals []string
for _, dir := range []string{"var/log/journal", "run/log/journal"} {
base := filepath.Join(h.root, dir)
for _, pat := range []string{"*.journal", "*/*.journal"} {
m, _ := filepath.Glob(filepath.Join(base, pat))
journals = append(journals, m...)
}
}
if len(journals) > 0 {
h.scanJournals(journals, notBefore)
return
}
h.scanFiles(notBefore.IsZero())
}
func (h *HostLogs) scanJournals(paths []string, notBefore time.Time) {
seen := map[string]bool{}
var firstErr error
for _, p := range paths {
f, err := os.Open(p)
if err != nil {
firstErr = keepFirst(firstErr, err)
continue
}
hdr, err := readJournalHeader(f)
f.Close()
if err != nil {
continue // file being created, or not a journal
}
key := "j:" + hdr.fileID
seen[key] = true
pos, known := h.pos[key]
if known && pos.Off < 0 {
continue
}
from, nb := pos.Off, time.Time{}
if !known {
if hdr.archived && !notBefore.IsZero() && time.UnixMicro(int64(hdr.tailRealtimeUS)).Before(notBefore) {
h.pos[key], h.dirty = hostPos{Off: -1}, true // archived before the backfill window
continue
}
nb = notBefore
}
next, hdr, err := readJournal(p, from, nb, h.emitJournal)
if err != nil {
firstErr = keepFirst(firstErr, err)
continue
}
if hdr.archived {
next = -1
}
if !known || next != pos.Off {
h.pos[key], h.dirty = hostPos{Off: next}, true
}
}
h.prune("j:", seen)
code, detail := "", ""
if firstErr != nil && len(seen) == 0 {
code, detail = errCode(firstErr), firstErr.Error()
}
h.setStatus("journal", len(seen), code, detail)
}
func keepFirst(first, err error) error {
if first != nil {
return first
}
return err
}
func errCode(err error) string {
if errors.Is(err, fs.ErrPermission) {
return "permission"
}
return "read"
}
func (h *HostLogs) prune(prefix string, seen map[string]bool) {
for k := range h.pos {
if strings.HasPrefix(k, prefix) && !seen[k] {
delete(h.pos, k)
h.dirty = true
}
}
}
// hostLogFile reports the classic text logs of /var/log (rotated and
// compressed copies are left out).
func hostLogFile(name string) bool {
return name == "syslog" || name == "messages" || strings.HasSuffix(name, ".log")
}
const maxHostRead = 4 << 20 // per file and per poll
// scanFiles follows the text files of /var/log, for hosts without journald.
// Files present at the first scan are read from their end (only new lines).
func (h *HostLogs) scanFiles(fromStart bool) {
dir := filepath.Join(h.root, "var/log")
ents, err := os.ReadDir(dir)
if err != nil {
code := errCode(err)
if errors.Is(err, fs.ErrNotExist) {
code = "not_mounted"
}
h.setStatus("", 0, code, err.Error())
return
}
seen := map[string]bool{}
var firstErr error
for _, de := range ents {
if !de.Type().IsRegular() || !hostLogFile(de.Name()) {
continue
}
name := de.Name()
key := "f:" + name
fi, err := de.Info()
if err != nil {
continue
}
var ino uint64
if st, ok := fi.Sys().(*syscall.Stat_t); ok {
ino = uint64(st.Ino)
}
pos, known := h.pos[key]
switch {
case !known && !fromStart:
pos = hostPos{Off: fi.Size(), Ino: ino}
case !known || pos.Ino != ino || fi.Size() < pos.Off:
pos = hostPos{Off: 0, Ino: ino} // new, rotated or truncated file
}
if fi.Size() > pos.Off {
off, err := h.readFile(filepath.Join(dir, name), name, pos.Off)
if err != nil {
firstErr = keepFirst(firstErr, err)
continue // not readable: not counted as followed
}
pos.Off = off
}
seen[key] = true
if old, ok := h.pos[key]; !ok || old != pos {
h.pos[key], h.dirty = pos, true
}
}
h.prune("f:", seen)
code, detail := "", ""
if firstErr != nil && len(seen) == 0 {
code, detail = errCode(firstErr), firstErr.Error()
}
mode := "files"
if len(seen) == 0 && code == "" {
mode, code = "", "empty"
}
h.setStatus(mode, len(seen), code, detail)
}
// readFile sends the complete lines written after off and returns the new offset.
func (h *HostLogs) readFile(path, name string, off int64) (int64, error) {
f, err := os.Open(path)
if err != nil {
return off, err
}
defer f.Close()
buf, err := io.ReadAll(io.NewSectionReader(f, off, maxHostRead))
if err != nil {
return off, err
}
end := bytes.LastIndexByte(buf, '\n')
if end < 0 {
if len(buf) < maxHostRead {
return off, nil // partial line: wait for its end
}
end = len(buf) - 1 // line longer than the read window: cut it
}
now := time.Now()
fac := fileFacility(name)
for _, line := range bytes.Split(buf[:end+1], []byte{'\n'}) {
if len(bytes.TrimSpace(line)) == 0 {
continue
}
e := ParseSyslog(line, "", "file", now)
if e.Host == "" {
e.Host = "localhost"
}
if n, ok := detectLevel(e.Message); ok {
e.SevNum, e.Severity = n, severityNames[n]
} else {
e.SevNum, e.Severity = 6, "info"
}
if fac >= 0 {
e.Facility = facilityNames[fac]
}
e.SourceType = "host"
e.Extra = map[string]string{"log_file": "/var/log/" + name}
e.Wait = true
h.sink(e)
h.count(1, 0)
}
return off + int64(end) + 1, nil
}
// fileFacility guesses the facility from the usual Debian/RHEL file names.
func fileFacility(name string) int {
switch strings.TrimSuffix(name, ".log") {
case "kern":
return 0
case "mail", "maillog":
return 2
case "daemon":
return 3
case "auth", "secure":
return 4
case "cron":
return 9
}
return -1
}
func (h *HostLogs) count(read, compressed uint64) {
h.mu.Lock()
h.read += read
h.compressed += compressed
h.mu.Unlock()
}
// emitJournal turns a journal entry into an Entry.
func (h *HostLogs) emitJournal(je *journalEntry) {
f := je.Fields
msg := f["MESSAGE"]
if msg == "" && je.Compressed > 0 {
msg = "(compressed journal entry: read it with journalctl)"
}
h.count(1, uint64(je.Compressed))
if strings.TrimSpace(msg) == "" {
return
}
if !utf8.ValidString(msg) {
msg = strings.ToValidUTF8(msg, "�")
}
sev := 6
if n, err := strconv.Atoi(f["PRIORITY"]); err == nil && n >= 0 && n <= 7 {
sev = n
}
fac := 1 // user
switch {
case f["_TRANSPORT"] == "kernel":
fac = 0
case f["_SYSTEMD_UNIT"] != "":
fac = 3 // daemon
}
if n, err := strconv.Atoi(f["SYSLOG_FACILITY"]); err == nil && n >= 0 && n < len(facilityNames) {
fac = n
}
app := firstNonEmpty(f["SYSLOG_IDENTIFIER"], f["_COMM"], strings.TrimSuffix(f["_SYSTEMD_UNIT"], ".service"))
host := firstNonEmpty(f["_HOSTNAME"], "localhost")
h.sink(&Entry{
Time: je.Realtime,
Received: time.Now(),
Host: host,
App: app,
ProcID: firstNonEmpty(f["SYSLOG_PID"], f["_PID"]),
Facility: facilityNames[fac],
Severity: severityNames[sev],
SevNum: sev,
Message: strings.TrimRight(msg, "\n"),
Proto: "journal",
SourceType: "host",
Extra: map[string]string{"unit": f["_SYSTEMD_UNIT"]},
Wait: true,
})
}
func firstNonEmpty(vals ...string) string {
for _, v := range vals {
if v != "" {
return v
}
}
return ""
}
// Configure saves and applies the configuration.
func (h *HostLogs) Configure(cfg hostLogsConfig) error {
h.mu.Lock()
if h.cfg.Enabled && !cfg.Enabled {
h.reset = true
h.mode, h.files, h.errCode, h.errDetail = "", 0, "", ""
}
h.cfg = cfg
h.mu.Unlock()
select {
case h.wake <- struct{}{}:
default:
}
return writeJSONFile(h.cfgPath, cfg)
}
// Status is the state shown in Settings > Sources.
func (h *HostLogs) Status() map[string]any {
h.mu.Lock()
defer h.mu.Unlock()
dirs := []string{}
for _, d := range []string{"var/log/journal", "run/log/journal", "var/log"} {
if _, err := os.Stat(filepath.Join(h.root, d)); err == nil {
dirs = append(dirs, "/"+d)
}
}
sort.Strings(dirs)
return map[string]any{
"enabled": h.cfg.Enabled,
"mode": h.mode,
"files": h.files,
"read": h.read,
"compressed": h.compressed,
"code": h.errCode,
"error": h.errDetail,
"mounted": dirs,
}
}