Compare commits
9
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f30c353b46 | ||
|
|
84f8b9f9ad | ||
|
|
30018e09e2 | ||
|
|
cb2c2c5200 | ||
|
|
fd1c0a2698 | ||
|
|
9c36e61c69 | ||
|
|
556d5ee767 | ||
|
|
eb51eb15b4 | ||
|
|
459d5cb79b |
No files matched your search
+19
-1
@@ -4,9 +4,21 @@ HTTP_PORT=8080
|
||||
TZ=Europe/Paris
|
||||
# How long logs are kept (e.g. 7d, 30d, 12w, 1y)
|
||||
RETENTION=30d
|
||||
# Web UI authentication (empty = disabled)
|
||||
# Web UI authentication: local (HTTP Basic below, or none) or oidc (OpenID Connect provider)
|
||||
AUTH_MODE=local
|
||||
# local mode: user and password (empty = no authentication)
|
||||
AUTH_USER=
|
||||
AUTH_PASS=
|
||||
# oidc mode: issuer URL exactly as the provider announces it
|
||||
# (Keycloak: https://sso.example.org/realms/<realm>, Authentik: https://auth.example.org/application/o/<slug>/)
|
||||
OIDC_ISSUER=
|
||||
OIDC_CLIENT_ID=
|
||||
OIDC_CLIENT_SECRET=
|
||||
# Callback URL of logstream, to register in the provider (path free, /auth/callback recommended)
|
||||
OIDC_REDIRECT_URL=https://logs.example.org/auth/callback
|
||||
# Requested scopes (openid is always added) and session lifetime (e.g. 8h, 24h)
|
||||
OIDC_SCOPES=openid profile email
|
||||
OIDC_SESSION_TTL=12h
|
||||
# Reverse DNS: show host names instead of IP addresses (on/off)
|
||||
RDNS=on
|
||||
# DNS server used for reverse lookups (e.g. your router: 192.168.1.1). Empty = system resolver
|
||||
@@ -19,3 +31,9 @@ EXPORT_MAX=100000
|
||||
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
|
||||
# Host system logs (enable them in Settings > Sources)
|
||||
# Group allowed to read the host logs: 4 = adm on Debian/Ubuntu; or the systemd-journal group
|
||||
# (getent group systemd-journal | cut -d: -f3)
|
||||
HOST_LOGS_GID=4
|
||||
# History read when the source is turned on (e.g. 30m, 1h, 24h; 0 = only new entries)
|
||||
HOST_LOGS_BACKFILL=1h
|
||||
+388
@@ -0,0 +1,388 @@
|
||||
[English](README.md) | Français
|
||||
|
||||
# Logstream
|
||||
|
||||
Un collecteur syslog simple : il reçoit les logs en UDP/TCP sur le port 514, les stocke dans
|
||||
[VictoriaLogs](https://docs.victoriametrics.com/victorialogs/) et sert une interface web
|
||||
soignée (thème clair/sombre, recherche instantanée, direct, tags de couleur, interface en
|
||||
anglais et en français).
|
||||
|
||||
```
|
||||
devices ──514 udp/tcp──▶ logstream (Go) ──HTTP batches──▶ VictoriaLogs
|
||||
▲ └── SSE (live) ──▶ browser
|
||||
└──── API / UI ◀─────────┘
|
||||
```
|
||||
|
||||
## Architecture
|
||||
|
||||

|
||||
|
||||
- **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.
|
||||
- **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 à
|
||||
VictoriaLogs.
|
||||
- **Frise** : `/api/histogram` (`histogram.go`) compte les logs par intervalle et par sévérité,
|
||||
alignés sur l'heure locale ; les bornes du zoom (`from`/`to`) s'appliquent aussi à la liste
|
||||
et à l'export.
|
||||
- **État** : les tags et les réglages des sources sont dans `/data` (volume `logstream-data`) ;
|
||||
les logs eux-mêmes dans le volume `vlogs-data`.
|
||||
|
||||
La source modifiable du schéma est
|
||||
[`docs/architecture.excalidraw`](docs/architecture.excalidraw)
|
||||
(à ouvrir sur [excalidraw.com](https://excalidraw.com)).
|
||||
|
||||
## Démarrage
|
||||
|
||||
```bash
|
||||
cp .env.example .env # optional
|
||||
docker compose up -d --build
|
||||
./tools/send-test-logs.sh # sends 100 test messages
|
||||
```
|
||||
|
||||
Ouvrez ensuite <http://localhost:8080>.
|
||||
|
||||
## Envoyer des logs
|
||||
|
||||
- **rsyslog** (Linux) : ajoutez `*.* @SERVER_IP:514` (UDP) ou `*.* @@SERVER_IP:514` (TCP)
|
||||
dans `/etc/rsyslog.d/90-logstream.conf`, puis lancez `systemctl restart rsyslog`.
|
||||
- **Équipements réseau, NAS, pare-feu** : indiquez l'IP du serveur et le port 514 dans leurs
|
||||
réglages « syslog distant » (remote syslog).
|
||||
- **Test manuel** : `logger -n 127.0.0.1 -P 514 -d "hello error"` (util-linux) ou
|
||||
`echo "<14>test ok" | nc -u -w1 127.0.0.1 514`.
|
||||
|
||||
**Paramètres > Sources > Syslog** active ou désactive la réception syslog et choisit les
|
||||
protocoles (UDP, TCP) sans redémarrage ; le choix est enregistré dans `/data/syslog.json`. Cet
|
||||
onglet affiche aussi l'état d'écoute (et l'erreur si le port est déjà utilisé). Le port
|
||||
lui-même est publié par docker-compose : modifiez `SYSLOG_PORT` dans `.env`, puis lancez
|
||||
`docker compose up -d`.
|
||||
|
||||
Formats pris en charge : RFC 3164 (BSD) et RFC 5424. En TCP, les deux découpages sont acceptés :
|
||||
« un message par ligne » et « octet counting » (RFC 6587).
|
||||
|
||||
## Recherche
|
||||
|
||||
**Mode simple** (par défaut) : chaque mot est cherché comme sous-chaîne, sans tenir compte de
|
||||
la casse, dans le message, l'hôte et l'application. Les mots sont combinés avec ET.
|
||||
|
||||
| Saisie | Signification |
|
||||
|---|---|
|
||||
| `error disk` | contient « error » **et** « disk » |
|
||||
| `"disk full"` | contient la phrase exacte |
|
||||
| `error -timeout` | contient « error » mais pas « timeout » |
|
||||
|
||||
**Mode LogsQL** (bouton `Simple` / `LogsQL`) : le langage de requête complet de VictoriaLogs,
|
||||
par exemple `error AND host:web-01`, `app:~"ssh|nginx"` ou `* | stats by (host) count()`.
|
||||
Le direct est désactivé dans ce mode.
|
||||
|
||||
Chaque ligne affiche, de gauche à droite : l'**heure de réception** (horloge du serveur),
|
||||
l'horodatage trouvé dans le message lui-même (`msg_time`), la sévérité, l'hôte, l'application
|
||||
et le message. Un clic sur un hôte ou une application filtre dessus.
|
||||
|
||||
Les logs sont indexés, recherchés et triés par **heure de réception** : les équipements dont
|
||||
l'horloge est fausse (par exemple des points d'accès dont le NTP échoue) apparaissent quand même
|
||||
dans la bonne plage de temps. Les logs stockés par les versions antérieures à ce changement sont
|
||||
indexés par l'horodatage de leur message ; purgez-les depuis les Paramètres pour repartir sur
|
||||
une base propre.
|
||||
|
||||
Les badges de sévérité `err`/`crit` et `warning` reprennent les couleurs des tags `error` et
|
||||
`warning`.
|
||||
|
||||
Raccourcis : `/` place le curseur dans la recherche, `Esc` la vide. Un clic sur une ligne
|
||||
affiche tous ses champs.
|
||||
|
||||
Les en-têtes de colonnes restent visibles pendant le défilement. Faites glisser le bord d'un
|
||||
en-tête (réception, heure message, sévérité, hôte, app) pour redimensionner la colonne, et
|
||||
double-cliquez dessus pour revenir à la largeur automatique ; le message occupe l'espace
|
||||
restant. Les largeurs sont mémorisées par le navigateur (**Paramètres > Interface >
|
||||
Réinitialiser les colonnes** les rétablit toutes). Sur téléphone, la liste garde sa
|
||||
présentation sur deux lignes, sans colonnes.
|
||||
|
||||
## Frise
|
||||
|
||||
La frise au-dessus de la liste montre le volume de logs par intervalle, compté selon l'**heure
|
||||
de réception** sur l'horloge du serveur (elle correspond donc aux heures affichées dans les
|
||||
lignes).
|
||||
|
||||
- Au survol d'un intervalle : ses bornes, son total et le détail par sévérité.
|
||||
- Un clic sur une barre zoome sur cet intervalle ; un glisser sur plusieurs barres zoome sur la
|
||||
sélection. La plage de temps affiche alors la période zoomée (« × Annuler le zoom » ou le
|
||||
choix d'une autre plage en sort) ; la liste, les compteurs et l'export CSV suivent le zoom,
|
||||
et le direct se met en pause.
|
||||
- En direct, le dernier intervalle grandit à l'arrivée des messages, et la frise se recharge à
|
||||
chaque nouvel intervalle. Rien n'est rafraîchi tant que l'onglet du navigateur est masqué ;
|
||||
la frise se met à jour dès qu'il redevient visible.
|
||||
|
||||
**Paramètres > Interface > Frise** (mémorisé par navigateur) : échelle (linéaire, √ par défaut,
|
||||
log), hauteur (S/M/L : 40/80/120 px), couleur (empilement par sévérité, intensité comparée à
|
||||
la médiane de la fenêtre : calme, rafale au-delà de 3×, anomalie au-delà de 10×, ou aucune),
|
||||
barres ou aire, division (automatique, environ 100 intervalles, ou fixe : 1 s, 10 s, 1 min,
|
||||
5 min, 1 h, 1 jour) et rafraîchissement (désactivé, 5 s, 15 s, 30 s, 1 min, ou à chaque nouvel
|
||||
intervalle). Une division fixe qui dépasserait 300 intervalles sur la plage est élargie
|
||||
(signalé par « élargi »). Les intervalles sont alignés sur l'heure locale du fuseau horaire
|
||||
choisi dans les Paramètres (les jours commencent à minuit, heure locale).
|
||||
|
||||
API : `curl 'localhost:8080/api/histogram?range=24h&step=auto&tz=Europe/Paris'` renvoie
|
||||
`step` (ms), `start` (ms), `count`, `now` (horloge du serveur) et les `buckets` non vides
|
||||
(`i` = numéro d'intervalle, `n` = total, `sev` = nombre par sévérité). Toutes les routes de
|
||||
logs acceptent aussi `from` / `to` (millisecondes Unix) à la place de `range`.
|
||||
|
||||
## Logs des conteneurs Docker
|
||||
|
||||
Logstream collecte aussi les logs des conteneurs Docker qui tournent sur la machine où il est
|
||||
installé (`DOCKER_LOGS=on`, le défaut dans `docker-compose.yml`). Ils sont recherchés, filtrés,
|
||||
colorés et exportés comme les messages syslog :
|
||||
|
||||
- **host** est le nom de l'hôte Docker, **app** le service compose (ou le nom du conteneur), et
|
||||
chaque log porte aussi `container`, `container_id`, `image`, `compose_project`,
|
||||
`compose_service` et `stream` (stdout/stderr), visibles dans le détail de la ligne.
|
||||
- Le filtre **Source** n'affiche que les logs syslog ou que les logs Docker ; les lignes Docker
|
||||
ont un petit cube devant le nom de l'application, de la couleur de son projet compose.
|
||||
- La sévérité vient de la ligne elle-même quand l'application l'écrit : JSON
|
||||
(`"level":"error"`), logfmt (`level=warn`), `[ERROR]`, ou un niveau en majuscules au début de
|
||||
la ligne (`ERROR`, `WARN`…). Sinon, elle vaut `info`. Les codes de couleur du terminal sont
|
||||
supprimés.
|
||||
- **Paramètres > Sources** affiche une étiquette par conteneur (`project/service`) dans un seul
|
||||
champ : les conteneurs suivis en couleur, puis les non suivis en gris ; un clic bascule une
|
||||
étiquette. La couleur identifie le projet compose, et les noms d'applications Docker de la
|
||||
liste utilisent la même couleur. Les conteneurs arrêtés sont masqués par défaut (« Afficher
|
||||
les conteneurs arrêtés » les montre, en pointillés, et ils restent modifiables). Une zone de
|
||||
filtre et « Tout activer » / « Tout désactiver » (appliqués aux étiquettes affichées) aident
|
||||
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
|
||||
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.
|
||||
|
||||
**Sécurité** : l'accès au socket Docker équivaut à un accès root sur la machine. Logstream passe
|
||||
donc par [docker-socket-proxy](https://github.com/Tecnativa/docker-socket-proxy), qui ne laisse
|
||||
passer que la liste des conteneurs, la lecture des logs, les événements et les informations du
|
||||
moteur (`GET` uniquement).
|
||||
|
||||
## Logs système de l'hôte
|
||||
|
||||
Logstream peut aussi collecter les logs système de la machine qui héberge la stack, sans rien
|
||||
configurer sur l'hôte. La source est **désactivée par défaut** : activez-la dans
|
||||
**Paramètres > Sources > Logs système de l'hôte** (enregistré dans `/data/hostlogs.json`).
|
||||
|
||||
- `docker-compose.yml` monte `/var/log` et `/run/log/journal` en lecture seule sous `/host`.
|
||||
- Si l'hôte utilise systemd, Logstream lit directement les fichiers du **journal systemd**
|
||||
(pas besoin de `journalctl` dans l'image) : **host** est le nom de la machine, **app** le
|
||||
programme (`SYSLOG_IDENTIFIER`), la sévérité et la facility viennent du journal, et l'unité
|
||||
systemd est conservée dans `unit`. Sinon, il suit les fichiers texte de `/var/log` (`syslog`,
|
||||
`messages`, `*.log`), analysés comme des lignes syslog, avec le nom du fichier dans `log_file`.
|
||||
- Ces logs ont la source `host` : le filtre **Source** permet de les afficher seuls.
|
||||
- À l'activation, la dernière heure est lue d'abord (`HOST_LOGS_BACKFILL`) ; la position
|
||||
atteinte est enregistrée dans `/data/hostlogs-state.json`, un redémarrage ne perd donc ni ne
|
||||
duplique d'entrées.
|
||||
- **Droits** : le conteneur tourne avec un utilisateur sans privilège et reçoit le groupe `adm`
|
||||
(gid 4), qui peut lire le journal et `/var/log` sur Debian et Ubuntu. Sur d'autres systèmes,
|
||||
réglez `HOST_LOGS_GID` dans `.env` sur le gid de `systemd-journal`
|
||||
(`getent group systemd-journal | cut -d: -f3`). Paramètres > Sources affiche un message clair
|
||||
si l'accès est refusé.
|
||||
- Limites : les champs du journal compressés par journald (messages de plus de 512 octets,
|
||||
compressés en zstd/lz4/xz) ne peuvent pas être décodés sans bibliothèque supplémentaire ; ils
|
||||
sont comptés dans les Paramètres et affichés comme « (compressed journal entry) ». Les fichiers
|
||||
texte tournés ou compressés (`*.1`, `*.gz`) ne sont pas lus.
|
||||
|
||||
## Export CSV
|
||||
|
||||
Le bouton **Exporter** (à côté du nombre de logs) télécharge tous les logs stockés qui
|
||||
correspondent aux filtres en cours (recherche, plage de temps, sévérité, hôte, application),
|
||||
du plus récent au plus ancien, jusqu'à `EXPORT_MAX` lignes (100 000 par défaut), et pas
|
||||
seulement les lignes à l'écran. Deux variantes :
|
||||
|
||||
- **CSV** : séparateur virgule, UTF-8.
|
||||
- **CSV pour Excel** : séparateur point-virgule avec un BOM UTF-8, pour qu'Excel en français
|
||||
l'ouvre directement avec les accents. Les cellules qui commencent par `=`, `+`, `-` ou `@`
|
||||
sont préfixées par `'` pour qu'un message de log forgé ne puisse pas s'exécuter comme une
|
||||
formule.
|
||||
|
||||
Colonnes : `received`, `message_time` (toutes deux au format `YYYY-MM-DD HH:MM:SS.mmm` dans le
|
||||
fuseau horaire choisi dans les Paramètres), `severity`, `facility`, `host`, `host_ip`, `app`,
|
||||
`pid`, `source_ip`, `proto`, `message`. En mode LogsQL, seule la partie filtre est prise en
|
||||
charge (pas de `| pipes`).
|
||||
|
||||
En ligne de commande : `curl -o logs.csv 'localhost:8080/api/export.csv?q=error&range=24h&tz=Europe/Paris'`.
|
||||
|
||||
## Paramètres
|
||||
|
||||
L'icône en forme d'engrenage ouvre les paramètres, organisés en onglets. Tout, sauf les tags de
|
||||
couleur, est mémorisé par navigateur.
|
||||
|
||||
- **Localisation**
|
||||
- *Langue* : anglais ou français.
|
||||
- *Date et heure* : fuseau horaire (celui du navigateur, UTC ou environ 80 fuseaux courants)
|
||||
et format d'affichage de l'heure de réception : `DD/MM/YYYY HH:MM:SS` (par défaut,
|
||||
l'affichage français habituel), avec millisecondes, `YYYY-MM-DD`, horloge sur 12 heures,
|
||||
ISO 8601 ou epoch Unix. Le fuseau horaire s'applique à toutes les dates affichées.
|
||||
- **Filtres** : tags de couleur. Chaque tag a un mot-clé, une couleur de fond (le texte passe
|
||||
automatiquement en noir ou en blanc pour rester lisible) et des options : mot entier, respect
|
||||
de la casse, expression régulière, actif. Les tags sont stockés sur le serveur dans
|
||||
`/data/tags.json` (volume `logstream-data`) : ils sont donc partagés par tous les navigateurs.
|
||||
Tags par défaut (pastel) : `warning` (orange), `error` (rouge), `ok` (vert). Les tags par
|
||||
défaut qui utilisent encore les couleurs des versions précédentes passent automatiquement aux
|
||||
couleurs pastel.
|
||||
- **Interface**
|
||||
- *Thème* : Système (suit la préférence de l'ordinateur ou du téléphone), Clair ou Sombre. Le
|
||||
bouton soleil/lune de l'en-tête bascule entre clair et sombre.
|
||||
- *Affichage des logs* : taille du texte (très petite, petite, moyenne, grande) et police :
|
||||
la police monospace du système, ou l'une des 12 polices libres conçues pour le texte dense
|
||||
(JetBrains Mono, Fira Code, Source Code Pro, IBM Plex Mono, Cascadia Code, Roboto Mono,
|
||||
Ubuntu Mono, Inconsolata, Red Hat Mono, Noto Sans Mono, Victor Mono, DM Mono). Elles sont
|
||||
chargées par le navigateur depuis [Bunny Fonts](https://fonts.bunny.net), un service
|
||||
européen de polices respectueux de la vie privée ; sans accès à internet, la police du
|
||||
système est utilisée. Les ligatures sont désactivées pour que `->` ou `!=` s'affichent tels
|
||||
quels.
|
||||
- **Données** : « Supprimer tous les logs » efface définitivement tous les logs stockés (il faut
|
||||
taper `PURGE` pour confirmer). Les tags et les paramètres sont conservés. VictoriaLogs doit
|
||||
être lancé avec `-delete.enable` (déjà présent dans `docker-compose.yml`) ; mettez
|
||||
`ALLOW_PURGE=false` pour désactiver la fonction. Toute personne qui peut ouvrir l'interface
|
||||
peut purger : activez l'[authentification](#authentification) si l'interface est accessible
|
||||
à d'autres.
|
||||
|
||||
## Authentification
|
||||
|
||||
`AUTH_MODE` choisit comment l'interface et l'API sont protégées (`/healthz` reste toujours ouvert) :
|
||||
|
||||
- **`local`** (par défaut) : authentification HTTP Basic avec `AUTH_USER` / `AUTH_PASS` ; laissez-les
|
||||
vides pour n'avoir aucune authentification (par exemple derrière un reverse proxy qui contrôle déjà).
|
||||
- **`oidc`** : connexion par un fournisseur OpenID Connect (Keycloak, Authentik, Authelia, Zitadel…),
|
||||
flux « authorization code » avec PKCE.
|
||||
|
||||
Pour utiliser OIDC :
|
||||
|
||||
1. Dans le fournisseur, créez un client **confidentiel** (avec secret) pour logstream et déclarez
|
||||
l'URL de retour `https://logs.example.org/auth/callback` (votre adresse).
|
||||
2. Dans `.env` :
|
||||
```bash
|
||||
AUTH_MODE=oidc
|
||||
OIDC_ISSUER=https://sso.example.org/realms/maison # exactement l'« issuer » du fournisseur
|
||||
OIDC_CLIENT_ID=logstream
|
||||
OIDC_CLIENT_SECRET=...
|
||||
OIDC_REDIRECT_URL=https://logs.example.org/auth/callback
|
||||
```
|
||||
3. `docker compose up -d`. Les logs affichent `oidc authentication enabled`, ou la raison pour
|
||||
laquelle le fournisseur n'a pas pu être lu (issuer incorrect, injoignable…).
|
||||
|
||||
Ouvrir l'interface renvoie vers la page de connexion du fournisseur, puis revient sur logstream.
|
||||
La session dure `OIDC_SESSION_TTL` (12 h par défaut) et survit aux redémarrages (sa clé de
|
||||
signature est dans `/data/session.key`) ; à son expiration, la page repasse par la connexion. Le
|
||||
bouton de déconnexion (en haut à droite) termine la session logstream, puis ouvre la page de
|
||||
déconnexion du fournisseur s'il en a une.
|
||||
|
||||
Tout utilisateur accepté par le fournisseur pour ce client peut se connecter : restreignez l'accès
|
||||
dans le fournisseur (Keycloak : rôles du client ou realm dédié ; Authentik : liaisons de
|
||||
l'application). Les connexions sont écrites dans les logs de logstream (`oidc: alice logged in`).
|
||||
Avec une URL de retour en `https`, les cookies ne sont envoyés qu'en HTTPS : logstream doit être
|
||||
joint à travers un reverse proxy TLS.
|
||||
|
||||
## Noms d'hôtes (DNS inverse)
|
||||
|
||||
Quand un équipement envoie son adresse IP comme nom d'hôte (ou pas de nom d'hôte du tout),
|
||||
Logstream cherche son nom DNS (enregistrement PTR) et stocke le nom dans `host` et l'IP dans
|
||||
`host_ip`. Les résultats sont mis en cache (1 heure, 10 minutes quand il n'y a pas de nom). Les
|
||||
logs stockés auparavant avec une IP sont résolus à l'affichage, et le filtre d'hôte affiche
|
||||
`name (IP)`.
|
||||
|
||||
Le conteneur utilise le DNS de Docker, qui relaie vers les résolveurs de l'hôte. Si vos noms
|
||||
locaux ne sont connus que de votre routeur ou d'un DNS local (Pi-hole, AdGuard, Unbound…),
|
||||
définissez `DNS_SERVER=192.168.1.1` (son adresse). Mettez `RDNS=off` pour désactiver les
|
||||
résolutions.
|
||||
|
||||
## Configuration
|
||||
|
||||
| Variable | Défaut | Rôle |
|
||||
|---|---|---|
|
||||
| `SYSLOG_PORT` | `514` | port syslog publié sur l'hôte |
|
||||
| `HTTP_PORT` | `8080` | port de l'interface web |
|
||||
| `RETENTION` | `30d` | durée de conservation des logs dans VictoriaLogs |
|
||||
| `AUTH_MODE` | `local` | `local` (HTTP Basic) ou `oidc`, voir [Authentification](#authentification) |
|
||||
| `AUTH_USER` / `AUTH_PASS` | vide | authentification HTTP Basic pour l'interface (mode `local`) |
|
||||
| `OIDC_ISSUER` | vide | URL de l'issuer du fournisseur OpenID Connect (mode `oidc`) |
|
||||
| `OIDC_CLIENT_ID` / `OIDC_CLIENT_SECRET` | vide | client déclaré dans le fournisseur |
|
||||
| `OIDC_REDIRECT_URL` | vide | URL de retour de logstream, ex. `https://logs.example.org/auth/callback` |
|
||||
| `OIDC_SCOPES` | `openid profile email` | scopes demandés |
|
||||
| `OIDC_SESSION_TTL` | `12h` | durée de la session |
|
||||
| `RDNS` | `on` | résoudre les hôtes donnés par leur IP en noms DNS |
|
||||
| `DNS_SERVER` | vide | serveur DNS pour les résolutions inverses (`ip` ou `ip:port`) |
|
||||
| `ALLOW_PURGE` | `true` | autoriser « Supprimer tous les logs » dans les Paramètres |
|
||||
| `EXPORT_MAX` | `100000` | nombre maximal de lignes dans un export CSV |
|
||||
| `DOCKER_LOGS` | `on` dans compose | collecter les logs des conteneurs Docker locaux |
|
||||
| `DOCKER_HOST` | `tcp://docker-proxy:2375` dans compose | adresse de l'API Docker (`unix:///var/run/docker.sock` hors compose) |
|
||||
| `DOCKER_BACKFILL` | `1h` | historique lu pour un conteneur vu pour la première fois |
|
||||
| `HOST_LOGS_GID` | `4` (adm) dans compose | groupe donné au conteneur pour lire les logs de l'hôte |
|
||||
| `HOST_LOGS_BACKFILL` | `1h` | historique lu à l'activation de la source « logs système de l'hôte » |
|
||||
| `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 |
|
||||
|
||||
## 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.
|
||||
- <http://localhost:9428/select/vmui> : l'interface de VictoriaLogs, pour essayer des requêtes
|
||||
LogsQL.
|
||||
- API :
|
||||
```bash
|
||||
curl 'localhost:8080/api/logs?q=error&range=1h&limit=5' # the response includes the generated LogsQL query
|
||||
curl localhost:8080/api/stats
|
||||
curl localhost:8080/api/tags
|
||||
```
|
||||
- Lancement hors Docker (Go 1.22+) : `VLOGS_URL=http://localhost:9428 DATA_DIR=./data SYSLOG_ADDR=:5514 go run .`
|
||||
|
||||
## Mise à jour
|
||||
|
||||
Toutes les versions d'images sont épinglées : un `docker compose pull` ou une reconstruction ne
|
||||
change jamais un composant à votre insu.
|
||||
|
||||
| Emplacement | Image | Version |
|
||||
|---|---|---|
|
||||
| `docker-compose.yml` | `victoriametrics/victoria-logs` | `v1.52.0` |
|
||||
| `docker-compose.yml` | `tecnativa/docker-socket-proxy` | `v0.5.0` |
|
||||
| `Dockerfile` (build) | `golang` | `1.27.1-alpine3.24` |
|
||||
| `Dockerfile` (exécution) | `alpine` | `3.24.2` |
|
||||
|
||||
Pour mettre à jour l'une d'elles :
|
||||
|
||||
1. Lisez les notes de version : [VictoriaLogs](https://docs.victoriametrics.com/victorialogs/changelog/),
|
||||
[docker-socket-proxy](https://github.com/Tecnativa/docker-socket-proxy/releases),
|
||||
[Go](https://go.dev/doc/devel/release), [Alpine](https://alpinelinux.org/releases/).
|
||||
VictoriaLogs conserve son format de stockage entre versions mineures ; lisez le changelog
|
||||
avant un changement de version majeure.
|
||||
2. Modifiez la version dans le fichier indiqué ci-dessus, puis lancez `docker compose up -d --build`.
|
||||
3. Vérifiez la barre du bas (reçus / stockés / perdus) et **Paramètres > Sources**. Pour revenir
|
||||
en arrière, remettez la version précédente et lancez la même commande.
|
||||
|
||||
## Organisation du code
|
||||
|
||||
| Fichier | Contenu |
|
||||
|---|---|
|
||||
| `main.go` | configuration, démarrage |
|
||||
| `auth.go` | authentification : HTTP Basic ou OpenID Connect (découverte, PKCE, contrôle de l'ID token, cookie de session) |
|
||||
| `syslog.go` | écoute UDP/TCP et analyse RFC 3164 / 5424 |
|
||||
| `store.go` | insertions par lots dans VictoriaLogs et requêtes LogsQL |
|
||||
| `query.go` | traduit les filtres de l'interface en LogsQL ; filtre du direct |
|
||||
| `histogram.go` | frise : division, alignement des intervalles, `/api/histogram` |
|
||||
| `hub.go` | envoie les nouveaux messages aux navigateurs (SSE) |
|
||||
| `rdns.go` | résolutions DNS inverses avec cache |
|
||||
| `export.go` | export CSV en flux |
|
||||
| `docker.go` | logs des conteneurs Docker (API, lecteurs, positions, détection du niveau) |
|
||||
| `syslogserver.go` | écoutes syslog ouvertes et fermées depuis Paramètres > Sources |
|
||||
| `hostlogs.go`, `journal.go` | logs système de l'hôte : lecteur du journal systemd (sans `journalctl`) et suivi de `/var/log` |
|
||||
| `tags.go` | stockage des tags de couleur |
|
||||
| `api.go` | routes HTTP `/api/*` |
|
||||
| `web/` | interface (HTML, CSS, JavaScript simple, sans étape de build), embarquée dans le binaire ; les traductions sont dans `web/app.js` (`I18N`) |
|
||||
|
||||
## Remarque
|
||||
|
||||
Avec Docker Desktop (macOS/Windows), l'IP source vue par le conteneur pour les paquets UDP est
|
||||
la passerelle Docker. Le champ `host` vient toujours de l'en-tête syslog, qui porte normalement
|
||||
le vrai nom de l'émetteur.
|
||||
@@ -1,3 +1,5 @@
|
||||
English | [Français](README.fr.md)
|
||||
|
||||
# Logstream
|
||||
|
||||
A simple syslog sink: receives logs over UDP/TCP on port 514, stores them in
|
||||
@@ -12,18 +14,20 @@ devices ──514 udp/tcp──▶ logstream (Go) ──HTTP batches──▶ Vi
|
||||
|
||||
## Architecture
|
||||
|
||||

|
||||

|
||||
|
||||
- **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.
|
||||
- **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.
|
||||
- **Timeline**: `/api/histogram` (`histogram.go`) counts the logs per interval and severity,
|
||||
aligned on the local time; the zoom bounds (`from`/`to`) also apply to the list and the export.
|
||||
- **State**: tags and source settings live in `/data` (`logstream-data` volume); the logs
|
||||
themselves in the `vlogs-data` volume.
|
||||
|
||||
The editable source of the diagram is
|
||||
[`docs/logstream-schema-logique.excalidraw`](docs/logstream-schema-logique.excalidraw)
|
||||
[`docs/architecture.excalidraw`](docs/architecture.excalidraw)
|
||||
(open it on [excalidraw.com](https://excalidraw.com)).
|
||||
|
||||
## Getting started
|
||||
@@ -143,6 +147,32 @@ colored and exported like syslog messages:
|
||||
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).
|
||||
|
||||
## Host system logs
|
||||
|
||||
Logstream can also collect the system logs of the machine hosting the stack, without
|
||||
configuring anything on the host. The source is **off by default**: turn it on in
|
||||
**Settings > Sources > Host system logs** (saved in `/data/hostlogs.json`).
|
||||
|
||||
- `docker-compose.yml` mounts `/var/log` and `/run/log/journal` read-only under `/host`.
|
||||
- When the host runs systemd, Logstream reads the **systemd journal** files directly (no
|
||||
`journalctl` needed in the image): **host** is the machine name, **app** the program
|
||||
(`SYSLOG_IDENTIFIER`), severity and facility come from the journal, and the systemd unit is
|
||||
kept in `unit`. Otherwise it follows the text files of `/var/log` (`syslog`, `messages`,
|
||||
`*.log`), parsed like syslog lines, with the file name in `log_file`.
|
||||
- These logs have the `host` source: the **Source** filter shows them alone.
|
||||
- When the source is turned on, the last hour is read first (`HOST_LOGS_BACKFILL`); the
|
||||
position reached is saved in `/data/hostlogs-state.json`, so a restart neither loses nor
|
||||
duplicates entries.
|
||||
- **Permissions**: the container runs as an unprivileged user and gets the `adm` group
|
||||
(gid 4), which can read the journal and `/var/log` on Debian and Ubuntu. On other systems,
|
||||
set `HOST_LOGS_GID` in `.env` to the gid of `systemd-journal`
|
||||
(`getent group systemd-journal | cut -d: -f3`). Settings > Sources shows a clear message when
|
||||
access is denied.
|
||||
- Limits: journal fields compressed by journald (messages longer than 512 bytes, compressed
|
||||
with zstd/lz4/xz) cannot be decoded without extra libraries; they are counted in Settings and
|
||||
shown as "(compressed journal entry)". Rotated or compressed text files (`*.1`, `*.gz`) are
|
||||
not read.
|
||||
|
||||
## CSV export
|
||||
|
||||
The **Export** button (next to the log count) downloads every stored log matching the
|
||||
@@ -189,8 +219,42 @@ remembered per browser.
|
||||
- **Data**: "Delete all logs" permanently erases every stored log (you must type
|
||||
`PURGE` to confirm). Tags and settings are kept. VictoriaLogs needs `-delete.enable`
|
||||
(already set in `docker-compose.yml`); set `ALLOW_PURGE=false` to disable the feature.
|
||||
Anyone who can open the UI can purge: set `AUTH_USER` / `AUTH_PASS` if the UI is reachable
|
||||
by others.
|
||||
Anyone who can open the UI can purge: turn on [authentication](#authentication) if the UI
|
||||
is reachable by others.
|
||||
|
||||
## Authentication
|
||||
|
||||
`AUTH_MODE` picks how the UI and the API are protected (`/healthz` always stays open):
|
||||
|
||||
- **`local`** (default): HTTP Basic authentication with `AUTH_USER` / `AUTH_PASS`; leave them
|
||||
empty to have no authentication (for instance behind a reverse proxy that already checks).
|
||||
- **`oidc`**: login through an OpenID Connect provider (Keycloak, Authentik, Authelia, Zitadel…),
|
||||
authorization code flow with PKCE.
|
||||
|
||||
To use OIDC:
|
||||
|
||||
1. In the provider, create a **confidential** client (with a secret) for logstream and register
|
||||
the redirect URL `https://logs.example.org/auth/callback` (your address).
|
||||
2. In `.env`:
|
||||
```bash
|
||||
AUTH_MODE=oidc
|
||||
OIDC_ISSUER=https://sso.example.org/realms/home # exactly the "issuer" of the provider
|
||||
OIDC_CLIENT_ID=logstream
|
||||
OIDC_CLIENT_SECRET=...
|
||||
OIDC_REDIRECT_URL=https://logs.example.org/auth/callback
|
||||
```
|
||||
3. `docker compose up -d`. The logs show `oidc authentication enabled`, or the reason the
|
||||
provider could not be read (wrong issuer, unreachable…).
|
||||
|
||||
Opening the UI sends you to the provider's login page, then back to logstream. The session
|
||||
lasts `OIDC_SESSION_TTL` (12 h by default) and survives restarts (its signing key is in
|
||||
`/data/session.key`); when it ends, the page goes through the login again. The log out button
|
||||
(top right) ends the logstream session, then opens the provider's log out page if it has one.
|
||||
|
||||
Every user the provider accepts for this client can log in: restrict access in the provider
|
||||
(Keycloak: client roles or a dedicated realm; Authentik: application bindings). Logins are written
|
||||
in the logstream logs (`oidc: alice logged in`). With an `https` redirect URL, the cookies are
|
||||
only sent over HTTPS: logstream must be reached through a TLS reverse proxy.
|
||||
|
||||
## Host names (reverse DNS)
|
||||
|
||||
@@ -202,11 +266,6 @@ on display, and the host filter shows `name (IP)`.
|
||||
The container uses Docker's DNS, which forwards to the host's resolvers. If your local names
|
||||
are only known by your router or a local DNS (Pi-hole, AdGuard, Unbound…), set
|
||||
`DNS_SERVER=192.168.1.1` (its address). Set `RDNS=off` to disable lookups.
|
||||
- **Color tags**: each tag has a keyword, a background color (the text automatically
|
||||
switches to black or white to stay readable) and options: whole word, match case,
|
||||
regular expression, active. Tags are stored in `/data/tags.json` (`logstream-data` volume).
|
||||
Default tags (pastel): `warning` (orange), `error` (red), `ok` (green). Default tags
|
||||
still using the colors of earlier versions are switched to the pastel ones automatically.
|
||||
|
||||
## Configuration
|
||||
|
||||
@@ -215,7 +274,13 @@ are only known by your router or a local DNS (Pi-hole, AdGuard, Unbound…), set
|
||||
| `SYSLOG_PORT` | `514` | syslog port published on the host |
|
||||
| `HTTP_PORT` | `8080` | web UI port |
|
||||
| `RETENTION` | `30d` | how long VictoriaLogs keeps logs |
|
||||
| `AUTH_USER` / `AUTH_PASS` | empty | HTTP Basic authentication for the UI |
|
||||
| `AUTH_MODE` | `local` | `local` (HTTP Basic) or `oidc`, see [Authentication](#authentication) |
|
||||
| `AUTH_USER` / `AUTH_PASS` | empty | HTTP Basic authentication for the UI (`local` mode) |
|
||||
| `OIDC_ISSUER` | empty | issuer URL of the OpenID Connect provider (`oidc` mode) |
|
||||
| `OIDC_CLIENT_ID` / `OIDC_CLIENT_SECRET` | empty | client registered in the provider |
|
||||
| `OIDC_REDIRECT_URL` | empty | callback URL of logstream, e.g. `https://logs.example.org/auth/callback` |
|
||||
| `OIDC_SCOPES` | `openid profile email` | requested scopes |
|
||||
| `OIDC_SESSION_TTL` | `12h` | session lifetime |
|
||||
| `RDNS` | `on` | resolve IP hosts to DNS names |
|
||||
| `DNS_SERVER` | empty | DNS server for reverse lookups (`ip` or `ip:port`) |
|
||||
| `ALLOW_PURGE` | `true` | allow "Delete all logs" in Settings |
|
||||
@@ -223,6 +288,9 @@ are only known by your router or a local DNS (Pi-hole, AdGuard, Unbound…), set
|
||||
| `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 |
|
||||
| `HOST_LOGS_GID` | `4` (adm) in compose | group given to the container to read the host logs |
|
||||
| `HOST_LOGS_BACKFILL` | `1h` | history read when the host system logs source is turned on |
|
||||
| `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 |
|
||||
|
||||
@@ -266,7 +334,8 @@ To update one of them:
|
||||
|
||||
| File | Contents |
|
||||
|---|---|
|
||||
| `main.go` | configuration, startup, authentication |
|
||||
| `main.go` | configuration, startup |
|
||||
| `auth.go` | authentication: HTTP Basic or OpenID Connect (discovery, PKCE, ID token checks, session cookie) |
|
||||
| `syslog.go` | UDP/TCP listeners and RFC 3164 / 5424 parsing |
|
||||
| `store.go` | batched inserts into VictoriaLogs and LogsQL queries |
|
||||
| `query.go` | turns UI filters into LogsQL; live-view filter |
|
||||
@@ -276,6 +345,7 @@ To update one of them:
|
||||
| `export.go` | streamed CSV export |
|
||||
| `docker.go` | Docker container logs (API, followers, positions, level detection) |
|
||||
| `syslogserver.go` | syslog listeners opened and closed from Settings > Sources |
|
||||
| `hostlogs.go`, `journal.go` | host system logs: systemd journal reader (no `journalctl`) and `/var/log` follower |
|
||||
| `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`) |
|
||||
|
||||
@@ -21,6 +21,7 @@ type API struct {
|
||||
exportMax int
|
||||
docker *DockerManager // nil when DOCKER_LOGS is off
|
||||
syslog *SyslogServer
|
||||
host *HostLogs
|
||||
}
|
||||
|
||||
func (a *API) Routes(mux *http.ServeMux) {
|
||||
@@ -42,6 +43,28 @@ func (a *API) Routes(mux *http.ServeMux) {
|
||||
mux.HandleFunc("PUT /api/syslog", a.syslogConfigure)
|
||||
mux.HandleFunc("GET /api/docker", a.dockerStatus)
|
||||
mux.HandleFunc("PUT /api/docker", a.dockerConfigure)
|
||||
mux.HandleFunc("GET /api/hostlogs", a.hostLogsStatus)
|
||||
mux.HandleFunc("PUT /api/hostlogs", a.hostLogsConfigure)
|
||||
}
|
||||
|
||||
// GET /api/hostlogs: state of the host system logs source (Settings > Sources).
|
||||
func (a *API) hostLogsStatus(w http.ResponseWriter, r *http.Request) {
|
||||
writeJSON(w, http.StatusOK, a.host.Status())
|
||||
}
|
||||
|
||||
// PUT /api/hostlogs {"enabled": bool}
|
||||
func (a *API) hostLogsConfigure(w http.ResponseWriter, r *http.Request) {
|
||||
var cfg hostLogsConfig
|
||||
if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 4096)).Decode(&cfg); err != nil {
|
||||
writeErr(w, http.StatusBadRequest, err)
|
||||
return
|
||||
}
|
||||
if err := a.host.Configure(cfg); err != nil {
|
||||
writeErr(w, http.StatusInternalServerError, err)
|
||||
return
|
||||
}
|
||||
log.Printf("host system logs changed from %s: %+v", r.RemoteAddr, cfg)
|
||||
writeJSON(w, http.StatusOK, a.host.Status())
|
||||
}
|
||||
|
||||
// GET /api/syslog: syslog reception state (Settings > Sources).
|
||||
|
||||
@@ -0,0 +1,631 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"crypto"
|
||||
"crypto/ecdsa"
|
||||
"crypto/elliptic"
|
||||
"crypto/hmac"
|
||||
"crypto/rand"
|
||||
"crypto/rsa"
|
||||
"crypto/sha256"
|
||||
_ "crypto/sha512" // SHA-384/512 for RS384, ES384, RS512…
|
||||
"crypto/subtle"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"math/big"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Web UI authentication. AUTH_MODE=local (default) keeps the optional HTTP Basic
|
||||
// authentication (AUTH_USER / AUTH_PASS); AUTH_MODE=oidc delegates the login to an
|
||||
// OpenID Connect provider (Keycloak, Authentik, Authelia…) with the authorization code
|
||||
// flow and PKCE. Only the standard library is used.
|
||||
|
||||
const (
|
||||
sessionCookie = "logstream_session"
|
||||
loginCookie = "logstream_login_" // + state: one cookie per login in progress
|
||||
loginTTL = 10 * time.Minute
|
||||
clockSkew = time.Minute
|
||||
)
|
||||
|
||||
type authConfig struct {
|
||||
mode string
|
||||
user, pass string // local mode
|
||||
issuer string
|
||||
clientID string
|
||||
clientSecret string
|
||||
redirectURL string
|
||||
scopes string
|
||||
sessionTTL time.Duration
|
||||
dataDir string
|
||||
}
|
||||
|
||||
// newAuth returns the middleware that protects the UI and the API (except /healthz).
|
||||
func newAuth(c authConfig, next http.Handler) (http.Handler, error) {
|
||||
switch strings.ToLower(c.mode) {
|
||||
case "", "local":
|
||||
return basicAuth(c.user, c.pass, next), nil
|
||||
case "oidc":
|
||||
o, err := newOIDC(c)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
o.next = next
|
||||
log.Printf("oidc authentication enabled (issuer %s)", c.issuer)
|
||||
return o, nil
|
||||
}
|
||||
return nil, fmt.Errorf("AUTH_MODE=%q: expected local or oidc", c.mode)
|
||||
}
|
||||
|
||||
// basicAuth protects the UI when AUTH_USER is set (except /healthz).
|
||||
func basicAuth(user, pass string, next http.Handler) http.Handler {
|
||||
if user == "" {
|
||||
return next
|
||||
}
|
||||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Path == "/healthz" {
|
||||
next.ServeHTTP(w, r)
|
||||
return
|
||||
}
|
||||
u, p, ok := r.BasicAuth()
|
||||
if !ok ||
|
||||
subtle.ConstantTimeCompare([]byte(u), []byte(user)) != 1 ||
|
||||
subtle.ConstantTimeCompare([]byte(p), []byte(pass)) != 1 {
|
||||
w.Header().Set("WWW-Authenticate", `Basic realm="logstream"`)
|
||||
http.Error(w, "authentication required", http.StatusUnauthorized)
|
||||
return
|
||||
}
|
||||
next.ServeHTTP(w, r)
|
||||
})
|
||||
}
|
||||
|
||||
type oidcMeta struct {
|
||||
Issuer string `json:"issuer"`
|
||||
AuthEndpoint string `json:"authorization_endpoint"`
|
||||
TokenEndpoint string `json:"token_endpoint"`
|
||||
JWKSURI string `json:"jwks_uri"`
|
||||
EndSession string `json:"end_session_endpoint"`
|
||||
TokenAuthMethods []string `json:"token_endpoint_auth_methods_supported"`
|
||||
}
|
||||
|
||||
type OIDC struct {
|
||||
cfg authConfig
|
||||
callback string // path of OIDC_REDIRECT_URL
|
||||
secure bool // cookies only sent over HTTPS
|
||||
key []byte // signs the session and login cookies
|
||||
client *http.Client
|
||||
next http.Handler
|
||||
|
||||
mu sync.Mutex
|
||||
meta *oidcMeta
|
||||
keys map[string]crypto.PublicKey
|
||||
keysAt time.Time
|
||||
}
|
||||
|
||||
func newOIDC(c authConfig) (*OIDC, error) {
|
||||
var missing []string
|
||||
for _, v := range [][2]string{
|
||||
{"OIDC_ISSUER", c.issuer}, {"OIDC_CLIENT_ID", c.clientID},
|
||||
{"OIDC_CLIENT_SECRET", c.clientSecret}, {"OIDC_REDIRECT_URL", c.redirectURL},
|
||||
} {
|
||||
if v[1] == "" {
|
||||
missing = append(missing, v[0])
|
||||
}
|
||||
}
|
||||
if len(missing) > 0 {
|
||||
return nil, fmt.Errorf("AUTH_MODE=oidc: missing %s", strings.Join(missing, ", "))
|
||||
}
|
||||
ru, err := url.Parse(c.redirectURL)
|
||||
if err != nil || ru.Host == "" || ru.Path == "" || ru.Path == "/" {
|
||||
return nil, fmt.Errorf("OIDC_REDIRECT_URL=%q: expected a full URL such as https://logs.example.org/auth/callback", c.redirectURL)
|
||||
}
|
||||
if c.scopes == "" {
|
||||
c.scopes = "openid profile email"
|
||||
}
|
||||
if !strings.Contains(" "+c.scopes+" ", " openid ") {
|
||||
c.scopes = "openid " + c.scopes
|
||||
}
|
||||
if c.sessionTTL <= 0 {
|
||||
c.sessionTTL = 12 * time.Hour
|
||||
}
|
||||
return &OIDC{
|
||||
cfg: c,
|
||||
callback: ru.Path,
|
||||
secure: ru.Scheme == "https",
|
||||
key: sessionKey(c.dataDir),
|
||||
client: &http.Client{Timeout: 10 * time.Second},
|
||||
}, nil
|
||||
}
|
||||
|
||||
// checkProvider reads the provider configuration at startup so a mistake shows in the logs.
|
||||
func (o *OIDC) checkProvider() {
|
||||
if _, err := o.discover(); err != nil {
|
||||
log.Printf("oidc: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// sessionKey is kept in DATA_DIR so sessions survive a restart.
|
||||
func sessionKey(dir string) []byte {
|
||||
path := filepath.Join(dir, "session.key")
|
||||
if k, err := os.ReadFile(path); err == nil && len(k) >= 32 {
|
||||
return k
|
||||
}
|
||||
k := make([]byte, 32)
|
||||
if _, err := rand.Read(k); err != nil {
|
||||
log.Fatalf("session key: %v", err)
|
||||
}
|
||||
if err := os.WriteFile(path, k, 0o600); err != nil {
|
||||
log.Printf("oidc: cannot save %s (%v): sessions end when logstream restarts", path, err)
|
||||
}
|
||||
return k
|
||||
}
|
||||
|
||||
type session struct {
|
||||
User string `json:"u"`
|
||||
Exp int64 `json:"e"`
|
||||
}
|
||||
|
||||
type loginState struct {
|
||||
Nonce string `json:"n"`
|
||||
Verifier string `json:"v"`
|
||||
Return string `json:"r"`
|
||||
Exp int64 `json:"e"`
|
||||
}
|
||||
|
||||
func (o *OIDC) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
switch r.URL.Path {
|
||||
case "/healthz":
|
||||
o.next.ServeHTTP(w, r)
|
||||
return
|
||||
case o.callback:
|
||||
o.handleCallback(w, r)
|
||||
return
|
||||
case "/auth/logout":
|
||||
o.handleLogout(w, r)
|
||||
return
|
||||
}
|
||||
var s session
|
||||
if c, err := r.Cookie(sessionCookie); err == nil && o.verifyCookie(c.Value, &s) && time.Now().Unix() < s.Exp {
|
||||
if r.URL.Path == "/auth/me" {
|
||||
writeJSON(w, http.StatusOK, map[string]string{"mode": "oidc", "user": s.User})
|
||||
return
|
||||
}
|
||||
o.next.ServeHTTP(w, r)
|
||||
return
|
||||
}
|
||||
// Not logged in: pages go to the provider, API calls get a 401 that the UI turns
|
||||
// into a reload (and so into a new login).
|
||||
if r.Method == http.MethodGet && !strings.HasPrefix(r.URL.Path, "/api/") && r.URL.Path != "/auth/me" {
|
||||
o.startLogin(w, r)
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(http.StatusUnauthorized)
|
||||
_, _ = w.Write([]byte(`{"error":"authentication required","code":"auth"}` + "\n"))
|
||||
}
|
||||
|
||||
func (o *OIDC) startLogin(w http.ResponseWriter, r *http.Request) {
|
||||
meta, err := o.discover()
|
||||
if err != nil {
|
||||
log.Printf("oidc: %v", err)
|
||||
http.Error(w, "identity provider unreachable, try again later", http.StatusBadGateway)
|
||||
return
|
||||
}
|
||||
state, nonce, verifier := randomString(), randomString(), randomString()+randomString()
|
||||
ret := r.URL.RequestURI()
|
||||
if !strings.HasPrefix(ret, "/") || strings.HasPrefix(ret, "//") {
|
||||
ret = "/"
|
||||
}
|
||||
http.SetCookie(w, &http.Cookie{
|
||||
Name: loginCookie + state,
|
||||
Value: o.signCookie(loginState{Nonce: nonce, Verifier: verifier, Return: ret, Exp: time.Now().Add(loginTTL).Unix()}),
|
||||
Path: "/",
|
||||
MaxAge: int(loginTTL.Seconds()),
|
||||
HttpOnly: true,
|
||||
Secure: o.secure,
|
||||
SameSite: http.SameSiteLaxMode, // sent back on the redirect from the provider
|
||||
})
|
||||
challenge := sha256.Sum256([]byte(verifier))
|
||||
q := url.Values{
|
||||
"response_type": {"code"},
|
||||
"client_id": {o.cfg.clientID},
|
||||
"redirect_uri": {o.cfg.redirectURL},
|
||||
"scope": {o.cfg.scopes},
|
||||
"state": {state},
|
||||
"nonce": {nonce},
|
||||
"code_challenge": {base64.RawURLEncoding.EncodeToString(challenge[:])},
|
||||
"code_challenge_method": {"S256"},
|
||||
}
|
||||
http.Redirect(w, r, addQuery(meta.AuthEndpoint, q), http.StatusFound)
|
||||
}
|
||||
|
||||
func (o *OIDC) handleCallback(w http.ResponseWriter, r *http.Request) {
|
||||
q := r.URL.Query()
|
||||
if e := q.Get("error"); e != "" {
|
||||
log.Printf("oidc: login refused by the provider: %s %s", e, q.Get("error_description"))
|
||||
http.Error(w, "login refused by the identity provider: "+e, http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
state := q.Get("state")
|
||||
var ls loginState
|
||||
c, err := r.Cookie(loginCookie + state)
|
||||
if state == "" || err != nil || !o.verifyCookie(c.Value, &ls) || time.Now().Unix() > ls.Exp {
|
||||
http.Error(w, "login expired or started in another browser: open logstream again", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
http.SetCookie(w, &http.Cookie{Name: loginCookie + state, Path: "/", MaxAge: -1, HttpOnly: true, Secure: o.secure})
|
||||
|
||||
user, err := o.exchange(r, q.Get("code"), ls)
|
||||
if err != nil {
|
||||
log.Printf("oidc: login failed: %v", err)
|
||||
http.Error(w, "login failed, see the logstream logs", http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
log.Printf("oidc: %s logged in", user)
|
||||
http.SetCookie(w, &http.Cookie{
|
||||
Name: sessionCookie,
|
||||
Value: o.signCookie(session{User: user, Exp: time.Now().Add(o.cfg.sessionTTL).Unix()}),
|
||||
Path: "/",
|
||||
MaxAge: int(o.cfg.sessionTTL.Seconds()),
|
||||
HttpOnly: true,
|
||||
Secure: o.secure,
|
||||
SameSite: http.SameSiteLaxMode,
|
||||
})
|
||||
http.Redirect(w, r, ls.Return, http.StatusFound)
|
||||
}
|
||||
|
||||
// The session ends here; the provider's own session ends on its logout page if it has one.
|
||||
func (o *OIDC) handleLogout(w http.ResponseWriter, r *http.Request) {
|
||||
http.SetCookie(w, &http.Cookie{Name: sessionCookie, Path: "/", MaxAge: -1, HttpOnly: true, Secure: o.secure})
|
||||
if meta, err := o.discover(); err == nil && meta.EndSession != "" {
|
||||
http.Redirect(w, r, addQuery(meta.EndSession, url.Values{"client_id": {o.cfg.clientID}}), http.StatusFound)
|
||||
return
|
||||
}
|
||||
http.Redirect(w, r, "/", http.StatusFound)
|
||||
}
|
||||
|
||||
// exchange trades the code for tokens and returns the user name from the verified ID token.
|
||||
func (o *OIDC) exchange(r *http.Request, code string, ls loginState) (string, error) {
|
||||
if code == "" {
|
||||
return "", errors.New("no code in the callback")
|
||||
}
|
||||
meta, err := o.discover()
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
form := url.Values{
|
||||
"grant_type": {"authorization_code"},
|
||||
"code": {code},
|
||||
"redirect_uri": {o.cfg.redirectURL},
|
||||
"code_verifier": {ls.Verifier},
|
||||
}
|
||||
// client_secret_basic is the default; some providers only accept client_secret_post.
|
||||
post := len(meta.TokenAuthMethods) > 0 && !contains(meta.TokenAuthMethods, "client_secret_basic") && contains(meta.TokenAuthMethods, "client_secret_post")
|
||||
if post {
|
||||
form.Set("client_id", o.cfg.clientID)
|
||||
form.Set("client_secret", o.cfg.clientSecret)
|
||||
}
|
||||
req, err := http.NewRequestWithContext(r.Context(), http.MethodPost, meta.TokenEndpoint, strings.NewReader(form.Encode()))
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
|
||||
req.Header.Set("Accept", "application/json")
|
||||
if !post {
|
||||
req.SetBasicAuth(url.QueryEscape(o.cfg.clientID), url.QueryEscape(o.cfg.clientSecret))
|
||||
}
|
||||
res, err := o.client.Do(req)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("token endpoint: %w", err)
|
||||
}
|
||||
defer res.Body.Close()
|
||||
body, _ := io.ReadAll(io.LimitReader(res.Body, 1<<20))
|
||||
if res.StatusCode != http.StatusOK {
|
||||
return "", fmt.Errorf("token endpoint: %s: %s", res.Status, bytes.TrimSpace(body))
|
||||
}
|
||||
var tok struct {
|
||||
IDToken string `json:"id_token"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &tok); err != nil || tok.IDToken == "" {
|
||||
return "", errors.New("token endpoint: no id_token in the response")
|
||||
}
|
||||
claims, err := o.verifyIDToken(tok.IDToken, ls.Nonce)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
for _, k := range []string{"preferred_username", "email", "name", "sub"} {
|
||||
if v, _ := claims[k].(string); v != "" {
|
||||
return v, nil
|
||||
}
|
||||
}
|
||||
return "", errors.New("id_token: no sub")
|
||||
}
|
||||
|
||||
// verifyIDToken checks the signature (keys from jwks_uri) and the claims of an ID token.
|
||||
func (o *OIDC) verifyIDToken(raw, nonce string) (map[string]any, error) {
|
||||
parts := strings.Split(raw, ".")
|
||||
if len(parts) != 3 {
|
||||
return nil, errors.New("id_token: not a JWT")
|
||||
}
|
||||
var hdr struct {
|
||||
Alg string `json:"alg"`
|
||||
Kid string `json:"kid"`
|
||||
}
|
||||
if err := decodeSegment(parts[0], &hdr); err != nil {
|
||||
return nil, fmt.Errorf("id_token header: %w", err)
|
||||
}
|
||||
sig, err := base64.RawURLEncoding.DecodeString(parts[2])
|
||||
if err != nil {
|
||||
return nil, errors.New("id_token: bad signature encoding")
|
||||
}
|
||||
key, err := o.keyFor(hdr.Kid)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := verifySignature(hdr.Alg, key, []byte(parts[0]+"."+parts[1]), sig); err != nil {
|
||||
return nil, fmt.Errorf("id_token: %w", err)
|
||||
}
|
||||
var claims map[string]any
|
||||
if err := decodeSegment(parts[1], &claims); err != nil {
|
||||
return nil, fmt.Errorf("id_token claims: %w", err)
|
||||
}
|
||||
if iss, _ := claims["iss"].(string); iss != o.cfg.issuer {
|
||||
return nil, fmt.Errorf("id_token: issuer %q, expected %q", iss, o.cfg.issuer)
|
||||
}
|
||||
var aud []string
|
||||
switch v := claims["aud"].(type) {
|
||||
case string:
|
||||
aud = []string{v}
|
||||
case []any:
|
||||
for _, a := range v {
|
||||
if s, ok := a.(string); ok {
|
||||
aud = append(aud, s)
|
||||
}
|
||||
}
|
||||
}
|
||||
if !contains(aud, o.cfg.clientID) {
|
||||
return nil, fmt.Errorf("id_token: audience %v does not include %q", aud, o.cfg.clientID)
|
||||
}
|
||||
if azp, ok := claims["azp"].(string); ok && len(aud) > 1 && azp != o.cfg.clientID {
|
||||
return nil, fmt.Errorf("id_token: azp %q", azp)
|
||||
}
|
||||
now := time.Now()
|
||||
exp, _ := claims["exp"].(float64)
|
||||
if exp == 0 || now.After(time.Unix(int64(exp), 0).Add(clockSkew)) {
|
||||
return nil, errors.New("id_token: expired (check the clocks)")
|
||||
}
|
||||
if iat, ok := claims["iat"].(float64); ok && time.Unix(int64(iat), 0).After(now.Add(clockSkew)) {
|
||||
return nil, errors.New("id_token: issued in the future (check the clocks)")
|
||||
}
|
||||
if n, _ := claims["nonce"].(string); subtle.ConstantTimeCompare([]byte(n), []byte(nonce)) != 1 {
|
||||
return nil, errors.New("id_token: wrong nonce")
|
||||
}
|
||||
return claims, nil
|
||||
}
|
||||
|
||||
func verifySignature(alg string, key crypto.PublicKey, signed, sig []byte) error {
|
||||
if len(alg) != 5 {
|
||||
return fmt.Errorf("unsupported algorithm %q", alg)
|
||||
}
|
||||
var h crypto.Hash
|
||||
switch alg[2:] {
|
||||
case "256":
|
||||
h = crypto.SHA256
|
||||
case "384":
|
||||
h = crypto.SHA384
|
||||
case "512":
|
||||
h = crypto.SHA512
|
||||
}
|
||||
if h == 0 {
|
||||
return fmt.Errorf("unsupported algorithm %q", alg)
|
||||
}
|
||||
hh := h.New()
|
||||
hh.Write(signed)
|
||||
digest := hh.Sum(nil)
|
||||
switch k := key.(type) {
|
||||
case *rsa.PublicKey:
|
||||
switch alg[:2] {
|
||||
case "RS":
|
||||
return rsa.VerifyPKCS1v15(k, h, digest, sig)
|
||||
case "PS":
|
||||
return rsa.VerifyPSS(k, h, digest, sig, &rsa.PSSOptions{SaltLength: rsa.PSSSaltLengthEqualsHash})
|
||||
}
|
||||
case *ecdsa.PublicKey:
|
||||
size := (k.Curve.Params().BitSize + 7) / 8
|
||||
if alg[:2] != "ES" || len(sig) != 2*size {
|
||||
break
|
||||
}
|
||||
r, s := new(big.Int).SetBytes(sig[:size]), new(big.Int).SetBytes(sig[size:])
|
||||
if ecdsa.Verify(k, digest, r, s) {
|
||||
return nil
|
||||
}
|
||||
return errors.New("bad signature")
|
||||
}
|
||||
return fmt.Errorf("algorithm %q does not match the key", alg)
|
||||
}
|
||||
|
||||
// discover reads the provider configuration once (and again after a failure).
|
||||
func (o *OIDC) discover() (*oidcMeta, error) {
|
||||
o.mu.Lock()
|
||||
defer o.mu.Unlock()
|
||||
if o.meta != nil {
|
||||
return o.meta, nil
|
||||
}
|
||||
u := strings.TrimSuffix(o.cfg.issuer, "/") + "/.well-known/openid-configuration"
|
||||
var m oidcMeta
|
||||
if err := o.getJSON(u, &m); err != nil {
|
||||
return nil, fmt.Errorf("discovery: %w", err)
|
||||
}
|
||||
if m.Issuer != o.cfg.issuer {
|
||||
return nil, fmt.Errorf("discovery: the provider says its issuer is %q, set OIDC_ISSUER to that exact value", m.Issuer)
|
||||
}
|
||||
if m.AuthEndpoint == "" || m.TokenEndpoint == "" || m.JWKSURI == "" {
|
||||
return nil, errors.New("discovery: incomplete provider configuration")
|
||||
}
|
||||
o.meta = &m
|
||||
return o.meta, nil
|
||||
}
|
||||
|
||||
// keyFor returns the signing key kid; the key set is reloaded when the provider rotates its keys.
|
||||
func (o *OIDC) keyFor(kid string) (crypto.PublicKey, error) {
|
||||
meta, err := o.discover()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
o.mu.Lock()
|
||||
defer o.mu.Unlock()
|
||||
pick := func() crypto.PublicKey {
|
||||
if k, ok := o.keys[kid]; ok {
|
||||
return k
|
||||
}
|
||||
if kid == "" && len(o.keys) == 1 {
|
||||
for _, k := range o.keys {
|
||||
return k
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if k := pick(); k != nil {
|
||||
return k, nil
|
||||
}
|
||||
if time.Since(o.keysAt) < 10*time.Second {
|
||||
return nil, fmt.Errorf("id_token: unknown key %q", kid)
|
||||
}
|
||||
var set struct {
|
||||
Keys []struct {
|
||||
Kty string `json:"kty"`
|
||||
Kid string `json:"kid"`
|
||||
Use string `json:"use"`
|
||||
N string `json:"n"`
|
||||
E string `json:"e"`
|
||||
Crv string `json:"crv"`
|
||||
X string `json:"x"`
|
||||
Y string `json:"y"`
|
||||
} `json:"keys"`
|
||||
}
|
||||
if err := o.getJSON(meta.JWKSURI, &set); err != nil {
|
||||
return nil, fmt.Errorf("jwks: %w", err)
|
||||
}
|
||||
keys := map[string]crypto.PublicKey{}
|
||||
for _, k := range set.Keys {
|
||||
if k.Use != "" && k.Use != "sig" {
|
||||
continue
|
||||
}
|
||||
switch k.Kty {
|
||||
case "RSA":
|
||||
n, e := decodeBig(k.N), decodeBig(k.E)
|
||||
if n != nil && e != nil && e.IsInt64() {
|
||||
keys[k.Kid] = &rsa.PublicKey{N: n, E: int(e.Int64())}
|
||||
}
|
||||
case "EC":
|
||||
var c elliptic.Curve
|
||||
switch k.Crv {
|
||||
case "P-256":
|
||||
c = elliptic.P256()
|
||||
case "P-384":
|
||||
c = elliptic.P384()
|
||||
case "P-521":
|
||||
c = elliptic.P521()
|
||||
}
|
||||
x, y := decodeBig(k.X), decodeBig(k.Y)
|
||||
if c != nil && x != nil && y != nil && c.IsOnCurve(x, y) {
|
||||
keys[k.Kid] = &ecdsa.PublicKey{Curve: c, X: x, Y: y}
|
||||
}
|
||||
}
|
||||
}
|
||||
o.keys, o.keysAt = keys, time.Now()
|
||||
if k := pick(); k != nil {
|
||||
return k, nil
|
||||
}
|
||||
return nil, fmt.Errorf("id_token: unknown key %q", kid)
|
||||
}
|
||||
|
||||
func (o *OIDC) getJSON(u string, v any) error {
|
||||
res, err := o.client.Get(u)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer res.Body.Close()
|
||||
if res.StatusCode != http.StatusOK {
|
||||
return fmt.Errorf("%s: %s", u, res.Status)
|
||||
}
|
||||
return json.NewDecoder(io.LimitReader(res.Body, 1<<20)).Decode(v)
|
||||
}
|
||||
|
||||
// Cookies are base64url(JSON) + "." + base64url(HMAC-SHA256).
|
||||
func (o *OIDC) signCookie(v any) string {
|
||||
b, _ := json.Marshal(v)
|
||||
p := base64.RawURLEncoding.EncodeToString(b)
|
||||
m := hmac.New(sha256.New, o.key)
|
||||
m.Write([]byte(p))
|
||||
return p + "." + base64.RawURLEncoding.EncodeToString(m.Sum(nil))
|
||||
}
|
||||
|
||||
func (o *OIDC) verifyCookie(s string, v any) bool {
|
||||
p, sig, ok := strings.Cut(s, ".")
|
||||
if !ok {
|
||||
return false
|
||||
}
|
||||
got, err := base64.RawURLEncoding.DecodeString(sig)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
m := hmac.New(sha256.New, o.key)
|
||||
m.Write([]byte(p))
|
||||
if !hmac.Equal(got, m.Sum(nil)) {
|
||||
return false
|
||||
}
|
||||
return decodeSegment(p, v) == nil
|
||||
}
|
||||
|
||||
func decodeSegment(s string, v any) error {
|
||||
b, err := base64.RawURLEncoding.DecodeString(s)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return json.Unmarshal(b, v)
|
||||
}
|
||||
|
||||
func decodeBig(s string) *big.Int {
|
||||
b, err := base64.RawURLEncoding.DecodeString(s)
|
||||
if err != nil || len(b) == 0 {
|
||||
return nil
|
||||
}
|
||||
return new(big.Int).SetBytes(b)
|
||||
}
|
||||
|
||||
func randomString() string {
|
||||
b := make([]byte, 16)
|
||||
if _, err := rand.Read(b); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
return base64.RawURLEncoding.EncodeToString(b)
|
||||
}
|
||||
|
||||
func addQuery(endpoint string, q url.Values) string {
|
||||
sep := "?"
|
||||
if strings.Contains(endpoint, "?") {
|
||||
sep = "&"
|
||||
}
|
||||
return endpoint + sep + q.Encode()
|
||||
}
|
||||
|
||||
func contains(list []string, s string) bool {
|
||||
for _, v := range list {
|
||||
if v == s {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
+230
@@ -0,0 +1,230 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"crypto"
|
||||
"crypto/ecdsa"
|
||||
"crypto/elliptic"
|
||||
"crypto/rand"
|
||||
"crypto/rsa"
|
||||
"crypto/sha256"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"math/big"
|
||||
"net/http"
|
||||
"net/http/cookiejar"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// hosts routes requests to in-memory handlers (no listening socket needed).
|
||||
type hosts map[string]http.Handler
|
||||
|
||||
func (h hosts) RoundTrip(r *http.Request) (*http.Response, error) {
|
||||
rec := httptest.NewRecorder()
|
||||
h[r.URL.Host].ServeHTTP(rec, r)
|
||||
res := rec.Result()
|
||||
res.Request = r
|
||||
return res, nil
|
||||
}
|
||||
|
||||
// fakeIdP is a minimal OpenID provider: it logs in "alice" without asking.
|
||||
type fakeIdP struct {
|
||||
mux *http.ServeMux
|
||||
rsaKey *rsa.PrivateKey
|
||||
ecKey *ecdsa.PrivateKey
|
||||
useEC bool
|
||||
codes map[string]url.Values // code -> authorize request
|
||||
claims func(map[string]any) // last-minute changes to the ID token
|
||||
tokenErr bool
|
||||
}
|
||||
|
||||
func newFakeIdP(t *testing.T) *fakeIdP {
|
||||
rk, _ := rsa.GenerateKey(rand.Reader, 2048)
|
||||
ek, _ := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
|
||||
p := &fakeIdP{rsaKey: rk, ecKey: ek, codes: map[string]url.Values{}}
|
||||
mux := http.NewServeMux()
|
||||
p.mux = mux
|
||||
iss := "http://idp.test/realm"
|
||||
mux.HandleFunc("/realm/.well-known/openid-configuration", func(w http.ResponseWriter, r *http.Request) {
|
||||
_ = json.NewEncoder(w).Encode(map[string]any{
|
||||
"issuer": iss, "authorization_endpoint": iss + "/auth", "token_endpoint": iss + "/token",
|
||||
"jwks_uri": iss + "/jwks", "end_session_endpoint": iss + "/logout",
|
||||
})
|
||||
})
|
||||
mux.HandleFunc("/realm/jwks", func(w http.ResponseWriter, r *http.Request) {
|
||||
b := func(i *big.Int) string { return base64.RawURLEncoding.EncodeToString(i.Bytes()) }
|
||||
_ = json.NewEncoder(w).Encode(map[string]any{"keys": []any{
|
||||
map[string]string{"kty": "RSA", "kid": "r1", "use": "sig", "n": b(rk.N), "e": "AQAB"},
|
||||
map[string]string{"kty": "EC", "kid": "e1", "crv": "P-256", "x": b(ek.X), "y": b(ek.Y)},
|
||||
}})
|
||||
})
|
||||
mux.HandleFunc("/realm/auth", func(w http.ResponseWriter, r *http.Request) {
|
||||
q := r.URL.Query()
|
||||
code := randomString()
|
||||
p.codes[code] = q
|
||||
http.Redirect(w, r, q.Get("redirect_uri")+"?code="+code+"&state="+q.Get("state"), http.StatusFound)
|
||||
})
|
||||
mux.HandleFunc("/realm/token", func(w http.ResponseWriter, r *http.Request) {
|
||||
_ = r.ParseForm()
|
||||
id, secret, _ := r.BasicAuth()
|
||||
authz, ok := p.codes[r.Form.Get("code")]
|
||||
sum := sha256.Sum256([]byte(r.Form.Get("code_verifier")))
|
||||
if p.tokenErr || !ok || id != "logstream" || secret != "s3cret" ||
|
||||
base64.RawURLEncoding.EncodeToString(sum[:]) != authz.Get("code_challenge") ||
|
||||
r.Form.Get("redirect_uri") != authz.Get("redirect_uri") {
|
||||
http.Error(w, `{"error":"invalid_grant"}`, http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
delete(p.codes, r.Form.Get("code"))
|
||||
c := map[string]any{
|
||||
"iss": iss, "aud": "logstream", "sub": "123", "preferred_username": "alice",
|
||||
"exp": time.Now().Add(5 * time.Minute).Unix(), "iat": time.Now().Unix(), "nonce": authz.Get("nonce"),
|
||||
}
|
||||
if p.claims != nil {
|
||||
p.claims(c)
|
||||
}
|
||||
_ = json.NewEncoder(w).Encode(map[string]string{"access_token": "x", "id_token": p.sign(c)})
|
||||
})
|
||||
return p
|
||||
}
|
||||
|
||||
func (p *fakeIdP) sign(claims map[string]any) string {
|
||||
alg, kid := "RS256", "r1"
|
||||
if p.useEC {
|
||||
alg, kid = "ES256", "e1"
|
||||
}
|
||||
h, _ := json.Marshal(map[string]string{"alg": alg, "kid": kid, "typ": "JWT"})
|
||||
c, _ := json.Marshal(claims)
|
||||
in := base64.RawURLEncoding.EncodeToString(h) + "." + base64.RawURLEncoding.EncodeToString(c)
|
||||
d := sha256.Sum256([]byte(in))
|
||||
var sig []byte
|
||||
if p.useEC {
|
||||
r, s, _ := ecdsa.Sign(rand.Reader, p.ecKey, d[:])
|
||||
sig = make([]byte, 64)
|
||||
r.FillBytes(sig[:32])
|
||||
s.FillBytes(sig[32:])
|
||||
} else {
|
||||
sig, _ = rsa.SignPKCS1v15(rand.Reader, p.rsaKey, crypto.SHA256, d[:])
|
||||
}
|
||||
return in + "." + base64.RawURLEncoding.EncodeToString(sig)
|
||||
}
|
||||
|
||||
// newOIDCApp puts logstream's auth in front of a handler that echoes "app" and returns
|
||||
// a browser (client with cookies) that reaches both the app and the provider.
|
||||
func newOIDCApp(t *testing.T, idp *fakeIdP) (string, *http.Client) {
|
||||
h, err := newAuth(authConfig{
|
||||
mode: "oidc", issuer: "http://idp.test/realm", clientID: "logstream", clientSecret: "s3cret",
|
||||
redirectURL: "http://app.test/auth/callback", dataDir: t.TempDir(),
|
||||
}, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { _, _ = w.Write([]byte("app " + r.URL.Path)) }))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
net := hosts{"app.test": h, "idp.test": idp.mux}
|
||||
h.(*OIDC).client.Transport = net
|
||||
jar, _ := cookiejar.New(nil)
|
||||
return "http://app.test", &http.Client{Jar: jar, Transport: net}
|
||||
}
|
||||
|
||||
func get(t *testing.T, c *http.Client, u string) (int, string) {
|
||||
t.Helper()
|
||||
res, err := c.Get(u)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer res.Body.Close()
|
||||
var b strings.Builder
|
||||
buf := make([]byte, 4096)
|
||||
for {
|
||||
n, err := res.Body.Read(buf)
|
||||
b.Write(buf[:n])
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
}
|
||||
return res.StatusCode, b.String()
|
||||
}
|
||||
|
||||
func TestOIDCLoginFlow(t *testing.T) {
|
||||
for _, ec := range []bool{false, true} {
|
||||
idp := newFakeIdP(t)
|
||||
idp.useEC = ec
|
||||
app, c := newOIDCApp(t, idp)
|
||||
|
||||
if code, _ := get(t, c, app+"/api/logs"); code != http.StatusUnauthorized {
|
||||
t.Fatalf("api without session: %d", code)
|
||||
}
|
||||
if code, body := get(t, c, app+"/healthz"); code != 200 || body != "app /healthz" {
|
||||
t.Fatalf("healthz: %d %q", code, body)
|
||||
}
|
||||
// A page goes through the provider and comes back to the page asked for.
|
||||
if code, body := get(t, c, app+"/index.html?x=1"); code != 200 || body != "app /index.html" {
|
||||
t.Fatalf("login (ec=%v): %d %q", ec, code, body)
|
||||
}
|
||||
if code, body := get(t, c, app+"/api/logs"); code != 200 || body != "app /api/logs" {
|
||||
t.Fatalf("api with session: %d %q", code, body)
|
||||
}
|
||||
if code, body := get(t, c, app+"/auth/me"); code != 200 || !strings.Contains(body, `"user":"alice"`) {
|
||||
t.Fatalf("me: %d %q", code, body)
|
||||
}
|
||||
// Logout drops the session (the fake provider has no logout page: 404).
|
||||
get(t, c, app+"/auth/logout")
|
||||
if code, _ := get(t, c, app+"/api/logs"); code != http.StatusUnauthorized {
|
||||
t.Fatalf("api after logout: %d", code)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestOIDCRejectsBadTokens(t *testing.T) {
|
||||
cases := map[string]func(map[string]any){
|
||||
"wrong nonce": func(c map[string]any) { c["nonce"] = "x" },
|
||||
"wrong audience": func(c map[string]any) { c["aud"] = "other" },
|
||||
"wrong issuer": func(c map[string]any) { c["iss"] = "https://evil" },
|
||||
"expired": func(c map[string]any) { c["exp"] = time.Now().Add(-time.Hour).Unix() },
|
||||
}
|
||||
for name, change := range cases {
|
||||
idp := newFakeIdP(t)
|
||||
idp.claims = change
|
||||
app, c := newOIDCApp(t, idp)
|
||||
if code, _ := get(t, c, app+"/"); code != http.StatusForbidden {
|
||||
t.Errorf("%s: login gave %d, expected 403", name, code)
|
||||
}
|
||||
if code, _ := get(t, c, app+"/api/logs"); code != http.StatusUnauthorized {
|
||||
t.Errorf("%s: session created", name)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestOIDCCallbackNeedsLoginCookie(t *testing.T) {
|
||||
idp := newFakeIdP(t)
|
||||
app, c := newOIDCApp(t, idp)
|
||||
if code, _ := get(t, c, app+"/auth/callback?code=abc&state=forged"); code != http.StatusBadRequest {
|
||||
t.Fatalf("forged callback: %d", code)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOIDCForgedSessionCookie(t *testing.T) {
|
||||
idp := newFakeIdP(t)
|
||||
app, c := newOIDCApp(t, idp)
|
||||
u, _ := url.Parse(app)
|
||||
payload := base64.RawURLEncoding.EncodeToString([]byte(`{"u":"mallory","e":9999999999}`))
|
||||
c.Jar.SetCookies(u, []*http.Cookie{{Name: sessionCookie, Value: payload + ".AAAA"}})
|
||||
if code, _ := get(t, c, app+"/api/logs"); code != http.StatusUnauthorized {
|
||||
t.Fatalf("forged session accepted: %d", code)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAuthModeConfig(t *testing.T) {
|
||||
next := http.NotFoundHandler()
|
||||
if _, err := newAuth(authConfig{mode: "oidc"}, next); err == nil || !strings.Contains(err.Error(), "OIDC_CLIENT_ID") {
|
||||
t.Errorf("missing variables not reported: %v", err)
|
||||
}
|
||||
if _, err := newAuth(authConfig{mode: "ldap"}, next); err == nil {
|
||||
t.Error("unknown mode accepted")
|
||||
}
|
||||
if h, err := newAuth(authConfig{mode: "local"}, next); err != nil || h == nil {
|
||||
t.Errorf("local mode: %v", err)
|
||||
}
|
||||
}
|
||||
+25
-15
@@ -12,24 +12,36 @@ services:
|
||||
- "${HTTP_PORT:-8080}:8080"
|
||||
environment:
|
||||
VLOGS_URL: http://victorialogs:9428
|
||||
SYSLOG_PUBLIC_PORT: ${SYSLOG_PORT:-514} # shown in Settings > Sources (the mapping is set above)
|
||||
SYSLOG_PUBLIC_PORT: ${SYSLOG_PORT:-514} # le port d'ecoute syslog par defaut (attention aux ports <1024)
|
||||
TZ: ${TZ:-Europe/Paris}
|
||||
AUTH_USER: ${AUTH_USER:-} # leave empty to disable authentication
|
||||
AUTH_MODE: ${AUTH_MODE:-local} # local (Basic Auth ci-dessous) ou oidc
|
||||
AUTH_USER: ${AUTH_USER:-} # vide = pas d'authentification, on delegue ca au reverse proxy traefik
|
||||
AUTH_PASS: ${AUTH_PASS:-}
|
||||
RDNS: ${RDNS:-on} # replace IP hosts with their DNS name (PTR)
|
||||
DNS_SERVER: ${DNS_SERVER:-} # e.g. 192.168.1.1 to query your LAN DNS; empty = system resolver
|
||||
OIDC_ISSUER: ${OIDC_ISSUER:-}
|
||||
OIDC_CLIENT_ID: ${OIDC_CLIENT_ID:-}
|
||||
OIDC_CLIENT_SECRET: ${OIDC_CLIENT_SECRET:-}
|
||||
OIDC_REDIRECT_URL: ${OIDC_REDIRECT_URL:-}
|
||||
OIDC_SCOPES: ${OIDC_SCOPES:-openid profile email}
|
||||
OIDC_SESSION_TTL: ${OIDC_SESSION_TTL:-12h}
|
||||
RDNS: ${RDNS:-on} # resol dns
|
||||
DNS_SERVER: ${DNS_SERVER:-} # si resolv directe
|
||||
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
|
||||
DOCKER_LOGS: ${DOCKER_LOGS:-on} # collecte des logs des conteneurs Docker
|
||||
DOCKER_HOST: tcp://docker-proxy:2375 # lecture seul de l'API Docker
|
||||
DOCKER_BACKFILL: ${DOCKER_BACKFILL:-1h}
|
||||
HOST_LOGS_BACKFILL: ${HOST_LOGS_BACKFILL:-1h} # historique lu a l'activation des logs systeme de l'hote
|
||||
group_add:
|
||||
- "${HOST_LOGS_GID:-4}" # groupe autorise a lire les logs de l'hote (4 = adm sur Debian/Ubuntu)
|
||||
volumes:
|
||||
- logstream-data:/data # tags.json, docker.json (container choices)
|
||||
- logstream-data:/data # tags.json, docker.json (les choix des conteneurs)
|
||||
# logs systeme de l'hote, en lecture seule (source a activer dans Reglages > Sources)
|
||||
- /var/log:/host/var/log:ro # journal systemd persistant et fichiers texte
|
||||
- /run/log/journal:/host/run/log/journal:ro # journal systemd volatile
|
||||
labels:
|
||||
logstream.exclude: "true" # never collect Logstream's own logs
|
||||
logstream.exclude: "true" # pas de collect des logs logstream
|
||||
|
||||
victorialogs:
|
||||
# Pinned version: see "Updating" in the README before changing it
|
||||
image: victoriametrics/victoria-logs:v1.52.0
|
||||
container_name: logstream-victorialogs
|
||||
restart: unless-stopped
|
||||
@@ -37,7 +49,7 @@ services:
|
||||
- -storageDataPath=/vlogs
|
||||
- -retentionPeriod=${RETENTION:-30d}
|
||||
- -httpListenAddr=:9428
|
||||
- -delete.enable # required by "Delete all logs" in Settings
|
||||
- -delete.enable # obliger pour autoriser la purge de la base via l'interface
|
||||
volumes:
|
||||
- vlogs-data:/vlogs
|
||||
ports:
|
||||
@@ -45,9 +57,7 @@ services:
|
||||
- "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.
|
||||
# Docker proxy pour eviter de solliciter directement l'API docker
|
||||
image: tecnativa/docker-socket-proxy:v0.5.0
|
||||
container_name: logstream-docker-proxy
|
||||
restart: unless-stopped
|
||||
@@ -59,7 +69,7 @@ services:
|
||||
volumes:
|
||||
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||
labels:
|
||||
logstream.exclude: "true" # its access log would only echo Logstream's own requests
|
||||
logstream.exclude: "true" # pour exclure les logs propres de logstream
|
||||
|
||||
volumes:
|
||||
logstream-data:
|
||||
|
||||
File diff suppressed because it is too large.
Load diff
Binary file not shown.
|
After Width: | Height: | Size: 884 KiB |
File diff suppressed because it is too large.
Load diff
Binary file not shown.
|
Before Width: | Height: | Size: 597 KiB |
+452
@@ -0,0 +1,452 @@
|
||||
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}
|
||||
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"]},
|
||||
})
|
||||
}
|
||||
|
||||
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,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,206 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
+174
@@ -0,0 +1,174 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"encoding/binary"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Minimal reader of systemd journal files (*.journal), enough to follow them
|
||||
// without journalctl: objects are walked in file order, which is the order
|
||||
// entries were written. Format: https://systemd.io/JOURNAL_FILE_FORMAT/
|
||||
|
||||
const (
|
||||
jHeaderMin = 208 // header fields used below end at offset 208
|
||||
jObjHeaderSize = 16
|
||||
|
||||
jObjData = 1
|
||||
jObjEntry = 3
|
||||
|
||||
jIncompatCompact = 1 << 4
|
||||
jStateArchived = 2
|
||||
jObjCompressed = 1<<0 | 1<<1 | 1<<2 // XZ, LZ4, ZSTD: not decoded (no stdlib codec)
|
||||
)
|
||||
|
||||
var jSignature = []byte("LPKSHHRH")
|
||||
|
||||
// jHeader holds the header fields the reader needs.
|
||||
type jHeader struct {
|
||||
compact bool
|
||||
archived bool
|
||||
fileID string
|
||||
headerSize uint64
|
||||
tailObject uint64 // offset of the last object
|
||||
tailSeqnum uint64 // seqnum of the last complete entry
|
||||
tailRealtimeUS uint64 // realtime of the last entry (µs since the epoch)
|
||||
}
|
||||
|
||||
func readJournalHeader(f io.ReaderAt) (jHeader, error) {
|
||||
b := make([]byte, jHeaderMin)
|
||||
if _, err := f.ReadAt(b, 0); err != nil {
|
||||
return jHeader{}, err
|
||||
}
|
||||
if !bytes.Equal(b[:8], jSignature) {
|
||||
return jHeader{}, errors.New("not a journal file")
|
||||
}
|
||||
le := binary.LittleEndian
|
||||
h := jHeader{
|
||||
compact: le.Uint32(b[12:])&jIncompatCompact != 0,
|
||||
archived: b[16] == jStateArchived,
|
||||
fileID: hex.EncodeToString(b[24:40]),
|
||||
headerSize: le.Uint64(b[88:]),
|
||||
tailObject: le.Uint64(b[136:]),
|
||||
tailSeqnum: le.Uint64(b[160:]),
|
||||
tailRealtimeUS: le.Uint64(b[192:]),
|
||||
}
|
||||
if h.headerSize < jHeaderMin {
|
||||
return h, errors.New("journal header too small")
|
||||
}
|
||||
return h, nil
|
||||
}
|
||||
|
||||
// journalEntry is one entry: its time and its "FIELD=value" pairs.
|
||||
type journalEntry struct {
|
||||
Realtime time.Time
|
||||
Fields map[string]string
|
||||
Compressed int // fields that could not be read (compressed data)
|
||||
}
|
||||
|
||||
// readJournal reads the entries complete after offset `from` (0 = start of
|
||||
// the file) and returns the offset to resume from. Entries older than
|
||||
// `notBefore` are skipped without decoding their fields.
|
||||
func readJournal(path string, from int64, notBefore time.Time, emit func(*journalEntry)) (next int64, h jHeader, err error) {
|
||||
f, err := os.Open(path)
|
||||
if err != nil {
|
||||
return from, h, err
|
||||
}
|
||||
defer f.Close()
|
||||
if h, err = readJournalHeader(f); err != nil {
|
||||
return from, h, err
|
||||
}
|
||||
if from < int64(h.headerSize) {
|
||||
from = int64(h.headerSize)
|
||||
}
|
||||
end := int64(h.tailObject)
|
||||
r := bufio.NewReaderSize(io.NewSectionReader(f, from, 1<<62), 64*1024)
|
||||
le := binary.LittleEndian
|
||||
hdr := make([]byte, jObjHeaderSize)
|
||||
off := from
|
||||
for off <= end {
|
||||
if _, err := io.ReadFull(r, hdr); err != nil {
|
||||
return off, h, nil // object still being written
|
||||
}
|
||||
typ, size := hdr[0], int64(le.Uint64(hdr[8:]))
|
||||
if size < jObjHeaderSize {
|
||||
return off, h, nil
|
||||
}
|
||||
body := size - jObjHeaderSize
|
||||
if typ == jObjEntry && body >= 48 && body < 1<<20 {
|
||||
obj := make([]byte, body)
|
||||
if _, err := io.ReadFull(r, obj); err != nil {
|
||||
return off, h, nil
|
||||
}
|
||||
seq := le.Uint64(obj[0:])
|
||||
if seq == 0 || seq > h.tailSeqnum {
|
||||
return off, h, nil // not linked yet: retry at the next poll
|
||||
}
|
||||
rt := time.UnixMicro(int64(le.Uint64(obj[8:])))
|
||||
if !rt.Before(notBefore) {
|
||||
emit(readEntryFields(f, h.compact, rt, obj[48:]))
|
||||
}
|
||||
} else if _, err := r.Discard(int(body)); err != nil {
|
||||
return off, h, nil
|
||||
}
|
||||
next := (off + size + 7) &^ 7 // objects are 8-byte aligned
|
||||
if pad := next - off - size; pad > 0 {
|
||||
if _, err := r.Discard(int(pad)); err != nil {
|
||||
return next, h, nil // the object is complete: resume after it
|
||||
}
|
||||
}
|
||||
off = next
|
||||
}
|
||||
return off, h, nil
|
||||
}
|
||||
|
||||
// readEntryFields resolves the data objects referenced by an entry.
|
||||
func readEntryFields(f io.ReaderAt, compact bool, rt time.Time, items []byte) *journalEntry {
|
||||
le := binary.LittleEndian
|
||||
e := &journalEntry{Realtime: rt, Fields: make(map[string]string, 16)}
|
||||
step := 16
|
||||
if compact {
|
||||
step = 4
|
||||
}
|
||||
payloadAt := int64(64)
|
||||
if compact {
|
||||
payloadAt = 72
|
||||
}
|
||||
hdr := make([]byte, jObjHeaderSize)
|
||||
for i := 0; i+step <= len(items); i += step {
|
||||
var off int64
|
||||
if compact {
|
||||
off = int64(le.Uint32(items[i:]))
|
||||
} else {
|
||||
off = int64(le.Uint64(items[i:]))
|
||||
}
|
||||
if off == 0 {
|
||||
continue
|
||||
}
|
||||
if _, err := f.ReadAt(hdr, off); err != nil || hdr[0] != jObjData {
|
||||
continue
|
||||
}
|
||||
if hdr[1]&jObjCompressed != 0 {
|
||||
e.Compressed++
|
||||
continue
|
||||
}
|
||||
size := int64(le.Uint64(hdr[8:]))
|
||||
n := size - payloadAt
|
||||
if n <= 0 || n > 1<<20 {
|
||||
continue
|
||||
}
|
||||
p := make([]byte, n)
|
||||
if _, err := f.ReadAt(p, off+payloadAt); err != nil {
|
||||
continue
|
||||
}
|
||||
if k := bytes.IndexByte(p, '='); k > 0 {
|
||||
e.Fields[string(p[:k])] = string(p[k+1:])
|
||||
}
|
||||
}
|
||||
return e
|
||||
}
|
||||
@@ -1,10 +1,9 @@
|
||||
// Logstream: a syslog (UDP/TCP) sink stored in VictoriaLogs,
|
||||
// Logstream: a syslog (UDP/TCP), Docker and host system logs sink stored in VictoriaLogs,
|
||||
// with a web interface to browse and search the logs.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/subtle"
|
||||
"embed"
|
||||
"errors"
|
||||
"io/fs"
|
||||
@@ -29,8 +28,7 @@ type config struct {
|
||||
httpAddr string
|
||||
vlogsURL string
|
||||
dataDir string
|
||||
authUser string
|
||||
authPass string
|
||||
auth authConfig
|
||||
rdns bool
|
||||
dnsServer string
|
||||
allowPurge bool
|
||||
@@ -38,6 +36,8 @@ type config struct {
|
||||
dockerLogs bool
|
||||
dockerHost string
|
||||
backfill time.Duration
|
||||
hostRoot string
|
||||
hostFill time.Duration
|
||||
batchSize int
|
||||
queueSize int
|
||||
flushEvery time.Duration
|
||||
@@ -80,8 +80,17 @@ func main() {
|
||||
httpAddr: getenv("HTTP_ADDR", ":8080"),
|
||||
vlogsURL: getenv("VLOGS_URL", "http://victorialogs:9428"),
|
||||
dataDir: getenv("DATA_DIR", "/data"),
|
||||
authUser: os.Getenv("AUTH_USER"),
|
||||
authPass: os.Getenv("AUTH_PASS"),
|
||||
auth: authConfig{
|
||||
mode: getenv("AUTH_MODE", "local"),
|
||||
user: os.Getenv("AUTH_USER"),
|
||||
pass: os.Getenv("AUTH_PASS"),
|
||||
issuer: os.Getenv("OIDC_ISSUER"),
|
||||
clientID: os.Getenv("OIDC_CLIENT_ID"),
|
||||
clientSecret: os.Getenv("OIDC_CLIENT_SECRET"),
|
||||
redirectURL: os.Getenv("OIDC_REDIRECT_URL"),
|
||||
scopes: os.Getenv("OIDC_SCOPES"),
|
||||
sessionTTL: getenvDuration("OIDC_SESSION_TTL", 12*time.Hour),
|
||||
},
|
||||
rdns: getenvBool("RDNS", true),
|
||||
dnsServer: os.Getenv("DNS_SERVER"),
|
||||
allowPurge: getenvBool("ALLOW_PURGE", true),
|
||||
@@ -89,11 +98,15 @@ func main() {
|
||||
dockerLogs: getenvBool("DOCKER_LOGS", false),
|
||||
dockerHost: getenv("DOCKER_HOST", "unix:///var/run/docker.sock"),
|
||||
backfill: getenvDuration("DOCKER_BACKFILL", time.Hour),
|
||||
hostRoot: getenv("HOST_LOGS_ROOT", "/host"),
|
||||
hostFill: getenvDuration("HOST_LOGS_BACKFILL", time.Hour),
|
||||
batchSize: getenvInt("BATCH_SIZE", 1000),
|
||||
queueSize: getenvInt("QUEUE_SIZE", 100000),
|
||||
flushEvery: time.Duration(getenvInt("FLUSH_MS", 1000)) * time.Millisecond,
|
||||
}
|
||||
|
||||
cfg.auth.dataDir = cfg.dataDir
|
||||
|
||||
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
||||
defer stop()
|
||||
|
||||
@@ -139,12 +152,22 @@ func main() {
|
||||
log.Printf("docker logs enabled through %s", cfg.dockerHost)
|
||||
}
|
||||
}
|
||||
// System logs of the host (journal or /var/log), off until enabled in Settings > Sources.
|
||||
api.host = NewHostLogs(cfg.hostRoot, cfg.dataDir, cfg.hostFill, sink)
|
||||
go api.host.Run(ctx)
|
||||
api.Routes(mux)
|
||||
mux.Handle("GET /", http.FileServer(http.FS(static)))
|
||||
|
||||
handler, err := newAuth(cfg.auth, mux)
|
||||
if err != nil {
|
||||
log.Fatalf("auth: %v", err)
|
||||
}
|
||||
if o, ok := handler.(*OIDC); ok {
|
||||
go o.checkProvider()
|
||||
}
|
||||
srv := &http.Server{
|
||||
Addr: cfg.httpAddr,
|
||||
Handler: basicAuth(cfg.authUser, cfg.authPass, mux),
|
||||
Handler: handler,
|
||||
ReadHeaderTimeout: 10 * time.Second,
|
||||
// Requests inherit the global context so SSE streams end on shutdown.
|
||||
BaseContext: func(net.Listener) context.Context { return ctx },
|
||||
@@ -163,25 +186,3 @@ func main() {
|
||||
<-storeDone
|
||||
log.Println("shutdown complete")
|
||||
}
|
||||
|
||||
// basicAuth protects the UI when AUTH_USER is set (except /healthz).
|
||||
func basicAuth(user, pass string, next http.Handler) http.Handler {
|
||||
if user == "" {
|
||||
return next
|
||||
}
|
||||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Path == "/healthz" {
|
||||
next.ServeHTTP(w, r)
|
||||
return
|
||||
}
|
||||
u, p, ok := r.BasicAuth()
|
||||
if !ok ||
|
||||
subtle.ConstantTimeCompare([]byte(u), []byte(user)) != 1 ||
|
||||
subtle.ConstantTimeCompare([]byte(p), []byte(pass)) != 1 {
|
||||
w.Header().Set("WWW-Authenticate", `Basic realm="logstream"`)
|
||||
http.Error(w, "authentication required", http.StatusUnauthorized)
|
||||
return
|
||||
}
|
||||
next.ServeHTTP(w, r)
|
||||
})
|
||||
}
|
||||
@@ -17,7 +17,7 @@ type Filter struct {
|
||||
Host string
|
||||
App string
|
||||
Severity int // highest severity number included (0 = emerg … 7 = debug), -1 = all
|
||||
Source string // "syslog", "docker" or "" for all
|
||||
Source string // "syslog", "docker", "host" or "" for all
|
||||
}
|
||||
|
||||
var rangeDurations = map[string]time.Duration{
|
||||
@@ -119,10 +119,10 @@ func (f Filter) filterExpr() string {
|
||||
parts = append(parts, "app:="+strconv.Quote(f.App))
|
||||
}
|
||||
switch f.Source {
|
||||
case "docker":
|
||||
parts = append(parts, `source_type:="docker"`)
|
||||
case "docker", "host":
|
||||
parts = append(parts, `source_type:=`+strconv.Quote(f.Source))
|
||||
case "syslog": // also matches logs stored before source_type existed
|
||||
parts = append(parts, `!(source_type:="docker")`)
|
||||
parts = append(parts, `!(source_type:in("docker","host"))`)
|
||||
}
|
||||
if f.Severity >= 0 && f.Severity < 7 {
|
||||
names := make([]string, 0, 8)
|
||||
@@ -209,7 +209,7 @@ 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 != "" {
|
||||
if m.f.Source != "" && m.f.Source != e.SourceType {
|
||||
return false
|
||||
}
|
||||
if m.f.Severity >= 0 && e.SevNum > m.f.Severity {
|
||||
|
||||
+92
-4
@@ -14,6 +14,7 @@ const I18N = {
|
||||
liveUnavailable: 'Live view is not available in LogsQL mode',
|
||||
settings: 'Settings',
|
||||
theme: 'Light / dark theme',
|
||||
logout: 'Log out',
|
||||
close: 'Close',
|
||||
rangeAria: 'Time range', severityAria: 'Severity', hostAria: 'Host', appAria: 'Application',
|
||||
histoAria: 'Log volume over time',
|
||||
@@ -67,6 +68,19 @@ const I18N = {
|
||||
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',
|
||||
srcHost: 'Host system',
|
||||
hostTitle: 'Host system logs',
|
||||
hostEnabled: 'Collect the system logs of this machine',
|
||||
hostOff: 'Off: the system logs of the machine hosting Logstream are not collected.',
|
||||
hostWaiting: 'Starting…',
|
||||
hostJournal: ({ n, r }) => `Reading the systemd journal (${n} files): ${r} entries since startup.`,
|
||||
hostFiles: ({ n, r }) => `Following ${n} files of /var/log: ${r} lines since startup.`,
|
||||
hostCompressed: (n) => ` ${n} compressed fields skipped (long messages, readable with journalctl).`,
|
||||
host_not_mounted: 'No host log directory found: mount /var/log and /run/log/journal under /host (see docker-compose.yml).',
|
||||
host_permission: 'Permission denied: give the container a group allowed to read the logs (HOST_LOGS_GID in .env, see the README). ',
|
||||
host_empty: 'No systemd journal and no log file found in /var/log.',
|
||||
host_read: 'Cannot read the host logs: ',
|
||||
hostHelp: 'Reads the systemd journal of the host (/var/log/journal, /run/log/journal), or the text files of /var/log (syslog, messages, *.log) when there is no journal. Directories are mounted read-only. When turned on, the last hour is read first (HOST_LOGS_BACKFILL).',
|
||||
syslogEnabled: 'Receive syslog messages', syslogProtocols: 'Protocols', syslogPort: 'Port',
|
||||
syslogListening: ({ port, protos }) => `Listening on port ${port} (${protos}).`,
|
||||
syslogOff: 'Syslog reception is off: syslog messages are ignored (Docker logs are still collected).',
|
||||
@@ -88,6 +102,7 @@ const I18N = {
|
||||
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',
|
||||
fUnit: 'systemd unit', fLogFile: 'log file',
|
||||
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.`,
|
||||
@@ -143,6 +158,7 @@ const I18N = {
|
||||
liveUnavailable: 'Le direct n\'est pas disponible en mode LogsQL',
|
||||
settings: 'Paramètres',
|
||||
theme: 'Thème clair / sombre',
|
||||
logout: 'Se déconnecter',
|
||||
close: 'Fermer',
|
||||
rangeAria: 'Période', severityAria: 'Sévérité', hostAria: 'Hôte', appAria: 'Application',
|
||||
histoAria: 'Volume de logs dans le temps',
|
||||
@@ -196,6 +212,19 @@ const I18N = {
|
||||
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',
|
||||
srcHost: 'Système hôte',
|
||||
hostTitle: 'Logs système de l\'hôte',
|
||||
hostEnabled: 'Collecter les logs système de cette machine',
|
||||
hostOff: 'Désactivé : les logs système de la machine qui héberge Logstream ne sont pas collectés.',
|
||||
hostWaiting: 'Démarrage…',
|
||||
hostJournal: ({ n, r }) => `Lecture du journal systemd (${n} fichiers) : ${r} entrées depuis le démarrage.`,
|
||||
hostFiles: ({ n, r }) => `Suivi de ${n} fichiers de /var/log : ${r} lignes depuis le démarrage.`,
|
||||
hostCompressed: (n) => ` ${n} champs compressés ignorés (messages longs, lisibles avec journalctl).`,
|
||||
host_not_mounted: 'Aucun répertoire de logs de l\'hôte trouvé : montez /var/log et /run/log/journal sous /host (voir docker-compose.yml).',
|
||||
host_permission: 'Accès refusé : donnez au conteneur un groupe autorisé à lire les logs (HOST_LOGS_GID dans .env, voir le README). ',
|
||||
host_empty: 'Ni journal systemd ni fichier de log trouvé dans /var/log.',
|
||||
host_read: 'Lecture des logs de l\'hôte impossible : ',
|
||||
hostHelp: 'Lit le journal systemd de l\'hôte (/var/log/journal, /run/log/journal), ou les fichiers texte de /var/log (syslog, messages, *.log) s\'il n\'y a pas de journal. Les répertoires sont montés en lecture seule. À l\'activation, la dernière heure est lue d\'abord (HOST_LOGS_BACKFILL).',
|
||||
syslogEnabled: 'Recevoir les messages syslog', syslogProtocols: 'Protocoles', syslogPort: 'Port',
|
||||
syslogListening: ({ port, protos }) => `Écoute sur le port ${port} (${protos}).`,
|
||||
syslogOff: 'Réception syslog désactivée : les messages syslog sont ignorés (les logs Docker sont toujours collectés).',
|
||||
@@ -217,6 +246,7 @@ const I18N = {
|
||||
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',
|
||||
fUnit: 'unité systemd', fLogFile: 'fichier de log',
|
||||
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.`,
|
||||
@@ -381,6 +411,8 @@ async function api(url, opts = {}) {
|
||||
|
||||
// Known error codes are translated; otherwise the server message is shown.
|
||||
function apiError(res, text, data) {
|
||||
// OIDC session expired: reloading the page goes through the login again.
|
||||
if (res.status === 401 && data && data.code === 'auth') location.reload();
|
||||
let msg = (data && data.error) || text || res.statusText;
|
||||
if (data && data.code && I18N[lang]['err_' + data.code]) {
|
||||
msg = t('err_' + data.code) + (data.detail ? (lang === 'fr' ? ' : ' : ': ') + data.detail : '');
|
||||
@@ -424,7 +456,7 @@ const list = $('#list');
|
||||
function applyLang() {
|
||||
document.documentElement.lang = lang;
|
||||
for (const el of document.querySelectorAll('[data-i18n]')) el.textContent = t(el.dataset.i18n);
|
||||
for (const el of document.querySelectorAll('[data-i18n-title]')) el.title = t(el.dataset.i18nTitle);
|
||||
for (const el of document.querySelectorAll('[data-i18n-title]')) el.title = t(el.dataset.i18nTitle) + (el.dataset.user ? ` (${el.dataset.user})` : '');
|
||||
for (const el of document.querySelectorAll('[data-i18n-aria]')) el.setAttribute('aria-label', t(el.dataset.i18nAria));
|
||||
for (const el of document.querySelectorAll('[data-i18n-ph]')) el.placeholder = t(el.dataset.i18nPh);
|
||||
for (const b of document.querySelectorAll('#langSwitch [data-lang]')) b.setAttribute('aria-checked', String(b.dataset.lang === lang));
|
||||
@@ -447,6 +479,7 @@ function setLang(next) {
|
||||
renderPurge();
|
||||
renderInterface();
|
||||
renderSyslog();
|
||||
renderHost();
|
||||
renderDocker();
|
||||
}
|
||||
|
||||
@@ -805,14 +838,14 @@ function updateCount(tookMs) {
|
||||
}
|
||||
|
||||
function detailsHTML(r) {
|
||||
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 order = ['received', 'msg_time', '_time', 'host', 'host_name', 'host_ip', 'source_type', 'container', 'image', 'compose_project', 'compose_service', 'stream', 'unit', 'log_file', '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 = {
|
||||
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'),
|
||||
compose_service: t('fService'), stream: t('fStream'), source_type: t('fSourceType'), unit: t('fUnit'), log_file: t('fLogFile'),
|
||||
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]))} <span class="muted">(${esc(r[k])})</span>`
|
||||
@@ -1726,6 +1759,48 @@ $('#syslogProtos').addEventListener('click', (ev) => {
|
||||
configureSyslog({ [b.dataset.proto]: !syslogState[b.dataset.proto] });
|
||||
});
|
||||
|
||||
/* ================= Host system logs ================= */
|
||||
|
||||
let hostState = null;
|
||||
|
||||
async function loadHost() {
|
||||
try { hostState = await api('/api/hostlogs'); } catch { hostState = null; }
|
||||
renderHost();
|
||||
}
|
||||
|
||||
function renderHost() {
|
||||
const s = hostState;
|
||||
const status = $('#hostStatus');
|
||||
if (!s) { status.textContent = ''; return; }
|
||||
$('#hostEnabled').checked = s.enabled;
|
||||
let text = t('hostWaiting');
|
||||
let cls = 'muted';
|
||||
if (!s.enabled) {
|
||||
text = t('hostOff');
|
||||
} else if (s.code) {
|
||||
text = t('host_' + s.code) + (s.code === 'read' || s.code === 'permission' ? s.error || '' : '');
|
||||
cls = 'bad';
|
||||
} else if (s.mode === 'journal') {
|
||||
text = t('hostJournal', { n: s.files, r: fmtNum(s.read) }) + (s.compressed ? t('hostCompressed', fmtNum(s.compressed)) : '');
|
||||
cls = 'good';
|
||||
} else if (s.mode === 'files') {
|
||||
text = t('hostFiles', { n: s.files, r: fmtNum(s.read) });
|
||||
cls = 'good';
|
||||
}
|
||||
status.textContent = text;
|
||||
status.className = 'docker-status ' + cls;
|
||||
}
|
||||
|
||||
async function configureHost(enabled) {
|
||||
hostState = { ...(hostState || {}), enabled };
|
||||
renderHost();
|
||||
try { hostState = await api('/api/hostlogs', { method: 'PUT', body: { enabled } }); } catch (e) { toast(e.message); }
|
||||
renderHost();
|
||||
setTimeout(loadHost, 1500); // the first scan runs in the background
|
||||
}
|
||||
|
||||
$('#hostEnabled').addEventListener('change', (ev) => configureHost(ev.target.checked));
|
||||
|
||||
/* ================= Docker source ================= */
|
||||
|
||||
let dockerState = null;
|
||||
@@ -1839,7 +1914,10 @@ for (const [id, on] of [['#dockerAll', true], ['#dockerNone', false]]) {
|
||||
}
|
||||
// Keep the list fresh while the Sources tab is open.
|
||||
setInterval(() => {
|
||||
if ($('#settingsDlg').open && !document.querySelector('[data-panel="sources"]').hidden) loadDocker();
|
||||
if ($('#settingsDlg').open && !document.querySelector('[data-panel="sources"]').hidden) {
|
||||
loadDocker();
|
||||
loadHost();
|
||||
}
|
||||
}, 4000);
|
||||
|
||||
/* ================= Purge ================= */
|
||||
@@ -2002,6 +2080,7 @@ $('#settingsBtn').addEventListener('click', () => {
|
||||
renderInterface();
|
||||
loadPurgeStatus();
|
||||
loadSyslog();
|
||||
loadHost();
|
||||
loadDocker();
|
||||
showSettingsTab(store.get('settingsTab', 'locale'));
|
||||
$('#settingsDlg').showModal();
|
||||
@@ -2024,6 +2103,15 @@ $('#range').value = store.get('range', '1h');
|
||||
if (!$('#range').value) $('#range').value = '1h';
|
||||
$('#severity').value = store.get('severity', '');
|
||||
|
||||
// With OIDC login, show who is logged in and the log out button.
|
||||
fetch('/auth/me').then((res) => (res.ok ? res.json() : null)).then((me) => {
|
||||
if (!me || !me.user) return;
|
||||
const btn = $('#logoutBtn');
|
||||
btn.hidden = false;
|
||||
btn.dataset.user = me.user;
|
||||
btn.title = `${t('logout')} (${me.user})`;
|
||||
}).catch(() => {});
|
||||
|
||||
(async () => {
|
||||
await loadTags();
|
||||
loadFacets();
|
||||
|
||||
@@ -41,6 +41,10 @@
|
||||
<svg class="sun" viewBox="0 0 24 24"><circle cx="12" cy="12" r="4"/><path d="M12 2v2M12 20v2M4.9 4.9l1.4 1.4M17.7 17.7l1.4 1.4M2 12h2M20 12h2M4.9 19.1l1.4-1.4M17.7 6.3l1.4-1.4"/></svg>
|
||||
<svg class="moon" viewBox="0 0 24 24"><path d="M21 12.8A9 9 0 1 1 11.2 3a7 7 0 0 0 9.8 9.8z"/></svg>
|
||||
</button>
|
||||
<a id="logoutBtn" class="icon-btn" href="/auth/logout" hidden data-i18n-title="logout" data-i18n-aria="logout">
|
||||
<!-- Log out icon (Lucide "log-out", ISC license) -->
|
||||
<svg viewBox="0 0 24 24"><path d="M9 21H5a2 2 0 0 1-2-2V5a2 2 0 0 1 2-2h4"/><path d="m16 17 5-5-5-5"/><path d="M21 12H9"/></svg>
|
||||
</a>
|
||||
</div>
|
||||
<div class="progress" aria-hidden="true"></div>
|
||||
</header>
|
||||
@@ -70,6 +74,7 @@
|
||||
<option value="" data-i18n="srcAll">All sources</option>
|
||||
<option value="syslog">Syslog</option>
|
||||
<option value="docker">Docker</option>
|
||||
<option value="host" data-i18n="srcHost">Host system</option>
|
||||
</select>
|
||||
<span class="spacer"></span>
|
||||
<span id="count" class="muted"></span>
|
||||
@@ -209,6 +214,16 @@
|
||||
</div>
|
||||
</section>
|
||||
|
||||
<section class="set-section">
|
||||
<h3 data-i18n="hostTitle">Host system logs</h3>
|
||||
<p id="hostStatus" class="docker-status"></p>
|
||||
<label class="switch-row">
|
||||
<input id="hostEnabled" type="checkbox" class="switch">
|
||||
<span data-i18n="hostEnabled">Collect the system logs of this machine</span>
|
||||
</label>
|
||||
<p class="muted small" data-i18n="hostHelp"></p>
|
||||
</section>
|
||||
|
||||
<section class="set-section">
|
||||
<h3 data-i18n="dockerTitle">Docker containers</h3>
|
||||
<p id="dockerStatus" class="docker-status"></p>
|
||||
|
||||
Reference in new issue
Block a user