Files
logstream/hostlogs_test.go
cedricandClaude Opus 5.5 30018e09e2 Host system logs source (systemd journal or /var/log)
New source, off by default and switched in Settings > Sources, that
collects the system logs of the machine hosting the stack:
- reads the systemd journal files directly (pure Go reader, no
  journalctl in the image), from /var/log/journal and /run/log/journal
  mounted read-only under /host;
- falls back to following the text files of /var/log (syslog,
  messages, *.log) on hosts without journald;
- positions saved in /data/hostlogs-state.json, HOST_LOGS_BACKFILL
  read when the source is turned on;
- source_type "host", selectable in the Source filter;
- compose mounts and group_add (HOST_LOGS_GID, adm by default), docs.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-10-03 10:05:38 +02:00

207 lines
6.4 KiB
Go

package main
import (
"encoding/binary"
"os"
"path/filepath"
"strings"
"testing"
"time"
)
// fakeJournal builds a minimal journal file: header, data objects, entries.
type fakeJournal struct {
compact bool
buf []byte
seq uint64
tail uint64
}
func newFakeJournal(compact bool) *fakeJournal {
j := &fakeJournal{compact: compact, buf: make([]byte, 272)}
copy(j.buf, jSignature)
if compact {
binary.LittleEndian.PutUint32(j.buf[12:], jIncompatCompact)
}
j.buf[16] = 1 // online
copy(j.buf[24:40], "0123456789abcdef")
binary.LittleEndian.PutUint64(j.buf[88:], 272)
return j
}
func (j *fakeJournal) object(typ, flags byte, body []byte) uint64 {
for len(j.buf)%8 != 0 {
j.buf = append(j.buf, 0)
}
off := uint64(len(j.buf))
h := make([]byte, jObjHeaderSize)
h[0], h[1] = typ, flags
binary.LittleEndian.PutUint64(h[8:], uint64(jObjHeaderSize+len(body)))
j.buf = append(append(j.buf, h...), body...)
j.tail = off
binary.LittleEndian.PutUint64(j.buf[136:], off)
return off
}
func (j *fakeJournal) data(field string, flags byte) uint64 {
n := 48
if j.compact {
n = 56
}
return j.object(jObjData, flags, append(make([]byte, n), field...))
}
// entry appends an entry; linked=false leaves it unfinished (tail seqnum not updated).
func (j *fakeJournal) entry(t time.Time, linked bool, fields ...string) {
var items []byte
for _, f := range fields {
flags := byte(0)
if strings.HasPrefix(f, "!") { // compressed data object
f, flags = f[1:], 4
}
off := j.data(f, flags)
if j.compact {
items = binary.LittleEndian.AppendUint32(items, uint32(off))
} else {
items = binary.LittleEndian.AppendUint64(items, off)
items = binary.LittleEndian.AppendUint64(items, 0)
}
}
j.seq++
body := make([]byte, 48)
binary.LittleEndian.PutUint64(body[0:], j.seq)
binary.LittleEndian.PutUint64(body[8:], uint64(t.UnixMicro()))
j.object(jObjEntry, 0, append(body, items...))
if linked {
binary.LittleEndian.PutUint64(j.buf[160:], j.seq)
binary.LittleEndian.PutUint64(j.buf[192:], uint64(t.UnixMicro()))
}
}
func (j *fakeJournal) write(t *testing.T, path string) {
t.Helper()
if err := os.WriteFile(path, j.buf, 0o644); err != nil {
t.Fatal(err)
}
}
func TestReadJournal(t *testing.T) {
for _, compact := range []bool{false, true} {
path := filepath.Join(t.TempDir(), "system.journal")
now := time.Now()
j := newFakeJournal(compact)
j.entry(now.Add(-2*time.Hour), true, "MESSAGE=too old", "PRIORITY=6")
j.entry(now.Add(-time.Minute), true, "MESSAGE=Started cron.", "PRIORITY=5", "SYSLOG_IDENTIFIER=systemd", "_HOSTNAME=srv1", "_PID=1")
j.write(t, path)
var got []*journalEntry
emit := func(e *journalEntry) { got = append(got, e) }
next, _, err := readJournal(path, 0, now.Add(-time.Hour), emit)
if err != nil {
t.Fatal(err)
}
if len(got) != 1 || got[0].Fields["MESSAGE"] != "Started cron." || got[0].Fields["_HOSTNAME"] != "srv1" {
t.Fatalf("compact=%v: got %+v", compact, got)
}
// A new complete entry and one still being written.
j.entry(now, true, "MESSAGE=second", "!MESSAGE=big")
j.entry(now, false, "MESSAGE=unfinished")
j.write(t, path)
got = nil
next2, _, err := readJournal(path, next, time.Time{}, emit)
if err != nil {
t.Fatal(err)
}
if len(got) != 1 || got[0].Fields["MESSAGE"] != "second" || got[0].Compressed != 1 {
t.Fatalf("compact=%v: resumed read got %d entries", compact, len(got))
}
// Once linked, the unfinished entry is read from where reading stopped.
binary.LittleEndian.PutUint64(j.buf[160:], j.seq)
j.write(t, path)
got = nil
if _, _, err := readJournal(path, next2, time.Time{}, emit); err != nil {
t.Fatal(err)
}
if len(got) != 1 || got[0].Fields["MESSAGE"] != "unfinished" {
t.Fatalf("compact=%v: unfinished entry got %+v", compact, got)
}
}
}
func TestHostLogsJournal(t *testing.T) {
root := t.TempDir()
dir := filepath.Join(root, "var/log/journal/machine")
if err := os.MkdirAll(dir, 0o755); err != nil {
t.Fatal(err)
}
j := newFakeJournal(true)
j.entry(time.Now(), true, "MESSAGE=Accepted publickey", "PRIORITY=6", "SYSLOG_FACILITY=10", "SYSLOG_IDENTIFIER=sshd", "_PID=42", "_HOSTNAME=srv1", "_SYSTEMD_UNIT=ssh.service")
j.entry(time.Now(), true, "MESSAGE=oops", "PRIORITY=3", "_TRANSPORT=kernel", "_HOSTNAME=srv1")
j.write(t, filepath.Join(dir, "system.journal"))
var got []*Entry
h := NewHostLogs(root, t.TempDir(), time.Hour, func(e *Entry) { got = append(got, e) })
h.scan(time.Now())
if len(got) != 2 {
t.Fatalf("got %d entries", len(got))
}
e := got[0]
if e.App != "sshd" || e.ProcID != "42" || e.Facility != "authpriv" || e.Severity != "info" || e.Host != "srv1" || e.SourceType != "host" || e.Extra["unit"] != "ssh.service" {
t.Fatalf("entry %+v", e)
}
if got[1].Facility != "kern" || got[1].Severity != "err" {
t.Fatalf("kernel entry %+v", got[1])
}
h.scan(time.Now()) // nothing new
if len(got) != 2 {
t.Fatalf("re-read: %d entries", len(got))
}
if st := h.Status(); st["mode"] != "journal" || st["files"] != 1 {
t.Fatalf("status %+v", st)
}
}
func TestHostLogsFiles(t *testing.T) {
root := t.TempDir()
dir := filepath.Join(root, "var/log")
if err := os.MkdirAll(dir, 0o755); err != nil {
t.Fatal(err)
}
auth := filepath.Join(dir, "auth.log")
if err := os.WriteFile(auth, []byte("Oct 3 07:00:00 srv1 sshd[1]: old line\n"), 0o644); err != nil {
t.Fatal(err)
}
_ = os.WriteFile(filepath.Join(dir, "auth.log.1"), []byte("rotated\n"), 0o644)
var got []*Entry
h := NewHostLogs(root, t.TempDir(), time.Hour, func(e *Entry) { got = append(got, e) })
h.scan(time.Now()) // existing content is skipped
if len(got) != 0 {
t.Fatalf("first scan read %d lines", len(got))
}
f, _ := os.OpenFile(auth, os.O_APPEND|os.O_WRONLY, 0)
_, _ = f.WriteString("Oct 3 08:00:00 srv1 sshd[7]: Failed password for root\nOct 3 08:00:01 srv1 sshd[7]: partial")
f.Close()
h.scan(time.Now())
if len(got) != 1 {
t.Fatalf("got %d lines", len(got))
}
e := got[0]
if e.Host != "srv1" || e.App != "sshd" || e.ProcID != "7" || e.Facility != "auth" || e.Extra["log_file"] != "/var/log/auth.log" || e.Proto != "file" {
t.Fatalf("entry %+v", e)
}
if st := h.Status(); st["mode"] != "files" || st["files"] != 1 {
t.Fatalf("status %+v", st)
}
}
func TestHostLogsNotMounted(t *testing.T) {
h := NewHostLogs(filepath.Join(t.TempDir(), "none"), t.TempDir(), time.Hour, func(*Entry) {})
h.scan(time.Now())
if st := h.Status(); st["code"] != "not_mounted" {
t.Fatalf("status %+v", st)
}
}