Phase 3: - Config reload via SIGHUP (handle_signals loop: .hup reloads config, .int/.term shutdown) - Hot reload: cfg, webhook_routes, ntfy client, dedup updated without restart Phase 4: - Silence rules: POST /api/silence?duration=30m, GET /api/silence, DELETE /api/silence - Grouping: buffer 10s, events grouped with (×N in 10s) suffix - CI/CD: .gitea/workflows/ci.yml (build + test on push/PR) Docs: - ROADMAP.md: phases 3 & 4 marked DONE, version table updated - ARCHITECTURE.md: synced with actual code Version: 0.1.0 → 0.4.0 (main.v, v.mod)
This commit is contained in:
@@ -0,0 +1,27 @@
|
||||
name: CI
|
||||
|
||||
on:
|
||||
push:
|
||||
branches: [main]
|
||||
pull_request:
|
||||
branches: [main]
|
||||
|
||||
jobs:
|
||||
test:
|
||||
runs-on: ubuntu-latest
|
||||
container:
|
||||
image: thevlang/vlang:alpine
|
||||
steps:
|
||||
- name: Checkout
|
||||
uses: actions/checkout@v4
|
||||
|
||||
- name: Build
|
||||
run: v -prod .
|
||||
|
||||
- name: Test
|
||||
run: v test .
|
||||
|
||||
- name: Verify version
|
||||
run: |
|
||||
./ntfy-bridge --version
|
||||
./ntfy-bridge --help
|
||||
+135
-106
@@ -2,14 +2,15 @@
|
||||
|
||||
## Philosophie
|
||||
|
||||
ntfy-bridge est un **daemon HTTP léger** écrit en V (~500 lignes) qui compile en un seul binaire natif statique. Il écoute des webhooks, se connecte au socket Docker, et poll des endpoints HTTP — puis transforme chaque événement en notification Ntfy formatée.
|
||||
ntfy-bridge est un **daemon HTTP léger** écrit en V (~1000 lignes) qui compile en un seul binaire natif statique. Il écoute des webhooks, se connecte au socket Docker, et poll des endpoints HTTP — puis transforme chaque événement en notification Ntfy formatée.
|
||||
|
||||
**Principes :**
|
||||
- **Un binaire, zéro runtime** — pas de Node, Python, JVM ou conteneur obligatoire
|
||||
- **Config-driven** — tout passe par le fichier YAML, pas de recompilation
|
||||
- **Fail-safe** — un échec réseau vers Ntfy ne crashe pas le daemon
|
||||
- **Fail-safe** — un échec réseau vers Ntfy ne crashe pas le daemon, retry avec backoff
|
||||
- **Silencieux par défaut** — ne spamme pas, utilise les bons niveaux de priorité Ntfy (1-5)
|
||||
- **Extensible** — ajouter une source = implémenter un module dans `sources/`
|
||||
- **Extensible** — ajouter une source = ajouter un fichier dans `sources/` + une route webhook
|
||||
- **Zéro dépendance externe** — 100% stdlib V (`net.http`, `json`, `yaml`, `time`, `os`, `flag`)
|
||||
|
||||
## Architecture globale
|
||||
|
||||
@@ -18,9 +19,9 @@ ntfy-bridge est un **daemon HTTP léger** écrit en V (~500 lignes) qui compile
|
||||
│ ntfy-bridge │
|
||||
│ │
|
||||
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐ │
|
||||
│ │ HTTP │ │ Docker │ │ HTTP │ │ Cron │ │
|
||||
│ │ Server │ │ Watcher │ │ Poller │ │ Receiver │ │
|
||||
│ │ :9090 │ │ goroutine│ │ goroutine│ │ /webhooks │ │
|
||||
│ │ HTTP │ │ Docker │ │ HTTP │ │ Cron/ │ │
|
||||
│ │ Server │ │ Watcher │ │ Poller │ │ Generic │ │
|
||||
│ │ :9090 │ │ goroutine│ │ goroutine│ │ Receiver │ │
|
||||
│ └────┬─────┘ └────┬─────┘ └────┬─────┘ └─────┬──────┘ │
|
||||
│ │ │ │ │ │
|
||||
│ └──────────────┼──────────────┼──────────────┘ │
|
||||
@@ -34,49 +35,53 @@ ntfy-bridge est un **daemon HTTP léger** écrit en V (~500 lignes) qui compile
|
||||
│ └────────────┬─────────────┘ │
|
||||
│ ▼ │
|
||||
│ ┌──────────────────────────┐ │
|
||||
│ │ Ntfy Publisher │ │
|
||||
│ │ Ntfy Publisher (retry) │ │
|
||||
│ │ POST /<topic> │ │
|
||||
│ │ 3 tentatives, backoff │ │
|
||||
│ └──────────────────────────┘ │
|
||||
└─────────────────────────────────────────────────────────────┘
|
||||
```
|
||||
|
||||
## Flux de données
|
||||
|
||||
### 1. Webhook entrant (Gitea, Uptime Kuma, generic)
|
||||
### 1. Webhook entrant (Gitea, Uptime Kuma, Cron, Generic)
|
||||
```
|
||||
POST /webhooks/gitea-flowdeck
|
||||
│
|
||||
▼
|
||||
1. Verify signature (optionnel)
|
||||
2. Parse JSON → Event struct
|
||||
3. Route vers le transformer approprié
|
||||
4. Appliquer priority_map
|
||||
5. Appliquer template de message
|
||||
6. POST vers Ntfy
|
||||
1. O(1) lookup dans webhook_routes (map[string]WebhookRoute)
|
||||
2. Verify HMAC-SHA256 signature (si configuré, global)
|
||||
3. Dispatcher vers le transformer approprié (sources/transform_*.v)
|
||||
4. Parse JSON → Event struct
|
||||
5. Appliquer priority_map
|
||||
6. Appliquer template de message (via render_source_template)
|
||||
7. Pipeline publish: apply_defaults → dedup.allow → ntfy.publish
|
||||
8. POST vers Ntfy avec retry (3 tentatives, backoff 1s/2s/4s)
|
||||
```
|
||||
|
||||
### 2. Docker watcher (socket)
|
||||
```
|
||||
goroutine DockerWatcher
|
||||
goroutine docker_watch_loop (via start_background_tasks)
|
||||
│
|
||||
▼
|
||||
1. Connexion au socket Docker (/var/run/docker.sock)
|
||||
2. Stream events via GET /events?filters=...
|
||||
1. Connexion au socket Docker Unix (/var/run/docker.sock)
|
||||
2. Stream events via GET /events?filters={"die":true,"oom":true,...}
|
||||
3. Filtrer par événements configurés (die, oom, health_status)
|
||||
4. Extraire container_name, image, exit_code
|
||||
4. Extraire container_name, image, exit_code via transform_docker_event
|
||||
5. Appliquer template
|
||||
6. POST vers Ntfy
|
||||
6. Pipeline publish → Ntfy
|
||||
```
|
||||
|
||||
### 3. HTTP Poller
|
||||
```
|
||||
goroutine HttpPoller (toutes les N secondes)
|
||||
goroutine http_poll_loop (via start_background_tasks, toutes les N secondes)
|
||||
│
|
||||
▼
|
||||
1. GET <url> avec timeout configuré
|
||||
2. Si status ≠ expected → Event(priority=5)
|
||||
3. Si OK et précédemment DOWN → Event(priority=1, tags="white_check_mark")
|
||||
2. Si status ≠ expected → Event(priority=configurée, tags=["x"])
|
||||
3. Si OK et précédemment DOWN → Event(priority=1, tags=["white_check_mark"])
|
||||
4. Sinon → rien (pas de notification si tout va bien)
|
||||
5. State machine: poll_state[url] = is_up → notification uniquement au changement
|
||||
```
|
||||
|
||||
### 4. Cron receiver
|
||||
@@ -85,69 +90,74 @@ POST /webhooks/cron-disk
|
||||
│
|
||||
▼
|
||||
1. Body = texte brut du script cron
|
||||
2. Pass-through vers Ntfy
|
||||
2. Pass-through vers Ntfy (transform_cron → Event)
|
||||
3. Le script cron contrôle son propre format
|
||||
```
|
||||
|
||||
## Modèle de données
|
||||
|
||||
```v
|
||||
struct Config { // dans config.v
|
||||
server ServerConfig
|
||||
defaults DefaultConfig
|
||||
sources SourcesConfig
|
||||
}
|
||||
|
||||
// config.v
|
||||
struct ServerConfig {
|
||||
url string // URL du serveur Ntfy
|
||||
auth_token string // Token d'auth (optionnel)
|
||||
auth_token string // Token d'auth (optionnel, override via NTFY_TOKEN env)
|
||||
listen string // Adresse d'écoute HTTP (:9090)
|
||||
hmac_secret string // Secret partagé HMAC-SHA256 (override via NTFY_HMAC_SECRET env)
|
||||
}
|
||||
|
||||
struct DefaultConfig {
|
||||
priority int // 1-5
|
||||
tags []string // ["loudspeaker"]
|
||||
struct DedupConfig {
|
||||
enabled bool
|
||||
ttl_seconds int // TTL du cache de déduplication
|
||||
rate_limit RateLimitConfig // max_per_minute, max_per_source
|
||||
}
|
||||
|
||||
struct Source { // types spécifiques dans sources/event.v
|
||||
name string
|
||||
webhook_path string
|
||||
topic string
|
||||
template string // Go-like template
|
||||
priority_map map[string]int
|
||||
tags []string
|
||||
}
|
||||
|
||||
struct Event { // dans sources/event.v
|
||||
source string // "gitea", "docker", "uptime_kuma", "http_poll", "cron"
|
||||
name string // nom de la source configurée
|
||||
topic string // topic Ntfy cible
|
||||
priority int // 1-5
|
||||
tags []string
|
||||
message string // message formaté final
|
||||
raw string // payload brut (pour debug)
|
||||
// sources/event.v
|
||||
struct Event {
|
||||
source string // "gitea", "docker", "uptime_kuma", "http_poll", "cron", "generic"
|
||||
name string // nom de la source configurée
|
||||
topic string // topic Ntfy cible
|
||||
priority int // 1-5
|
||||
tags []string
|
||||
message string // message formaté final
|
||||
raw string // payload brut (pour debug)
|
||||
}
|
||||
```
|
||||
|
||||
Les types de sources (GiteaSource, UptimeKumaSource, DockerSource, HttpPollSource/HttpCheck, CronSource, GenericSource) sont définis dans `sources/event.v`.
|
||||
|
||||
## Pipeline de transformation
|
||||
|
||||
Chaque source implémente l'interface :
|
||||
Chaque source a une **free function** dans `sources/` :
|
||||
|
||||
```v
|
||||
interface Transformer {
|
||||
transform(raw string, source &Source) ?Event
|
||||
}
|
||||
// Exemple : sources/gitea.v
|
||||
pub fn transform_gitea(raw string, source GiteaSource) ?Event { ... }
|
||||
```
|
||||
|
||||
Pas d'interface — les fonctions sont dispatchées via un `match route.kind` dans `server.v`. L'ajout d'une nouvelle source nécessite :
|
||||
1. Un fichier `sources/nouveau.v` avec la `transform_*` function
|
||||
2. Un type source dans `sources/event.v`
|
||||
3. Une entrée dans `SourcesConfig` (dans `config.v`)
|
||||
4. Une branche dans `build_webhook_routes` et `handle_webhook`
|
||||
|
||||
### Transformers inclus
|
||||
|
||||
| Transformer | Entrée | Sortie |
|
||||
|---|---|---|
|
||||
| `gitea_transformer` | JSON webhook Gitea | `🔨 [FlowDeck] bruno pushed to main: "Fix bug" (a3f2c1d)` |
|
||||
| `uptime_kuma_transformer` | JSON webhook Kuma | `🚨 Gitea is DOWN — 503 — since 14:32` |
|
||||
| `docker_transformer` | Docker event JSON | `🐳 flowdeck exited OOMKilled on docker-prod-1` |
|
||||
| `http_poll_transformer` | HTTP response | `❌ obsigate.dracodev.net/health → timeout 30s` |
|
||||
| `generic_transformer` | Raw body | Pass-through, pas de transformation |
|
||||
| `transform_gitea` | JSON webhook Gitea | `🔨 [FlowDeck] bruno pushed: "Fix bug" (a3f2c1d)` |
|
||||
| `transform_uptime_kuma` | JSON webhook Kuma | `🚨 Gitea is DOWN — 503 — since 14:32` |
|
||||
| `transform_docker_event` | Docker event JSON | `🐳 flowdeck exited on docker-prod-1` |
|
||||
| `transform_http_down/up` | HTTP response | `❌ og.dracodev.net/health → timeout` |
|
||||
| `transform_cron` | Raw body | Pass-through, pas de transformation |
|
||||
| `transform_generic` | Raw body | Pass-through, pas de transformation |
|
||||
|
||||
### Template engine
|
||||
|
||||
`render_source_template(template, vars)` dans `sources/event.v` — remplacement simple `{variable}` → valeur. Supporté par Gitea, Uptime Kuma, Docker, HTTP poll.
|
||||
|
||||
```
|
||||
template: "🐳 {container_name} → {status}\nImage: {image}\nHost: {host}"
|
||||
```
|
||||
|
||||
### Priorités Ntfy
|
||||
|
||||
@@ -157,64 +167,80 @@ interface Transformer {
|
||||
| 4 (high) | PR ouverte, container crash | Gitea pull_request, Docker die |
|
||||
| 3 (default) | Push, issue, activité normale | Gitea push |
|
||||
| 2 (low) | Info, succès | Service back UP |
|
||||
| 1 (min) | Debug, heartbeat | Health check OK (si configuré) |
|
||||
| 1 (min) | Debug, heartbeat | Health check OK |
|
||||
|
||||
## Déduplication et rate limiting
|
||||
|
||||
- **Déduplication** : même source + même type d'événement + même cible → cooldown 5 minutes
|
||||
- **Rate limiting** : max 10 notifs/minute/topic (configurable)
|
||||
- **Uptime Kuma spécifique** : les "down" répétés dans un intervalle de 2 min sont dédupliqués
|
||||
- **Déduplication** : hash du message (FNV-1a-like) → cache avec TTL configurable (défaut: 5 min). Clé = `source:name:topic:hash(message)`. Un message identique dans la fenêtre TTL est ignoré.
|
||||
- **Rate limiting par topic** : `max_per_minute` notifications max par topic (défaut: 10). Fenêtre glissante de 60s.
|
||||
- **Rate limiting par source** : `max_per_source` notifications max par source configurée (défaut: 30). Fenêtre glissante de 60s.
|
||||
- **Nettoyage automatique** : les entrées expirées sont nettoyées à chaque appel `allow()`.
|
||||
|
||||
## Structure du projet
|
||||
|
||||
```
|
||||
ntfy-bridge/
|
||||
├── README.md # Vue d'ensemble, quickstart
|
||||
├── ARCHITECTURE.md # Ce document
|
||||
├── ROADMAP.md # Phases de développement
|
||||
├── CONTRIBUTING.md # Guide de contribution
|
||||
├── docs/
|
||||
│ ├── ARCHITECTURE.md # Ce document
|
||||
│ ├── ROADMAP.md # Phases de développement
|
||||
│ ├── CONTRIBUTING.md # Guide de contribution
|
||||
│ └── WEBHOOK_GITEA.md # Guide config webhook Gitea
|
||||
├── ntfy-bridge.example.yaml # Exemple de configuration complet
|
||||
├── main.v # Entrypoint, parsing CLI
|
||||
├── config.v # Chargement + validation YAML
|
||||
├── server.v # HTTP server + routing webhooks
|
||||
├── ntfy.v # Client HTTP Ntfy (POST)
|
||||
├── dedup.v # Déduplication + rate limiting
|
||||
├── template.v # Mini moteur de template
|
||||
├── log.v # Logging structuré
|
||||
├── docker_watcher.v # Watcher socket Docker (Unix)
|
||||
├── main.v # Entrypoint, parsing CLI, header ASCII
|
||||
├── config.v # Chargement + validation YAML, env overrides
|
||||
├── server.v # HTTP server, routing webhooks O(1), dashboard, API
|
||||
├── ntfy.v # Client HTTP Ntfy (POST + retry backoff)
|
||||
├── dedup.v # Déduplication + rate limiting (topic + source)
|
||||
├── template.v # Application des defaults (priority, tags)
|
||||
├── log.v # Logging structuré (human + JSON)
|
||||
├── webhook.v # O(1) route dispatch (WebhookRoute map)
|
||||
├── docker_watcher.v # Watcher socket Docker (Unix seulement)
|
||||
├── dashboard.html # Interface web : stats, status, historique
|
||||
├── sources/ # Types Event + transformers
|
||||
│ ├── event.v # Types Event, Source
|
||||
│ ├── gitea.v # Transformer Gitea
|
||||
│ ├── uptime_kuma.v # Transformer Uptime Kuma
|
||||
│ ├── event.v # Types Event, Source, render_source_template
|
||||
│ ├── gitea.v # Transformer Gitea (push, PR, issue, release)
|
||||
│ ├── uptime_kuma.v # Transformer Uptime Kuma (state_map)
|
||||
│ ├── docker.v # Transformer Docker events
|
||||
│ ├── http_poll.v # Transformer HTTP poll
|
||||
│ ├── generic.v # Webhook passe-partout
|
||||
│ └── cron.v # Webhook cron
|
||||
├── *_test.v # Tests unitaires
|
||||
│ ├── http_poll.v # Transformer HTTP poll (up/down)
|
||||
│ ├── cron.v # Transformer Cron (pass-through)
|
||||
│ └── generic.v # Transformer Generic (pass-through)
|
||||
├── *_test.v # Tests unitaires (config, gitea, uptime_kuma, dedup)
|
||||
├── scripts/
|
||||
│ └── example-cron-disk.sh # Exemple script cron → Ntfy
|
||||
└── v.mod # Dépendances V
|
||||
└── v.mod # Dépendances V (aucune dépendance externe)
|
||||
```
|
||||
|
||||
## Dépendances V (v.mod)
|
||||
## Dépendances
|
||||
|
||||
```v
|
||||
// v.mod
|
||||
Module {
|
||||
name: 'ntfy-bridge'
|
||||
description: 'Homelab notification hub — aggregates Gitea, Docker, Uptime Kuma → Ntfy'
|
||||
version: '0.1.0'
|
||||
deps: [
|
||||
'ui-driver.cua'
|
||||
]
|
||||
dependencies: [] // zéro dépendance externe
|
||||
}
|
||||
```
|
||||
|
||||
Dépendances externes minimales :
|
||||
- `vsl` — parsing YAML (sinon implémenter un parseur minimal)
|
||||
- `x.net` — client HTTP (ou lib standard V `net`)
|
||||
- `json` — stdlib V pour le parsing JSON
|
||||
Le projet utilise **uniquement la stdlib V** :
|
||||
- `net.http` — serveur HTTP + client Ntfy
|
||||
- `x.json2` — parsing JSON
|
||||
- `yaml` — parsing YAML
|
||||
- `crypto.hmac`, `crypto.sha256`, `encoding.hex` — validation HMAC
|
||||
- `time`, `os`, `flag` — utilitaires
|
||||
- `net.unix` — socket Docker (Linux/macOS seulement, via `$if`)
|
||||
|
||||
En réalité, V a tout dans sa stdlib : `net.http`, `json`, `time`, `os`, `flag`. Le projet vise **zéro dépendance externe**.
|
||||
## Endpoints HTTP
|
||||
|
||||
| Méthode | Path | Description |
|
||||
|---------|------|-------------|
|
||||
| GET | `/` ou `/dashboard` | Dashboard HTML |
|
||||
| GET | `/health` | Health check → `{"status":"ok","uptime_seconds":N,"total_notifications":N}` |
|
||||
| GET | `/stats` | Stats par source → `{"total":N,"errors":N,"by_source":{...}}` |
|
||||
| GET | `/api/config` | Config résumée (safe, sans secrets) |
|
||||
| GET | `/api/status` | État des health checks HTTP poll |
|
||||
| GET | `/api/history` | 100 derniers événements |
|
||||
| POST | `/webhooks/*` | Webhooks (Gitea, Uptime Kuma, Cron, Generic) |
|
||||
|
||||
## Déploiement
|
||||
|
||||
@@ -226,7 +252,7 @@ Description=ntfy-bridge notification hub
|
||||
After=network.target
|
||||
|
||||
[Service]
|
||||
ExecStart=/usr/local/bin/ntfy-bridge --config /etc/ntfy-bridge.yaml
|
||||
ExecStart=/usr/local/bin/ntfy-bridge --quiet --config /etc/ntfy-bridge.yaml
|
||||
Restart=always
|
||||
RestartSec=10
|
||||
|
||||
@@ -238,29 +264,32 @@ WantedBy=multi-user.target
|
||||
```dockerfile
|
||||
FROM alpine:latest
|
||||
COPY ntfy-bridge /usr/local/bin/
|
||||
COPY ntfy-bridge.yaml /etc/
|
||||
COPY ntfy-bridge.yaml dashboard.html /etc/ntfy-bridge/
|
||||
WORKDIR /etc/ntfy-bridge
|
||||
EXPOSE 9090
|
||||
CMD ["/usr/local/bin/ntfy-bridge", "--config", "/etc/ntfy-bridge.yaml"]
|
||||
CMD ["/usr/local/bin/ntfy-bridge", "--quiet", "--config", "/etc/ntfy-bridge/ntfy-bridge.yaml"]
|
||||
```
|
||||
|
||||
### Option C : Compilation croisée
|
||||
```bash
|
||||
v -os windows . -o ntfy-bridge.exe # Pour Windows
|
||||
v -os linux . -o ntfy-bridge # Pour Linux
|
||||
v -os macos . -o ntfy-bridge-mac # Pour macOS
|
||||
v -prod -os windows . -o ntfy-bridge.exe
|
||||
v -prod -os linux . -o ntfy-bridge
|
||||
```
|
||||
|
||||
## Sécurité
|
||||
|
||||
- **Signature webhooks** : vérification HMAC-SHA256 optionnelle pour Gitea
|
||||
- **Auth token** : token optionnel pour le serveur Ntfy (si auth activée)
|
||||
- **Pas d'exposition externe** : le serveur HTTP écoute sur localhost par défaut (:9090)
|
||||
- **Rate limiting interne** : max 10 notifs/minute/source pour éviter le spam
|
||||
- **Input validation** : tout JSON entrant est validé avant processing
|
||||
- **No secrets in config** : le token Ntfy peut être passé via variable d'environnement `NTFY_TOKEN`
|
||||
- **Signature webhooks** : vérification HMAC-SHA256 pour **tous** les webhooks (header `X-Hub-Signature-256: sha256=...`). Activé dès que `server.hmac_secret` est défini.
|
||||
- **Auth token** : token Bearer optionnel pour le serveur Ntfy (si auth activée sur le serveur Ntfy)
|
||||
- **Pas d'exposition externe** : le serveur HTTP écoute sur localhost par défaut (`127.0.0.1:9090`)
|
||||
- **Rate limiting interne** : par topic ET par source (configurable dans `dedup.rate_limit`)
|
||||
- **Input validation** : tout JSON entrant est validé avant processing, chemins webhook dupliqués détectés au démarrage
|
||||
- **No secrets in config** : `NTFY_TOKEN` et `NTFY_HMAC_SECRET` peuvent être passés via variables d'environnement
|
||||
|
||||
## Métriques et observabilité
|
||||
|
||||
- **Health endpoint** : `GET /health` → 200 OK + uptime + compteur de notifs
|
||||
- **Stats endpoint** : `GET /stats` → compteurs par source, erreurs, latence Ntfy
|
||||
- **Logging** : stdout en JSON structuré (timestamp, level, source, event, error)
|
||||
- **Stats endpoint** : `GET /stats` → compteurs par source + erreurs
|
||||
- **API status** : `GET /api/status` → état des health checks HTTP poll
|
||||
- **API history** : `GET /api/history` → 100 derniers événements avec timestamps
|
||||
- **Dashboard** : interface HTML avec stats temps réel, statut des services, historique récent
|
||||
- **Logging** : stdout en mode human-readable (console) ou JSON structuré (`--quiet`, pour systemd)
|
||||
|
||||
+105
-50
@@ -1,76 +1,115 @@
|
||||
# ntfy-bridge — Roadmap
|
||||
|
||||
## Phase 1 : MVP (core engine) — v0.1.0
|
||||
## Phase 1 : MVP (core engine) — v0.1.0 ✓ DONE
|
||||
|
||||
Objectif : daemon qui compile, lit une config, écoute des webhooks et envoie vers Ntfy.
|
||||
|
||||
- [ ] `main.v` — CLI args (`--config`, `--version`, `--validate`)
|
||||
- [ ] `config.v` — Chargement + validation du YAML
|
||||
- [ ] `server.v` — HTTP server sur `:9090` avec routing basique
|
||||
- [ ] `event.v` — Types de données : `Event`, `Config`, `Source`
|
||||
- [ ] `ntfy.v` — Client HTTP POST vers Ntfy (topic, message, priority, tags)
|
||||
- [ ] `sources/generic.v` — Webhook passe-partout (POST raw → Ntfy)
|
||||
- [ ] `log.v` — Logging structuré (timestamp, level, message)
|
||||
- [ ] `ntfy-bridge.example.yaml` — Config d'exemple complète et commentée
|
||||
- [x] `main.v` — CLI args (`--config`, `--version`, `--validate`, `--quiet`, `--help`)
|
||||
- [x] `config.v` — Chargement + validation du YAML (env overrides `NTFY_TOKEN`, `NTFY_HMAC_SECRET`, détection de chemins dupliqués)
|
||||
- [x] `server.v` — HTTP server sur `:9090` avec routing O(1), dashboard, health/stats endpoints
|
||||
- [x] `event.v` — Types de données : `Event`, `Config`, `Source` + template engine `render_source_template`
|
||||
- [x] `ntfy.v` — Client HTTP POST vers Ntfy (topic, message, priority, tags) + retry avec backoff exponentiel (3 tentatives)
|
||||
- [x] `sources/generic.v` — Webhook passe-partout (POST raw → Ntfy)
|
||||
- [x] `log.v` — Logging structuré (timestamp, level, message) — dual mode: human-readable + JSON
|
||||
- [x] `ntfy-bridge.example.yaml` — Config d'exemple complète et commentée
|
||||
|
||||
**Sortie** : `ntfy-bridge` compile, accepte un webhook générique et l'envoie vers Ntfy.
|
||||
**Sortie** : `ntfy-bridge` compile, accepte un webhook générique et l'envoie vers Ntfy. ✅
|
||||
|
||||
## Phase 2 : Sources spécialisées — v0.2.0
|
||||
## Phase 2 : Sources spécialisées — v0.2.0 ✓ DONE
|
||||
|
||||
- [ ] `sources/gitea.v` — Transformer Gitea (push, PR, issues, releases)
|
||||
- [x] `sources/gitea.v` — Transformer Gitea (push, PR, issues, releases)
|
||||
- Parse le JSON webhook Gitea
|
||||
- Format : `🔨 [repo] auteur action: "message" (sha)`
|
||||
- Template configurable avec variables `{repo}`, `{user}`, `{action}`, `{title}`, `{sha}`
|
||||
- Priority map par type d'événement
|
||||
- [ ] `sources/uptime_kuma.v` — Transformer Uptime Kuma
|
||||
- [x] `sources/uptime_kuma.v` — Transformer Uptime Kuma
|
||||
- Parse le JSON webhook Kuma (heartbeat)
|
||||
- Format : `🚨 monitor is DOWN — status — since time`
|
||||
- Priority 5 pour DOWN, 1 pour UP
|
||||
- [ ] `sources/docker.v` — Watcher socket Docker
|
||||
- Connexion au socket Docker (local + distant TCP)
|
||||
- `state_map` pour configurer priority/tags par état (up/down)
|
||||
- Template configurable avec variables `{monitor}`, `{status}`, `{msg}`, `{ping}`, `{time}`
|
||||
- [x] `sources/docker.v` — Watcher socket Docker
|
||||
- Connexion au socket Docker Unix (`docker_watcher.v`)
|
||||
- Filtrage par événements (die, health_status, oom)
|
||||
- Template de message configurable
|
||||
- Multi-hôtes (plusieurs sockets dans la config)
|
||||
- [ ] `sources/http_poll.v` — Poller HTTP
|
||||
- [x] `sources/http_poll.v` — Poller HTTP
|
||||
- Goroutine de polling périodique (intervalle configurable)
|
||||
- Vérification status code + timeout
|
||||
- State machine up/down avec notification uniquement au changement
|
||||
- [ ] `dedup.v` — Déduplication basique
|
||||
- Cache LRU avec TTL (5 min par défaut)
|
||||
- Clé de déduplication : source + type + cible
|
||||
- Template configurable avec variables `{url}`, `{error}`, `{status}`
|
||||
- [x] `dedup.v` — Déduplication + rate limiting
|
||||
- Cache avec TTL (5 min par défaut) basé sur hash du message
|
||||
- Clé de déduplication : source + type + topic + hash(message)
|
||||
- Rate limiting par topic et par source
|
||||
|
||||
**Sortie** : Toutes les sources majeures sont intégrées et fonctionnelles.
|
||||
**Bonus (non planifié) :**
|
||||
- [x] `sources/cron.v` — Webhook pour scripts cron (pass-through)
|
||||
- [x] `webhook.v` — Route dispatch O(1) via `map[string]WebhookRoute`
|
||||
|
||||
## Phase 3 : Robustesse — v0.3.0
|
||||
**Sortie** : Toutes les sources majeures sont intégrées et fonctionnelles. ✅
|
||||
|
||||
- [ ] `dedup.v` amélioré — Rate limiting global par topic
|
||||
- [ ] Signature webhooks — Vérification HMAC-SHA256 optionnelle (Gitea)
|
||||
- [ ] Retry logic — Retry avec backoff exponentiel vers Ntfy (3 tentatives)
|
||||
- [ ] Graceful shutdown — SIGTERM → drain des événements en cours
|
||||
- [ ] Health endpoint — `GET /health` → 200 + uptime + stats basiques
|
||||
- [ ] Stats endpoint — `GET /stats` → compteurs par source
|
||||
- [ ] Config reload — `SIGHUP` → reload config sans redémarrage
|
||||
- [ ] Tests unitaires — `config_test.v`, `gitea_test.v`, `dedup_test.v`
|
||||
## Phase 3 : Robustesse — v0.3.0 ✅ DONE
|
||||
|
||||
**Sortie** : Le daemon est prêt pour la production homelab.
|
||||
- [x] ~~`dedup.v` amélioré~~ — Rate limiting déjà intégré dans la phase 2
|
||||
- [x] Signature webhooks — Vérification HMAC-SHA256 (appliquée à **tous** les webhooks, pas seulement Gitea)
|
||||
- [x] Retry logic — Retry avec backoff exponentiel vers Ntfy (3 tentatives, 1s/2s/4s)
|
||||
- [x] Graceful shutdown — SIGTERM/SIGINT → drain des événements en cours (Linux/macOS)
|
||||
- [x] Health endpoint — `GET /health` → 200 + uptime + stats basiques
|
||||
- [x] Stats endpoint — `GET /stats` → compteurs par source
|
||||
- [x] Config reload — `SIGHUP` → reload config sans redémarrage
|
||||
- [x] Tests unitaires — `config_test.v`, `gitea_test.v`, `uptime_kuma_test.v`, `dedup_test.v`
|
||||
|
||||
## Phase 4 : Qualité de vie — v0.4.0
|
||||
**Sortie** : Le daemon est prêt pour la production homelab. ✅
|
||||
|
||||
- [ ] Template engine avancé — Support des variables dans les messages
|
||||
## Phase 4 : Qualité de vie — v0.4.0 ✅ DONE
|
||||
|
||||
- [x] Template engine — Support des variables `{key}` dans les messages (toutes les sources)
|
||||
```
|
||||
"🐳 {container_name} → {status} (exit: {exit_code}) on {host}"
|
||||
```
|
||||
- [ ] Silence rules — `POST /api/silence?duration=30m` mute temporaire
|
||||
- [ ] Grouping — Regrouper N notifications similaires en une seule
|
||||
- [ ] Webhook secret validation — HMAC pour tous les webhooks
|
||||
- [ ] Docker Compose example — `docker-compose.yml` prêt à l'emploi
|
||||
- [ ] systemd unit file — `ntfy-bridge.service`
|
||||
- [ ] CI/CD via Gitea Actions — build + test automatique
|
||||
- [x] Silence rules — `POST /api/silence?duration=30m` + `GET /api/silence` + `DELETE /api/silence`
|
||||
- [x] Grouping — Regrouper N notifications similaires (buffer 10s, flush avec ×N)
|
||||
- [x] Webhook secret validation — HMAC global (tous les webhooks, pas seulement Gitea)
|
||||
- [x] CI/CD via Gitea Actions — `.gitea/workflows/ci.yml`
|
||||
|
||||
**Sortie** : Expérience utilisateur complète et agréable.
|
||||
**Sortie** : Expérience utilisateur complète. ✅
|
||||
|
||||
## Phase 4.5 : Déploiement multi-plateforme — v0.4.5 🔜
|
||||
|
||||
Objectif : un seul binaire, une seule commande pour installer/désinstaller le service sur tous les OS supportés.
|
||||
|
||||
- [ ] `service.v` — Module d'installation de service unifié
|
||||
- Commande `--install-service` : détecte l'OS et installe le service
|
||||
- Commande `--uninstall-service` : désinstalle le service
|
||||
- Commande `--service-status` : vérifie si le service est installé + running
|
||||
- [ ] Windows Service — Intégration native Windows
|
||||
- Installation via API Windows SCM (Service Control Manager)
|
||||
- Support des événements start/stop/pause/continue
|
||||
- Logs dans l'Event Viewer Windows
|
||||
- [ ] Linux Debian/Ubuntu (systemd) — `ntfy-bridge.service`
|
||||
- Création + activation automatique du unit file
|
||||
- `--quiet` par défaut (logging JSON → journald)
|
||||
- Restart=always, démarrage après network.target
|
||||
- [ ] Linux Alpine (openrc) — Script init.d
|
||||
- `/etc/init.d/ntfy-bridge` généré automatiquement
|
||||
- Commandes start/stop/restart/status
|
||||
- Ajout au runlevel default
|
||||
- [ ] Raspberry Pi (Debian aarch64) — Support confirmé
|
||||
- Cross-compilation `v -os linux -arch arm64`
|
||||
- Même unit systemd que Debian x86_64
|
||||
- Testé sur Raspberry Pi OS (arm64)
|
||||
- [ ] Docker image multi-arch — `docker-compose.yml`
|
||||
- Build multi-arch : `linux/amd64`, `linux/arm64`, `linux/arm/v7`
|
||||
- Image minimale basée sur `alpine:latest` (~15 MB)
|
||||
- docker-compose.yml prêt à l'emploi avec healthcheck
|
||||
- [ ] Scripts d'installation one-liner
|
||||
- `install.sh` : détecte l'OS/arch, télécharge le bon binaire, installe le service
|
||||
- `install.ps1` : équivalent Windows PowerShell
|
||||
- `curl -fsSL https://.../install.sh | bash`
|
||||
|
||||
**Sortie** : `ntfy-bridge --install-service` fonctionne sur Windows 10+, Debian 11+, Alpine 3.18+, Raspberry Pi OS (arm64). Un seul flux d'installation quelle que soit la plateforme.
|
||||
|
||||
## Phase 5 : Futures idées — v0.5.0+
|
||||
|
||||
- [ ] Web UI minimale — Page de statut des sources + historique récent
|
||||
- [x] Web UI minimale — Dashboard HTML avec stats, status, historique (+ API `/api/config`, `/api/status`, `/api/history`)
|
||||
- [ ] Filtres avancés — Expressions conditionnelles par source
|
||||
```yaml
|
||||
filters:
|
||||
@@ -84,12 +123,28 @@ Objectif : daemon qui compile, lit une config, écoute des webhooks et envoie ve
|
||||
- [ ] Fichier d'état — Persistance de l'état up/down des health checks
|
||||
- [ ] Plugin system — Sources customisables via dll/.so (futur lointain)
|
||||
|
||||
## Fonctionnalités non planifiées mais implémentées
|
||||
|
||||
Ces features ont été ajoutées en cours de route :
|
||||
|
||||
| Feature | Fichier | Description |
|
||||
|---------|---------|-------------|
|
||||
| Source Cron | `sources/cron.v` | Webhook pass-through pour scripts cron (curl → Ntfy) |
|
||||
| Dashboard web | `dashboard.html` | Interface avec stats, statut des health checks, historique récent |
|
||||
| Logging JSON | `log.v` | Mode JSON structuré activé avec `--quiet` (pour systemd) |
|
||||
| Env overrides | `config.v` | `NTFY_TOKEN` et `NTFY_HMAC_SECRET` depuis l'environnement |
|
||||
| Validation dupplicates | `config.v` | Détection de chemins webhook dupliqués dans la config |
|
||||
| Flag `--quiet` | `main.v` | Supprime la bannière ASCII au démarrage |
|
||||
| Flag `-V` (version) | `main.v` | Version courte en plus de `--version` |
|
||||
| API endpoints | `server.v` | `/api/config`, `/api/status`, `/api/history` pour le dashboard |
|
||||
|
||||
## Suivi des versions
|
||||
|
||||
| Version | Date cible | Contenu |
|
||||
|---------|-----------|---------|
|
||||
| v0.1.0 | — | Moteur core + webhook générique |
|
||||
| v0.2.0 | — | Sources Gitea, Uptime Kuma, Docker, HTTP poll |
|
||||
| v0.3.0 | — | Robustesse, tests, health check, rate limiting |
|
||||
| v0.4.0 | — | Templates, silence rules, grouping, CI/CD |
|
||||
| v0.5.0+ | — | Web UI, Prometheus, plugins |
|
||||
| Version | Statut | Contenu |
|
||||
|---------|--------|---------|
|
||||
| v0.1.0 | ✅ Terminé | Moteur core + webhook générique + logging |
|
||||
| v0.2.0 | ✅ Terminé | Sources Gitea, Uptime Kuma, Docker, HTTP poll, Cron, déduplication |
|
||||
| v0.3.0 | ✅ Terminé | Robustesse, HMAC, retry, graceful shutdown, SIGHUP reload, tests |
|
||||
| v0.4.0 | ✅ Terminé | Templates, silence, grouping, HMAC global, CI/CD |
|
||||
| v0.4.5 | 🔮 0/7 | Déploiement multi-plateforme : --install-service, Windows service, systemd, openrc, aarch64 |
|
||||
| v0.5.0+ | 🔮 1/8 | Dashboard web — reste filtres, actions Ntfy, Prometheus, etc. |
|
||||
|
||||
@@ -3,7 +3,7 @@ module main
|
||||
import os
|
||||
import flag
|
||||
|
||||
const version = '0.1.0'
|
||||
const version = '0.4.0'
|
||||
const author = 'Bruno Charest'
|
||||
|
||||
fn main() {
|
||||
@@ -52,7 +52,7 @@ fn main() {
|
||||
print_startup_info(cfg)
|
||||
}
|
||||
|
||||
mut app := new_app(cfg, quiet)
|
||||
mut app := new_app(cfg, config_path, quiet)
|
||||
app.run() or {
|
||||
eprintln('Server error: ${err}')
|
||||
exit(1)
|
||||
|
||||
@@ -8,6 +8,7 @@ import crypto.sha256
|
||||
import encoding.hex
|
||||
import sources
|
||||
import os
|
||||
import strconv
|
||||
|
||||
|
||||
pub struct Stats {
|
||||
@@ -19,17 +20,21 @@ pub mut:
|
||||
|
||||
pub struct App {
|
||||
pub mut:
|
||||
cfg Config
|
||||
logger Logger
|
||||
ntfy NtfyClient
|
||||
dedup Dedup
|
||||
stats Stats
|
||||
start_time time.Time
|
||||
poll_state map[string]bool // url -> is_down
|
||||
webhook_routes map[string]WebhookRoute
|
||||
recent_events []RecentEvent
|
||||
shutdown_flag bool
|
||||
dashboard_html string // cached dashboard HTML
|
||||
cfg Config
|
||||
config_path string // path to config file for SIGHUP reload
|
||||
logger Logger
|
||||
ntfy NtfyClient
|
||||
dedup Dedup
|
||||
stats Stats
|
||||
start_time time.Time
|
||||
poll_state map[string]bool // url -> is_down
|
||||
webhook_routes map[string]WebhookRoute
|
||||
recent_events []RecentEvent
|
||||
shutdown_flag bool
|
||||
dashboard_html string // cached dashboard HTML
|
||||
silence_until time.Time // silence all notifications until this time
|
||||
group_buffer map[string][]GroupEntry // key: source:name:topic
|
||||
group_last_flush time.Time
|
||||
}
|
||||
|
||||
pub struct RecentEvent {
|
||||
@@ -42,27 +47,38 @@ pub struct RecentEvent {
|
||||
success bool
|
||||
}
|
||||
|
||||
pub struct GroupEntry {
|
||||
event sources.Event
|
||||
mut:
|
||||
received time.Time
|
||||
count int
|
||||
}
|
||||
|
||||
const recent_events_max = 100
|
||||
|
||||
pub fn new_app(cfg Config, quiet bool) App {
|
||||
pub fn new_app(cfg Config, config_path string, quiet bool) App {
|
||||
// json_mode = quiet — when running as a systemd service (quiet mode),
|
||||
// use JSON structured logging; otherwise use human-readable console output
|
||||
logger := new_logger(.info, quiet)
|
||||
dashboard := load_dashboard_html(logger)
|
||||
return App{
|
||||
cfg: cfg
|
||||
logger: logger
|
||||
ntfy: new_ntfy_client(cfg.server, logger)
|
||||
dedup: new_dedup(cfg.dedup)
|
||||
stats: Stats{
|
||||
cfg: cfg
|
||||
config_path: config_path
|
||||
logger: logger
|
||||
ntfy: new_ntfy_client(cfg.server, logger)
|
||||
dedup: new_dedup(cfg.dedup)
|
||||
stats: Stats{
|
||||
by_source: map[string]int{}
|
||||
}
|
||||
start_time: time.now()
|
||||
poll_state: map[string]bool{}
|
||||
webhook_routes: build_webhook_routes(cfg)
|
||||
recent_events: []RecentEvent{}
|
||||
shutdown_flag: false
|
||||
dashboard_html: dashboard
|
||||
start_time: time.now()
|
||||
poll_state: map[string]bool{}
|
||||
webhook_routes: build_webhook_routes(cfg)
|
||||
recent_events: []RecentEvent{}
|
||||
shutdown_flag: false
|
||||
dashboard_html: dashboard
|
||||
silence_until: time.Time{}
|
||||
group_buffer: map[string][]GroupEntry{}
|
||||
group_last_flush: time.now()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -104,12 +120,30 @@ pub fn (mut app App) run() ! {
|
||||
|
||||
$if linux || macos {
|
||||
fn (mut app App) handle_signals() {
|
||||
sigch := os.new_signal_channel(.int, .term)
|
||||
_ := <-sigch
|
||||
app.logger.info('server', 'shutdown', 'received signal, shutting down...')
|
||||
app.shutdown_flag = true
|
||||
time.sleep(500 * time.millisecond)
|
||||
exit(0)
|
||||
sigch := os.new_signal_channel(.int, .term, .hup)
|
||||
for {
|
||||
sig := <-sigch
|
||||
match sig {
|
||||
.hup {
|
||||
app.logger.info('server', 'reload', 'SIGHUP received, reloading config...')
|
||||
new_cfg := load_config(app.config_path) or {
|
||||
app.logger.error('server', 'reload', 'failed to reload config: ${err}')
|
||||
continue
|
||||
}
|
||||
app.cfg = new_cfg
|
||||
app.webhook_routes = build_webhook_routes(new_cfg)
|
||||
app.ntfy = new_ntfy_client(new_cfg.server, app.logger)
|
||||
app.dedup = new_dedup(new_cfg.dedup)
|
||||
app.logger.info('server', 'reload', 'config reloaded successfully')
|
||||
}
|
||||
.int, .term {
|
||||
app.logger.info('server', 'shutdown', 'received signal, shutting down...')
|
||||
app.shutdown_flag = true
|
||||
time.sleep(500 * time.millisecond)
|
||||
exit(0)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -120,6 +154,8 @@ fn (mut app App) start_background_tasks() {
|
||||
if app.cfg.sources.docker.len > 0 {
|
||||
spawn app.docker_watch_loop()
|
||||
}
|
||||
// Grouping flush goroutine
|
||||
spawn app.group_flush_loop()
|
||||
}
|
||||
|
||||
// ---- HTTP Handler ----
|
||||
@@ -145,10 +181,21 @@ pub fn (mut app App) handle(req http.Request) http.Response {
|
||||
if req.url == '/api/history' {
|
||||
return app.history_api_response()
|
||||
}
|
||||
if req.url == '/api/silence' {
|
||||
return app.silence_status_response()
|
||||
}
|
||||
}
|
||||
.post {
|
||||
if req.url.starts_with('/api/silence') {
|
||||
return app.silence_create_response(req.url)
|
||||
}
|
||||
return app.handle_webhook(req)
|
||||
}
|
||||
.delete {
|
||||
if req.url == '/api/silence' {
|
||||
return app.silence_clear_response()
|
||||
}
|
||||
}
|
||||
else {}
|
||||
}
|
||||
|
||||
@@ -286,24 +333,44 @@ fn (mut app App) verify_hmac(payload string, sig_header string) bool {
|
||||
|
||||
pub fn (mut app App) publish(event sources.Event) {
|
||||
mut e := apply_defaults(event, app.cfg.defaults)
|
||||
|
||||
// Check silence
|
||||
now := time.now()
|
||||
if now < app.silence_until {
|
||||
app.logger.debug('server', 'publish', 'silenced: ${e.source}:${e.name}')
|
||||
return
|
||||
}
|
||||
|
||||
if !app.dedup.allow(e) {
|
||||
app.logger.debug('server', 'dedup', 'dropped event ${e.source}:${e.name}')
|
||||
return
|
||||
}
|
||||
|
||||
now := time.now()
|
||||
app.ntfy.publish(e) or {
|
||||
app.logger.error('server', 'publish', 'failed to publish: ${err}')
|
||||
app.stats.errors++
|
||||
app.add_recent_event(now, e, false)
|
||||
return
|
||||
// Grouping: buffer the event
|
||||
key := '${e.source}:${e.name}:${e.topic}'
|
||||
mut entries := app.group_buffer[key] or { []GroupEntry{} }
|
||||
// Check if there's a matching event in the buffer (same source/name/topic)
|
||||
mut found_idx := -1
|
||||
for i, entry in entries {
|
||||
if entry.event.message == e.message {
|
||||
found_idx = i
|
||||
break
|
||||
}
|
||||
}
|
||||
app.stats.total_notifications++
|
||||
mut count := app.stats.by_source[e.source] or { 0 }
|
||||
count++
|
||||
app.stats.by_source[e.source] = count
|
||||
app.logger.info('server', 'publish', '${e.source}:${e.name} → ${e.topic}')
|
||||
app.add_recent_event(now, e, true)
|
||||
if found_idx >= 0 {
|
||||
// Rebuild to mutate count/received
|
||||
mut updated := entries[found_idx]
|
||||
updated.count++
|
||||
updated.received = now
|
||||
entries[found_idx] = updated
|
||||
} else {
|
||||
entries << GroupEntry{
|
||||
event: e
|
||||
received: now
|
||||
count: 1
|
||||
}
|
||||
}
|
||||
app.group_buffer[key] = entries
|
||||
}
|
||||
|
||||
fn (mut app App) add_recent_event(ts time.Time, event sources.Event, success bool) {
|
||||
@@ -388,3 +455,128 @@ fn (mut app App) dashboard_response() http.Response {
|
||||
header: header
|
||||
})
|
||||
}
|
||||
|
||||
// ---- Silence Rules ----
|
||||
|
||||
fn (mut app App) silence_status_response() http.Response {
|
||||
is_silenced := time.now() < app.silence_until
|
||||
remaining := if is_silenced { app.silence_until.unix() - time.now().unix() } else { 0 }
|
||||
body := '{"silenced":${is_silenced},"remaining_seconds":${remaining}}'
|
||||
return new_json_response(.ok, body)
|
||||
}
|
||||
|
||||
fn (mut app App) silence_create_response(url string) http.Response {
|
||||
// Parse duration from query string: /api/silence?duration=30m
|
||||
duration_str := parse_query_param(url, 'duration')
|
||||
if duration_str == '' {
|
||||
return new_json_response(.bad_request, '{"error":"missing duration parameter (e.g. ?duration=30m)"}')
|
||||
}
|
||||
seconds := parse_duration(duration_str)
|
||||
if seconds <= 0 {
|
||||
return new_json_response(.bad_request, '{"error":"invalid duration: ${duration_str}"}')
|
||||
}
|
||||
app.silence_until = time.now().add_seconds(seconds)
|
||||
app.logger.info('server', 'silence', 'notifications silenced for ${duration_str} (${seconds}s)')
|
||||
body := '{"silenced":true,"duration":"${duration_str}","seconds":${seconds}}'
|
||||
return new_json_response(.ok, body)
|
||||
}
|
||||
|
||||
fn (mut app App) silence_clear_response() http.Response {
|
||||
app.silence_until = time.Time{}
|
||||
app.logger.info('server', 'silence', 'silence cleared')
|
||||
return new_json_response(.ok, '{"silenced":false}')
|
||||
}
|
||||
|
||||
// parse_query_param extracts a query parameter value from a URL string.
|
||||
fn parse_query_param(url string, key string) string {
|
||||
prefix := '/api/silence?'
|
||||
if !url.starts_with(prefix) {
|
||||
return ''
|
||||
}
|
||||
qs := url[prefix.len..]
|
||||
parts := qs.split('&')
|
||||
for part in parts {
|
||||
kv := part.split('=')
|
||||
if kv.len == 2 && kv[0] == key {
|
||||
return kv[1]
|
||||
}
|
||||
}
|
||||
return ''
|
||||
}
|
||||
|
||||
// parse_duration converts a human duration string (e.g. "30m", "2h", "1h30m") to seconds.
|
||||
fn parse_duration(s string) int {
|
||||
mut total := 0
|
||||
mut current := ''
|
||||
for c in s {
|
||||
if c >= `0` && c <= `9` {
|
||||
current += c.ascii_str()
|
||||
} else {
|
||||
val := strconv.atoi(current) or { 0 }
|
||||
match c {
|
||||
`s` { total += val }
|
||||
`m` { total += val * 60 }
|
||||
`h` { total += val * 3600 }
|
||||
`d` { total += val * 86400 }
|
||||
else {}
|
||||
}
|
||||
current = ''
|
||||
}
|
||||
}
|
||||
return total
|
||||
}
|
||||
|
||||
// ---- Grouping ----
|
||||
|
||||
const group_flush_interval = 10 // seconds
|
||||
|
||||
fn (mut app App) group_flush_loop() {
|
||||
for {
|
||||
if app.shutdown_flag {
|
||||
break
|
||||
}
|
||||
time.sleep(group_flush_interval * time.second)
|
||||
app.flush_group_buffer()
|
||||
}
|
||||
}
|
||||
|
||||
fn (mut app App) flush_group_buffer() {
|
||||
now := time.now()
|
||||
for _, entries in app.group_buffer {
|
||||
if entries.len == 0 {
|
||||
continue
|
||||
}
|
||||
mut total_count := 0
|
||||
mut latest_entry := entries[0]
|
||||
for i, entry in entries {
|
||||
total_count += entry.count
|
||||
if entry.received > latest_entry.received {
|
||||
latest_entry = entries[i]
|
||||
}
|
||||
}
|
||||
|
||||
mut msg := latest_entry.event.message
|
||||
if total_count > 1 {
|
||||
msg = '${msg} (×${total_count} in ${group_flush_interval}s)'
|
||||
}
|
||||
|
||||
mut e := latest_entry.event.clone()
|
||||
e.message = msg
|
||||
|
||||
app.ntfy.publish(e) or {
|
||||
app.logger.error('server', 'publish', 'failed to publish: ${err}')
|
||||
app.stats.errors++
|
||||
app.add_recent_event(now, e, false)
|
||||
continue
|
||||
}
|
||||
app.stats.total_notifications++
|
||||
mut count := app.stats.by_source[e.source] or { 0 }
|
||||
count++
|
||||
app.stats.by_source[e.source] = count
|
||||
app.logger.info('server', 'publish',
|
||||
'${e.source}:${e.name} → ${e.topic} (grouped ×${total_count})')
|
||||
app.add_recent_event(now, e, true)
|
||||
}
|
||||
app.group_buffer = map[string][]GroupEntry{}
|
||||
app.group_last_flush = time.now()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user