diff --git a/.env.example b/.env.example index 8fa42db..a48e89b 100644 --- a/.env.example +++ b/.env.example @@ -15,3 +15,7 @@ DNS_SERVER= ALLOW_PURGE=true # Maximum number of rows in a CSV export EXPORT_MAX=100000 +# Collect the logs of the Docker containers of this machine (on/off) +DOCKER_LOGS=on +# History read from a container seen for the first time (e.g. 30m, 1h, 24h; 0 = only new lines) +DOCKER_BACKFILL=1h diff --git a/README.md b/README.md index ed061a1..605acbb 100644 --- a/README.md +++ b/README.md @@ -59,6 +59,34 @@ Severity badges `err`/`crit` and `warning` use the colors of the `error` and `wa Shortcuts: `/` focuses the search box, `Esc` clears it. Clicking a row shows all its fields. +## Docker container logs + +Logstream also collects the logs of the Docker containers running on the machine where it is +installed (`DOCKER_LOGS=on`, the default in `docker-compose.yml`). They are searched, filtered, +colored and exported like syslog messages: + +- **host** is the Docker host name, **app** the compose service (or the container name), and + each log also carries `container`, `container_id`, `image`, `compose_project`, + `compose_service` and `stream` (stdout/stderr), visible in the row details. +- The **Source** filter shows only syslog or only Docker logs; Docker rows have a small cube + before the app name. +- The severity comes from the line itself when the application writes it: JSON + (`"level":"error"`), logfmt (`level=warn`), `[ERROR]`, or an upper-case level word at the + start of the line (`ERROR`, `WARN`…). Otherwise it is `info`. Terminal color codes are removed. +- **Settings > Sources** lists the containers grouped by compose project, with a switch for + each one, "Enable all" / "Disable all", and whether new containers are followed + automatically (on by default). 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 + 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. + +**Security**: access to the Docker socket is equivalent to root on the machine. Logstream +therefore goes through [docker-socket-proxy](https://github.com/Tecnativa/docker-socket-proxy), +which only lets through listing containers, reading logs, events and engine info (`GET` only). + ## CSV export The **Export** button (next to the log count) downloads every stored log matching the @@ -136,6 +164,9 @@ are only known by your router or a local DNS (Pi-hole, AdGuard, Unbound…), set | `DNS_SERVER` | empty | DNS server for reverse lookups (`ip` or `ip:port`) | | `ALLOW_PURGE` | `true` | allow "Delete all logs" in Settings | | `EXPORT_MAX` | `100000` | maximum number of rows in a CSV export | +| `DOCKER_LOGS` | `on` in compose | collect the logs of the local Docker containers | +| `DOCKER_HOST` | `tcp://docker-proxy:2375` in compose | Docker API address (`unix:///var/run/docker.sock` outside compose) | +| `DOCKER_BACKFILL` | `1h` | history read from a container seen for the first time | | `TZ` | `Europe/Paris` | time zone for RFC 3164 timestamps (which carry none) | | `BATCH_SIZE`, `FLUSH_MS`, `QUEUE_SIZE` | `1000`, `1000`, `100000` | ingestion tuning | @@ -163,6 +194,7 @@ are only known by your router or a local DNS (Pi-hole, AdGuard, Unbound…), set | `hub.go` | pushes new messages to browsers (SSE) | | `rdns.go` | cached reverse DNS lookups | | `export.go` | streamed CSV export | +| `docker.go` | Docker container logs (API, followers, positions, level detection) | | `tags.go` | color tag storage | | `api.go` | `/api/*` HTTP routes | | `web/` | UI (HTML, CSS, plain JavaScript, no build step), embedded in the binary; translations live in `web/app.js` (`I18N`) | diff --git a/api.go b/api.go index e3e86a7..b920e2e 100644 --- a/api.go +++ b/api.go @@ -19,6 +19,7 @@ type API struct { rdns *ReverseDNS allowPurge bool exportMax int + docker *DockerManager // nil when DOCKER_LOGS is off } func (a *API) Routes(mux *http.ServeMux) { @@ -36,6 +37,38 @@ func (a *API) Routes(mux *http.ServeMux) { mux.HandleFunc("GET /api/purge", a.purgeStatus) mux.HandleFunc("POST /api/purge", a.purge) mux.HandleFunc("GET /api/export.csv", a.exportCSV) + mux.HandleFunc("GET /api/docker", a.dockerStatus) + mux.HandleFunc("PUT /api/docker", a.dockerConfigure) +} + +// GET /api/docker: Docker source state and containers (Settings > Sources). +func (a *API) dockerStatus(w http.ResponseWriter, r *http.Request) { + if a.docker == nil { + writeJSON(w, http.StatusOK, map[string]any{"enabled": false}) + return + } + writeJSON(w, http.StatusOK, a.docker.Status()) +} + +// PUT /api/docker {"defaultEnabled": bool, "containers": {"project/service": bool}} +func (a *API) dockerConfigure(w http.ResponseWriter, r *http.Request) { + if a.docker == nil { + writeErr(w, http.StatusConflict, &codedError{code: "docker_off", msg: "the Docker source is disabled (DOCKER_LOGS=off)"}) + return + } + var body struct { + DefaultEnabled *bool `json:"defaultEnabled"` + Containers map[string]bool `json:"containers"` + } + if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 256*1024)).Decode(&body); err != nil { + writeErr(w, http.StatusBadRequest, err) + return + } + if err := a.docker.Configure(body.DefaultEnabled, body.Containers); err != nil { + writeErr(w, http.StatusInternalServerError, err) + return + } + writeJSON(w, http.StatusOK, a.docker.Status()) } func writeJSON(w http.ResponseWriter, status int, v any) { diff --git a/docker-compose.yml b/docker-compose.yml index 8c55c02..4cd4d19 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -5,6 +5,7 @@ services: restart: unless-stopped depends_on: - victorialogs + - docker-proxy ports: - "${SYSLOG_PORT:-514}:5514/udp" - "${SYSLOG_PORT:-514}:5514/tcp" @@ -18,8 +19,13 @@ services: DNS_SERVER: ${DNS_SERVER:-} # e.g. 192.168.1.1 to query your LAN DNS; empty = system resolver ALLOW_PURGE: ${ALLOW_PURGE:-true} EXPORT_MAX: ${EXPORT_MAX:-100000} + DOCKER_LOGS: ${DOCKER_LOGS:-on} # collect the logs of this machine's containers + DOCKER_HOST: tcp://docker-proxy:2375 # read-only Docker API gateway (below) + DOCKER_BACKFILL: ${DOCKER_BACKFILL:-1h} # history read from a container seen for the first time volumes: - - logstream-data:/data # tags.json + - logstream-data:/data # tags.json, docker.json (container choices) + labels: + logstream.exclude: "true" # never collect Logstream's own logs victorialogs: # Pin a specific version in production (see hub.docker.com/r/victoriametrics/victoria-logs/tags) @@ -37,6 +43,23 @@ services: # VictoriaLogs debug UI (http://localhost:9428/select/vmui), localhost only - "127.0.0.1:9428:9428" + docker-proxy: + # Read-only gateway to the Docker API: Logstream can only list containers, + # read their logs, receive events and engine info. Anything else (start, + # stop, exec, images, volumes…) is refused. + image: tecnativa/docker-socket-proxy:latest + container_name: logstream-docker-proxy + restart: unless-stopped + environment: + CONTAINERS: 1 + EVENTS: 1 + INFO: 1 + POST: 0 + volumes: + - /var/run/docker.sock:/var/run/docker.sock:ro + labels: + logstream.exclude: "true" # its access log would only echo Logstream's own requests + volumes: logstream-data: vlogs-data: diff --git a/docker.go b/docker.go new file mode 100644 index 0000000..eb80209 --- /dev/null +++ b/docker.go @@ -0,0 +1,633 @@ +package main + +import ( + "bufio" + "bytes" + "context" + "encoding/binary" + "encoding/json" + "errors" + "fmt" + "io" + "log" + "net" + "net/http" + "net/url" + "os" + "path/filepath" + "regexp" + "sort" + "strings" + "sync" + "time" + "unicode/utf8" +) + +// DockerContainer is a container as listed in Settings > Sources. +type DockerContainer struct { + ID string `json:"id"` + Key string `json:"key"` // stable across re-creations: "project/service", or the name + Name string `json:"name"` + Image string `json:"image"` + Project string `json:"project,omitempty"` + Service string `json:"service,omitempty"` + State string `json:"state"` // running, exited, … + Enabled bool `json:"enabled"` + Locked bool `json:"locked"` // Logstream itself, or label logstream.exclude=true + Following bool `json:"following"` +} + +// dockerConfig is saved in /data/docker.json and edited from the UI. +type dockerConfig struct { + DefaultEnabled bool `json:"defaultEnabled"` + Containers map[string]bool `json:"containers"` // key -> enabled, overrides the default +} + +type follower struct{ cancel context.CancelFunc } + +// DockerManager follows the logs of the containers of the local Docker engine +// (through the Docker API) and sends each line to the sink as an Entry. +type DockerManager struct { + client *http.Client + base string + sink func(*Entry) + backfill time.Duration + selfID string + cfgPath string + statePath string + syncCh chan struct{} + + mu sync.Mutex + 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 + dirty bool + connected bool + lastErr string +} + +func NewDockerManager(host, dataDir string, backfill time.Duration, sink func(*Entry)) (*DockerManager, error) { + client, base, err := dockerHTTP(host) + if err != nil { + return nil, err + } + self, _ := os.Hostname() // in a container: the short container ID + m := &DockerManager{ + client: client, + base: base, + sink: sink, + backfill: backfill, + selfID: self, + cfgPath: filepath.Join(dataDir, "docker.json"), + statePath: filepath.Join(dataDir, "docker-state.json"), + syncCh: make(chan struct{}, 1), + cfg: dockerConfig{DefaultEnabled: true, Containers: map[string]bool{}}, + followers: map[string]*follower{}, + checkpoint: map[string]time.Time{}, + } + m.load() + return m, nil +} + +// dockerHTTP returns an HTTP client for DOCKER_HOST: unix:///path, tcp://host:port or http(s)://… +func dockerHTTP(host string) (*http.Client, string, error) { + switch { + case strings.HasPrefix(host, "unix://"): + path := strings.TrimPrefix(host, "unix://") + tr := &http.Transport{DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) { + var d net.Dialer + return d.DialContext(ctx, "unix", path) + }} + return &http.Client{Transport: tr}, "http://docker", nil + case strings.HasPrefix(host, "tcp://"): + return &http.Client{}, "http://" + strings.TrimPrefix(host, "tcp://"), nil + case strings.HasPrefix(host, "http://"), strings.HasPrefix(host, "https://"): + return &http.Client{}, strings.TrimRight(host, "/"), nil + } + return nil, "", fmt.Errorf("unsupported DOCKER_HOST %q", host) +} + +func (m *DockerManager) load() { + if b, err := os.ReadFile(m.cfgPath); err == nil { + var cfg dockerConfig + if err := json.Unmarshal(b, &cfg); err == nil { + if cfg.Containers == nil { + cfg.Containers = map[string]bool{} + } + m.cfg = cfg + } + } + if b, err := os.ReadFile(m.statePath); err == nil { + var st map[string]string + if err := json.Unmarshal(b, &st); err == nil { + for id, s := range st { + if t, err := time.Parse(time.RFC3339Nano, s); err == nil { + m.checkpoint[id] = t + } + } + } + } +} + +func writeJSONFile(path string, v any) error { + b, err := json.MarshalIndent(v, "", " ") + if err != nil { + return err + } + if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { + return err + } + tmp := path + ".tmp" + if err := os.WriteFile(tmp, b, 0o644); err != nil { + return err + } + return os.Rename(tmp, path) +} + +// Run keeps the followed containers in line with the Docker engine and the +// configuration until ctx is cancelled. +func (m *DockerManager) Run(ctx context.Context) { + go m.watchEvents(ctx) + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + save := time.NewTicker(5 * time.Second) + defer save.Stop() + m.sync(ctx) + for { + select { + case <-ctx.Done(): + m.saveState() + return + case <-ticker.C: + m.sync(ctx) + case <-m.syncCh: + m.sync(ctx) + case <-save.C: + m.saveState() + } + } +} + +func (m *DockerManager) requestSync() { + select { + case m.syncCh <- struct{}{}: + default: + } +} + +func (m *DockerManager) get(ctx context.Context, path string) (*http.Response, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, m.base+path, nil) + if err != nil { + return nil, err + } + resp, err := m.client.Do(req) + if err != nil { + return nil, err + } + if resp.StatusCode != http.StatusOK { + msg, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) + resp.Body.Close() + return nil, fmt.Errorf("docker %s: HTTP %d %s", strings.SplitN(path, "?", 2)[0], resp.StatusCode, strings.TrimSpace(string(msg))) + } + return resp, nil +} + +func (m *DockerManager) getJSON(ctx context.Context, path string, v any) error { + ctx, cancel := context.WithTimeout(ctx, 15*time.Second) + defer cancel() + resp, err := m.get(ctx, path) + if err != nil { + return err + } + defer resp.Body.Close() + return json.NewDecoder(resp.Body).Decode(v) +} + +func (m *DockerManager) setError(err error) { + m.mu.Lock() + defer m.mu.Unlock() + if m.lastErr != err.Error() { + log.Printf("docker: %v", err) + } + m.connected = false + m.lastErr = err.Error() +} + +// enabled tells whether a container's logs must be collected; the caller holds the lock. +func (m *DockerManager) enabled(c DockerContainer) bool { + if c.Locked { + return false + } + if v, ok := m.cfg.Containers[c.Key]; ok { + return v + } + return m.cfg.DefaultEnabled +} + +// sync lists the containers and starts or stops their followers. +func (m *DockerManager) sync(ctx context.Context) { + m.mu.Lock() + needInfo := m.hostName == "" + m.mu.Unlock() + if needInfo { + var info struct{ Name string } + if err := m.getJSON(ctx, "/info", &info); err == nil && info.Name != "" { + m.mu.Lock() + m.hostName = info.Name + m.mu.Unlock() + } + } + + var list []struct { + ID string `json:"Id"` + Names []string + Image string + State string + Labels map[string]string + } + if err := m.getJSON(ctx, "/containers/json?all=1", &list); err != nil { + if ctx.Err() == nil { + m.setError(err) + } + return + } + + m.mu.Lock() + defer m.mu.Unlock() + if !m.connected { + log.Printf("docker: connected, %d containers", len(list)) + } + m.connected, m.lastErr = true, "" + seen := map[string]bool{} + cs := make([]DockerContainer, 0, len(list)) + for _, c := range list { + dc := DockerContainer{ID: c.ID, Image: c.Image, State: c.State, + Project: c.Labels["com.docker.compose.project"], Service: c.Labels["com.docker.compose.service"]} + if len(c.Names) > 0 { + dc.Name = strings.TrimPrefix(c.Names[0], "/") + } + dc.Key = dc.Name + if dc.Project != "" && dc.Service != "" { + dc.Key = dc.Project + "/" + dc.Service + } + dc.Locked = c.Labels["logstream.exclude"] == "true" || (m.selfID != "" && strings.HasPrefix(c.ID, m.selfID)) + dc.Enabled = m.enabled(dc) + seen[c.ID] = true + + want := dc.Enabled && c.State == "running" + f, following := m.followers[c.ID] + switch { + case want && !following: + m.startFollower(ctx, dc) + case !want && following: + f.cancel() + delete(m.followers, c.ID) + } + cs = append(cs, dc) + } + for id, f := range m.followers { + if !seen[id] { + f.cancel() + delete(m.followers, id) + } + } + for id := range m.checkpoint { + if !seen[id] { + delete(m.checkpoint, id) + m.dirty = true + } + } + sort.Slice(cs, func(i, j int) bool { + if cs[i].Project != cs[j].Project { + return cs[i].Project < cs[j].Project + } + return cs[i].Name < cs[j].Name + }) + m.containers = cs +} + +// startFollower launches the log reader of a container; the caller holds the lock. +func (m *DockerManager) startFollower(parent context.Context, c DockerContainer) { + ctx, cancel := context.WithCancel(parent) + f := &follower{cancel: cancel} + m.followers[c.ID] = f + go func() { + defer func() { + cancel() + m.mu.Lock() + if m.followers[c.ID] == f { + delete(m.followers, c.ID) + } + m.mu.Unlock() + }() + m.follow(ctx, c) + }() +} + +// follow reads the logs of a container while it runs, reconnecting when the +// stream is cut (proxy timeout, Docker restart) without losing or repeating lines. +func (m *DockerManager) follow(ctx context.Context, c DockerContainer) { + for { + var inspect struct { + Config struct{ Tty bool } + State struct{ Running bool } + } + if err := m.getJSON(ctx, "/containers/"+c.ID+"/json", &inspect); err != nil { + if ctx.Err() != nil { + return + } + log.Printf("docker %s: %v", c.Name, err) + sleepCtx(ctx, 5*time.Second) + continue + } + if !inspect.State.Running { + return + } + err := m.stream(ctx, c, inspect.Config.Tty, m.since(c.ID)) + if ctx.Err() != nil { + return + } + if err != nil && !errors.Is(err, io.EOF) && !errors.Is(err, io.ErrUnexpectedEOF) { + log.Printf("docker %s: %v", c.Name, err) + } + sleepCtx(ctx, time.Second) + } +} + +func sleepCtx(ctx context.Context, d time.Duration) { + t := time.NewTimer(d) + defer t.Stop() + select { + case <-ctx.Done(): + case <-t.C: + } +} + +// since returns where to resume reading: just after the last line read, or +// DOCKER_BACKFILL ago for a container seen for the first time. +func (m *DockerManager) since(id string) time.Time { + m.mu.Lock() + defer m.mu.Unlock() + if t, ok := m.checkpoint[id]; ok { + return t.Add(time.Nanosecond) + } + return time.Now().Add(-m.backfill) +} + +func (m *DockerManager) setCheckpoint(id string, t time.Time) { + m.mu.Lock() + if t.After(m.checkpoint[id]) { + m.checkpoint[id] = t + m.dirty = true + } + m.mu.Unlock() +} + +func (m *DockerManager) saveState() { + m.mu.Lock() + if !m.dirty { + m.mu.Unlock() + return + } + st := make(map[string]string, len(m.checkpoint)) + for id, t := range m.checkpoint { + st[id] = t.UTC().Format(time.RFC3339Nano) + } + m.dirty = false + m.mu.Unlock() + if err := writeJSONFile(m.statePath, st); err != nil { + log.Printf("docker: saving positions: %v", err) + } +} + +func (m *DockerManager) stream(ctx context.Context, c DockerContainer, tty bool, since time.Time) error { + q := url.Values{ + "follow": {"1"}, + "stdout": {"1"}, + "stderr": {"1"}, + "timestamps": {"1"}, + "since": {fmt.Sprintf("%d.%09d", since.Unix(), since.Nanosecond())}, + } + resp, err := m.get(ctx, "/containers/"+c.ID+"/logs?"+q.Encode()) + if err != nil { + return err + } + defer resp.Body.Close() + emit := func(stream string, line []byte) { m.emit(c, stream, line) } + if tty { + return readLines(resp.Body, "stdout", emit) + } + return readMultiplexed(resp.Body, emit) +} + +// readLines reads a TTY container stream: plain lines. +func readLines(r io.Reader, stream string, emit func(string, []byte)) error { + br := bufio.NewReaderSize(r, 64*1024) + for { + line, err := br.ReadSlice('\n') + if len(line) > 0 { + emit(stream, line) + } + if err == bufio.ErrBufferFull { + continue + } + if err != nil { + return err + } + } +} + +// readMultiplexed reads a non-TTY container stream: frames with an 8-byte +// header (stream type, 3 zero bytes, big-endian payload size). +func readMultiplexed(r io.Reader, emit func(string, []byte)) error { + br := bufio.NewReaderSize(r, 64*1024) + hdr := make([]byte, 8) + names := [2]string{"stdout", "stderr"} + var partial [2][]byte + for { + if _, err := io.ReadFull(br, hdr); err != nil { + return err + } + size := int(binary.BigEndian.Uint32(hdr[4:8])) + idx := 0 + if hdr[0] == 2 { + idx = 1 + } + if size > maxFrame { + if _, err := io.CopyN(io.Discard, br, int64(size)); err != nil { + return err + } + continue + } + buf := make([]byte, size) + if _, err := io.ReadFull(br, buf); err != nil { + return err + } + data := append(partial[idx], buf...) + for { + i := bytes.IndexByte(data, '\n') + if i < 0 { + break + } + emit(names[idx], data[:i]) + data = data[i+1:] + } + if len(data) > maxFrame { + emit(names[idx], data) + data = nil + } + partial[idx] = append([]byte(nil), data...) + } +} + +// Terminal color and cursor codes, frequent in container output. +var ansiRe = regexp.MustCompile(`\x1b\[[0-9;?]*[ -/]*[@-~]`) + +func (m *DockerManager) emit(c DockerContainer, stream string, line []byte) { + now := time.Now() + ts := now + s := strings.TrimRight(string(line), "\r\n") + // With timestamps=1, Docker prefixes each line with its RFC 3339 timestamp. + if sp := strings.IndexByte(s, ' '); sp > 0 { + if t, err := time.Parse(time.RFC3339Nano, s[:sp]); err == nil { + ts, s = t, s[sp+1:] + } + } 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)) + } + if strings.TrimSpace(s) == "" { + return + } + sev := 6 // info + if n, ok := detectLevel(s); ok { + sev = n + } + app := c.Service + if app == "" { + app = c.Name + } + m.mu.Lock() + host := m.hostName + m.mu.Unlock() + if host == "" { + host = "docker" + } + short := c.ID + if len(short) > 12 { + short = short[:12] + } + m.sink(&Entry{ + Time: ts, + Received: now, + Host: host, + App: app, + SevNum: sev, + Severity: severityNames[sev], + Message: s, + Proto: "docker", + SourceType: "docker", + Extra: map[string]string{ + "container": c.Name, + "container_id": short, + "image": c.Image, + "compose_project": c.Project, + "compose_service": c.Service, + "stream": stream, + }, + }) +} + +// Log level written by the application: JSON ("level":"error"), logfmt +// (level=warn), [ERROR], or an upper-case level word near the start. +var ( + levelJSON = regexp.MustCompile(`(?i)"(?:level|lvl|severity|log\.level)"\s*:\s*"([a-z]+)"`) + levelLogfmt = regexp.MustCompile(`(?i)\b(?:level|lvl|severity)="?([a-z]+)`) + levelBracket = regexp.MustCompile(`(?i)\[\s*(emerg|alert|crit|critical|fatal|panic|err|error|warn|warning|notice|info|debug|trace)\s*\]`) + levelWord = regexp.MustCompile(`\b(EMERG|ALERT|CRIT|CRITICAL|FATAL|PANIC|ERR|ERROR|WARN|WARNING|NOTICE|INFO|DEBUG|TRACE)\b`) +) + +func levelNum(word string) (int, bool) { + switch strings.ToLower(word) { + case "emerg", "emergency": + return 0, true + case "alert": + return 1, true + case "crit", "critical", "fatal", "panic": + return 2, true + case "err", "error": + return 3, true + case "warn", "warning": + return 4, true + case "notice": + return 5, true + case "info", "information": + return 6, true + case "debug", "trace": + return 7, true + } + return 0, false +} + +func detectLevel(s string) (int, bool) { + for _, re := range []*regexp.Regexp{levelJSON, levelLogfmt, levelBracket} { + if m := re.FindStringSubmatch(s); m != nil { + if n, ok := levelNum(m[1]); ok { + return n, true + } + } + } + head := s + if len(head) > 60 { + head = head[:60] + } + if m := levelWord.FindStringSubmatch(head); m != nil { + return levelNum(m[1]) + } + return 0, false +} + +// Status is the state shown in Settings > Sources. +func (m *DockerManager) Status() map[string]any { + m.mu.Lock() + defer m.mu.Unlock() + cs := make([]DockerContainer, len(m.containers)) + copy(cs, m.containers) + for i := range cs { + cs[i].Enabled = m.enabled(cs[i]) + _, cs[i].Following = m.followers[cs[i].ID] + } + return map[string]any{ + "enabled": true, + "connected": m.connected, + "error": m.lastErr, + "host": m.hostName, + "defaultEnabled": m.cfg.DefaultEnabled, + "containers": cs, + } +} + +// Configure changes the default and per-container choices, then applies them. +func (m *DockerManager) Configure(defaultEnabled *bool, containers map[string]bool) error { + m.mu.Lock() + if defaultEnabled != nil { + m.cfg.DefaultEnabled = *defaultEnabled + } + for k, v := range containers { + m.cfg.Containers[k] = v + } + cfg := dockerConfig{DefaultEnabled: m.cfg.DefaultEnabled, Containers: map[string]bool{}} + for k, v := range m.cfg.Containers { + cfg.Containers[k] = v + } + m.mu.Unlock() + m.requestSync() + return writeJSONFile(m.cfgPath, cfg) +} diff --git a/export.go b/export.go index 2d6b1a9..5ff7466 100644 --- a/export.go +++ b/export.go @@ -22,6 +22,9 @@ var exportColumns = []struct{ field, name string }{ {"procid", "pid"}, {"source", "source_ip"}, {"proto", "proto"}, + {"source_type", "source_type"}, + {"container", "container"}, + {"image", "image"}, {"_msg", "message"}, } diff --git a/main.go b/main.go index 6cad740..b10c5ab 100644 --- a/main.go +++ b/main.go @@ -35,6 +35,9 @@ type config struct { dnsServer string allowPurge bool exportMax int + dockerLogs bool + dockerHost string + backfill time.Duration batchSize int queueSize int flushEvery time.Duration @@ -64,6 +67,13 @@ func getenvBool(key string, def bool) bool { return def } +func getenvDuration(key string, def time.Duration) time.Duration { + if d, err := time.ParseDuration(os.Getenv(key)); err == nil && d >= 0 { + return d + } + return def +} + func main() { cfg := config{ syslogAddr: getenv("SYSLOG_ADDR", ":5514"), @@ -76,6 +86,9 @@ func main() { dnsServer: os.Getenv("DNS_SERVER"), allowPurge: getenvBool("ALLOW_PURGE", true), exportMax: getenvInt("EXPORT_MAX", 100000), + dockerLogs: getenvBool("DOCKER_LOGS", false), + dockerHost: getenv("DOCKER_HOST", "unix:///var/run/docker.sock"), + backfill: getenvDuration("DOCKER_BACKFILL", time.Hour), batchSize: getenvInt("BATCH_SIZE", 1000), queueSize: getenvInt("QUEUE_SIZE", 100000), flushEvery: time.Duration(getenvInt("FLUSH_MS", 1000)) * time.Millisecond, @@ -117,6 +130,16 @@ func main() { } mux := http.NewServeMux() api := &API{store: store, hub: hub, tags: tags, rdns: rdns, allowPurge: cfg.allowPurge, exportMax: cfg.exportMax} + if cfg.dockerLogs { + dm, err := NewDockerManager(cfg.dockerHost, cfg.dataDir, cfg.backfill, sink) + if err != nil { + log.Printf("docker logs disabled: %v", err) + } else { + go dm.Run(ctx) + api.docker = dm + log.Printf("docker logs enabled through %s", cfg.dockerHost) + } + } api.Routes(mux) mux.Handle("GET /", http.FileServer(http.FS(static))) diff --git a/query.go b/query.go index d05ea45..0b3fc7a 100644 --- a/query.go +++ b/query.go @@ -14,7 +14,8 @@ type Filter struct { Range string // 5m, 15m, 1h, 6h, 24h, 7d, 30d; anything else = no time limit Host string App string - Severity int // highest severity number included (0 = emerg … 7 = debug), -1 = all + Severity int // highest severity number included (0 = emerg … 7 = debug), -1 = all + Source string // "syslog", "docker" or "" for all } var validRanges = map[string]bool{"5m": true, "15m": true, "1h": true, "6h": true, "24h": true, "7d": true, "30d": true} @@ -27,6 +28,7 @@ func FilterFromRequest(r *http.Request) Filter { Range: q.Get("range"), Host: q.Get("host"), App: q.Get("app"), + Source: q.Get("source"), Severity: -1, } if v, err := strconv.Atoi(q.Get("severity")); err == nil && v >= 0 && v <= 7 { @@ -90,6 +92,12 @@ func (f Filter) filterExpr() string { if f.App != "" { parts = append(parts, "app:="+strconv.Quote(f.App)) } + switch f.Source { + case "docker": + parts = append(parts, `source_type:="docker"`) + case "syslog": // also matches logs stored before source_type existed + parts = append(parts, `!(source_type:="docker")`) + } if f.Severity >= 0 && f.Severity < 7 { names := make([]string, 0, 8) for i := 0; i <= f.Severity; i++ { @@ -173,6 +181,9 @@ func (m *Matcher) Match(e *Entry) bool { if m.f.App != "" && e.App != m.f.App { return false } + if (m.f.Source == "docker") != (e.SourceType == "docker") && m.f.Source != "" { + return false + } if m.f.Severity >= 0 && e.SevNum > m.f.Severity { return false } diff --git a/syslog.go b/syslog.go index 72ba7c4..707ee3c 100644 --- a/syslog.go +++ b/syslog.go @@ -35,6 +35,9 @@ type Entry struct { Source string // sender IP address Proto string // udp or tcp HostIP string // original host value when it was an IP resolved through DNS + + SourceType string // "syslog" or "docker" + Extra map[string]string // additional fields (docker: container, image, …) } // Record returns the entry in the shape stored in VictoriaLogs and returned by the API. @@ -65,6 +68,14 @@ func (e *Entry) Record() map[string]string { if e.HostIP != "" { r["host_ip"] = e.HostIP } + if e.SourceType != "" { + r["source_type"] = e.SourceType + } + for k, v := range e.Extra { + if v != "" { + r[k] = v + } + } return r } @@ -81,7 +92,7 @@ func ParseSyslog(raw []byte, source, proto string, now time.Time) *Entry { if !utf8.ValidString(s) { s = strings.ToValidUTF8(s, "\uFFFD") } - e := &Entry{Time: now, Received: now, Host: source, Source: source, Proto: proto} + e := &Entry{Time: now, Received: now, Host: source, Source: source, Proto: proto, SourceType: "syslog"} e.setPri(13) // user.notice, the RFC 3164 default if len(s) > 2 && s[0] == '<' { diff --git a/web/app.js b/web/app.js index a7257f8..828ea45 100644 --- a/web/app.js +++ b/web/app.js @@ -66,6 +66,19 @@ const I18N = { msgTimeTitle: 'Timestamp from the message', fHostName: 'host name (DNS)', fHostIp: 'host IP', clickHost: 'Click to filter on this host', clickApp: 'Click to filter on this app', + sourceAria: 'Source', srcAll: 'All sources', tabSources: 'Sources', + dockerTitle: 'Docker containers', + dockerDefault: 'Follow new containers automatically', + dockerAll: 'Enable all', dockerNone: 'Disable all', + dockerOff: 'The Docker source is disabled on this server. Set DOCKER_LOGS=on (see docker-compose.yml).', + dockerErr: 'Cannot reach Docker: ', + dockerOk: ({ h, n, f }) => `Connected to Docker on ${h}: ${n} containers, ${f} followed.`, + dockerStandalone: 'Standalone containers', + stFollowing: 'followed', stIgnored: 'running, ignored', stStopped: 'stopped', + stLocked: 'excluded', stLockedTitle: 'Logstream itself, or label logstream.exclude=true', + dockerHelp: 'Logs are read through docker-socket-proxy, a read-only gateway: Logstream can list containers and read their logs, nothing else. Choices apply per compose service (or container name), so they survive container re-creations.', + fContainer: 'container', fContainerId: 'container ID', fImage: 'image', fProject: 'compose project', + fService: 'compose service', fStream: 'stream', fSourceType: 'source', export: 'Export', exportCsvHint: 'Comma separated, UTF-8', exportExcel: 'CSV for Excel', exportExcelHint: 'Semicolon separated, for Excel in French', exportNote: (n) => `Up to ${n} logs matching the current filters, newest first. Dates use the time zone chosen in Settings.`, @@ -154,6 +167,19 @@ const I18N = { msgTimeTitle: 'Horodatage contenu dans le message', fHostName: 'nom d\'hôte (DNS)', fHostIp: 'IP de l\'hôte', clickHost: 'Cliquer pour filtrer sur cet hôte', clickApp: 'Cliquer pour filtrer sur cette appli', + sourceAria: 'Source', srcAll: 'Toutes les sources', tabSources: 'Sources', + dockerTitle: 'Conteneurs Docker', + dockerDefault: 'Suivre automatiquement les nouveaux conteneurs', + dockerAll: 'Tout activer', dockerNone: 'Tout désactiver', + dockerOff: 'La source Docker est désactivée sur ce serveur. Réglez DOCKER_LOGS=on (voir docker-compose.yml).', + dockerErr: 'Docker injoignable : ', + dockerOk: ({ h, n, f }) => `Connecté à Docker sur ${h} : ${n} conteneurs, ${f} suivis.`, + dockerStandalone: 'Conteneurs isolés', + stFollowing: 'suivi', stIgnored: 'actif, ignoré', stStopped: 'arrêté', + stLocked: 'exclu', stLockedTitle: 'Logstream lui-même, ou étiquette logstream.exclude=true', + dockerHelp: 'Les logs sont lus via docker-socket-proxy, une passerelle en lecture seule : Logstream peut lister les conteneurs et lire leurs logs, rien d\'autre. Les choix s\'appliquent par service compose (ou nom de conteneur), ils survivent donc à la recréation des conteneurs.', + fContainer: 'conteneur', fContainerId: 'ID du conteneur', fImage: 'image', fProject: 'projet compose', + fService: 'service compose', fStream: 'flux', fSourceType: 'source', export: 'Exporter', exportCsvHint: 'Séparateur virgule, UTF-8', exportExcel: 'CSV pour Excel', exportExcelHint: 'Séparateur point-virgule, pour Excel en français', exportNote: (n) => `Jusqu'à ${n} logs correspondant aux filtres, du plus récent au plus ancien. Dates dans le fuseau choisi dans Paramètres.`, @@ -362,6 +388,7 @@ function setLang(next) { renderTimeSettings(); renderPurge(); renderInterface(); + renderDocker(); } $('#langSwitch').addEventListener('click', (ev) => { @@ -498,7 +525,7 @@ function params() { if (q) p.set('q', q); p.set('mode', state.mode); p.set('range', $('#range').value); - for (const id of ['severity', 'host', 'app']) { + for (const id of ['severity', 'host', 'app', 'source']) { const v = $('#' + id).value; if (v !== '') p.set(id, v); } @@ -518,7 +545,7 @@ function syncControls() { btn.disabled = !liveAvailable(); btn.classList.toggle('on', state.live && liveAvailable()); btn.title = liveAvailable() ? t('liveTitle') : t('liveUnavailable'); - for (const id of ['severity', 'host', 'app']) $('#' + id).classList.toggle('set', $('#' + id).value !== ''); + for (const id of ['severity', 'host', 'app', 'source']) $('#' + id).classList.toggle('set', $('#' + id).value !== ''); } /* ================= Tags and highlighting ================= */ @@ -618,12 +645,14 @@ function rowHTML(r, isNew) { ? `${esc(r.host_name || r.host)}` : '') + (r.app - ? `${esc(r.app)}` + ? `${r.source_type === 'docker' ? DOCKER_ICON : ''}${esc(r.app)}` : '') + `
${highlight(r._msg)}
` + ''; } +const DOCKER_ICON = ''; + // Tooltip of the host cell: name and IP when both are known. function hostTitle(r) { const ip = r.host_ip || (r.host_name ? r.host : ''); @@ -700,12 +729,15 @@ function updateCount(tookMs) { } function detailsHTML(r) { - const order = ['received', 'msg_time', '_time', 'host', 'host_name', 'host_ip', 'app', 'procid', 'severity', 'facility', 'source', 'proto', 'sd', '_msg']; + const order = ['received', 'msg_time', '_time', 'host', 'host_name', 'host_ip', 'source_type', 'container', 'image', 'compose_project', 'compose_service', 'stream', 'app', 'procid', 'severity', 'facility', 'source', 'proto', 'sd', '_msg']; const hidden = new Set(['_stream_id', '_stream', 'sevnum']); if (r.msg_time) hidden.add('_time'); // same as the reception time const keys = order.filter((k) => r[k] != null && r[k] !== '' && !hidden.has(k)) .concat(Object.keys(r).filter((k) => !order.includes(k) && !hidden.has(k)).sort()); - const label = { msg_time: t('fTime'), host_name: t('fHostName'), host_ip: t('fHostIp'), received: t('fReceived'), _time: t('fTime'), _msg: t('fMsg'), procid: t('fPid'), sd: t('fSd') }; + const label = { + container: t('fContainer'), container_id: t('fContainerId'), image: t('fImage'), compose_project: t('fProject'), + compose_service: t('fService'), stream: t('fStream'), source_type: t('fSourceType'), + msg_time: t('fTime'), host_name: t('fHostName'), host_ip: t('fHostIp'), received: t('fReceived'), _time: t('fTime'), _msg: t('fMsg'), procid: t('fPid'), sd: t('fSd') }; const val = (k) => (k === '_time' || k === 'received' || k === 'msg_time' ? `${esc(fmtFull(r[k]))} (${esc(r[k])})` : esc(r[k])); @@ -847,7 +879,7 @@ $('#q').addEventListener('keydown', (ev) => { if (ev.key === 'Enter') refresh(); if (ev.key === 'Escape') { ev.target.value = ''; refresh(); } }); -for (const id of ['range', 'severity', 'host', 'app']) $('#' + id).addEventListener('change', onFilterChange); +for (const id of ['range', 'severity', 'host', 'app', 'source']) $('#' + id).addEventListener('change', onFilterChange); $('#mode').addEventListener('click', () => { state.mode = state.mode === 'simple' ? 'logsql' : 'simple'; @@ -1048,6 +1080,90 @@ $('#exportMenu').addEventListener('click', (ev) => { document.addEventListener('click', (ev) => { if (!ev.target.closest('.export')) setExportMenu(false); }); document.addEventListener('keydown', (ev) => { if (ev.key === 'Escape') setExportMenu(false); }); +/* ================= Docker source ================= */ + +let dockerState = null; + +async function loadDocker() { + try { dockerState = await api('/api/docker'); } catch (e) { dockerState = { enabled: true, connected: false, error: e.message }; } + renderDocker(); +} + +function renderDocker() { + const status = $('#dockerStatus'); + const body = $('#dockerBody'); + const d = dockerState; + if (!d) { status.textContent = ''; body.hidden = true; return; } + if (!d.enabled) { + status.textContent = t('dockerOff'); + status.className = 'docker-status muted'; + body.hidden = true; + return; + } + const cs = d.containers || []; + if (!d.connected) { + status.textContent = t('dockerErr') + (d.error || ''); + status.className = 'docker-status bad'; + } else { + status.textContent = t('dockerOk', { h: d.host || 'docker', n: cs.length, f: cs.filter((c) => c.following).length }); + status.className = 'docker-status good'; + } + body.hidden = !cs.length && !d.connected; + $('#dockerDefault').checked = !!d.defaultEnabled; + + // Group by compose project; standalone containers last. + const groups = new Map(); + for (const c of cs) { + const g = c.project || ''; + if (!groups.has(g)) groups.set(g, []); + groups.get(g).push(c); + } + const names = [...groups.keys()].sort((a, b) => (a === '') - (b === '') || a.localeCompare(b)); + let html = ''; + for (const g of names) { + html += `
${esc(g || t('dockerStandalone'))}
`; + for (const c of groups.get(g)) { + let st = 'stopped'; + let label = t('stStopped'); + if (c.locked) { st = 'locked'; label = t('stLocked'); } + else if (c.following) { st = 'following'; label = t('stFollowing'); } + else if (c.state === 'running') { st = 'ignored'; label = t('stIgnored'); } + html += ``; + } + html += '
'; + } + $('#dockerList').innerHTML = html; +} + +async function configureDocker(body) { + try { + dockerState = await api('/api/docker', { method: 'PUT', body }); + renderDocker(); + setTimeout(loadDocker, 1500); // followers start or stop in the background + } catch (e) { toast(e.message); } +} + +$('#dockerDefault').addEventListener('change', (ev) => configureDocker({ defaultEnabled: ev.target.checked })); +$('#dockerList').addEventListener('change', (ev) => { + const key = ev.target.dataset.key; + if (key) configureDocker({ containers: { [key]: ev.target.checked } }); +}); +for (const [id, on] of [['#dockerAll', true], ['#dockerNone', false]]) { + $(id).addEventListener('click', () => { + const containers = {}; + for (const c of (dockerState && dockerState.containers) || []) if (!c.locked) containers[c.key] = on; + configureDocker({ containers }); + }); +} +// Keep the list fresh while the Sources tab is open. +setInterval(() => { + if ($('#settingsDlg').open && !document.querySelector('[data-panel="sources"]').hidden) loadDocker(); +}, 4000); + /* ================= Purge ================= */ const purge = { allowed: true, running: 0, error: null, code: null, wasRunning: false, timer: null }; @@ -1207,6 +1323,7 @@ $('#settingsBtn').addEventListener('click', () => { renderTagList(); renderInterface(); loadPurgeStatus(); + loadDocker(); showSettingsTab(store.get('settingsTab', 'locale')); $('#settingsDlg').showModal(); }); diff --git a/web/index.html b/web/index.html index d916a75..f8ca1c2 100644 --- a/web/index.html +++ b/web/index.html @@ -66,6 +66,11 @@ +
@@ -121,6 +126,10 @@ Filters +
+ + +