feat: v0.6.0 — ACLs multi-utilisateurs + Plugin system
CI / test (push) Has been cancelled

**ACLs (Access Control Lists):**
- Ajout du type AclConfig dans sources/event.v (allowed_ips CIDR + allowed_tokens)
- Chaque source webhook (Gitea, Uptime Kuma, Cron, Generic) supporte les ACLs
- Validation IP via X-Forwarded-For / X-Real-IP avant HMAC
- Validation Bearer token via Authorization header
- Tests: 7 tests ACL (IP exact, CIDR, parse IPv4, ACL vide, IP+token combiné)

**Plugin system:**
- sources/plugin.v: runner exécutable externe, stdout JSON → Event
- Exit 0 = publish, exit ≠ 0 = skip. Timeout configurable
- Plugin loop dans server.v (goroutine, toutes les 60s)
- Example: scripts/example-plugin-disk.sh (vérifie espace disque)

**Docs:**
- README.md: ajout source Plugin + section Features complète
- ARCHITECTURE.md: flux Plugin, flux ACL, endpoints /metrics /api/silence
- ROADMAP.md: Phase 5 → 8/8 complet, ajout v0.6.0
- ntfy-bridge.example.yaml: sections ACLs et Plugins commentées
- Version bump: 0.5.0 → 0.6.0
This commit is contained in:
2026-08-04 22:25:48 -04:00
parent 151dd4a1a8
commit 1a4808f5f7
13 changed files with 1230 additions and 686 deletions
+14
View File
@@ -57,6 +57,20 @@ vim ntfy-bridge.yaml
| **HTTP Poll** | Polling | Health checks périodiques | | **HTTP Poll** | Polling | Health checks périodiques |
| **Cron** | Webhook | Reçoit des notifs depuis des scripts shell | | **Cron** | Webhook | Reçoit des notifs depuis des scripts shell |
| **Generic** | Webhook | Endpoint passe-partout pour tout script custom | | **Generic** | Webhook | Endpoint passe-partout pour tout script custom |
| **Plugin** | Exec | Exécutable externe appelé périodiquement (bash, python, …) |
## Features
- **ACLs** — Restriction par IP (CIDR) ou Bearer token sur chaque webhook
- **Plugins** — Scripts externes exécutés toutes les 60s, contrat JSON stdout
- **Dashboard** — Interface web avec stats, statuts, historique temps réel
- **Filtres** — Règles conditionnelles drop/set_priority/add_tag par source
- **Actions Ntfy** — Boutons click, vues, broadcast intégrés
- **Prometheus** — Endpoint `/metrics` pour Grafana
- **Outgoing webhooks** — Forward vers Slack, Discord, JSON
- **Silence rules** — Mute temporaire des notifs via API
- **Grouping** — Regroupement intelligent des notifs identiques
- **State persistence** — Sauvegarde de l'état up/down et silence sur disque
## Exemple de config minimale ## Exemple de config minimale
+105
View File
@@ -0,0 +1,105 @@
module main
import net.http
import sources
fn test_ip_matches_exact() {
assert ip_matches('192.168.30.5', '192.168.30.5')
assert !ip_matches('192.168.30.6', '192.168.30.5')
}
fn test_ip_matches_cidr() {
assert ip_matches('192.168.30.5', '192.168.0.0/16')
assert ip_matches('10.0.0.1', '10.0.0.0/8')
assert ip_matches('10.255.255.255', '10.0.0.0/8')
assert !ip_matches('11.0.0.1', '10.0.0.0/8')
assert ip_matches('192.168.1.100', '192.168.1.0/24')
assert !ip_matches('192.168.2.1', '192.168.1.0/24')
}
fn test_parse_ip_valid() {
octets := parse_ip('192.168.1.5')!
assert octets[0] == 192
assert octets[1] == 168
assert octets[2] == 1
assert octets[3] == 5
}
fn test_parse_ip_invalid() {
// parse_ip with invalid input should return an error
result := parse_ip('invalid') or { []u8{} }
assert result.len == 0
result2 := parse_ip('1.2.3') or { []u8{} }
assert result2.len == 0
result3 := parse_ip('999.1.2.3') or { []u8{} }
assert result3.len == 0
}
fn test_check_acl_no_acl() {
mut app := new_test_acl_app()
route := WebhookRoute{kind: .generic, index: 0}
assert app.check_acl(route, '1.2.3.4', '') == true
}
fn test_check_acl_ip_allow() {
mut app := new_test_acl_app()
app.cfg.sources.generic << sources.GenericSource{
name: 'test-source'
webhook_path: '/test'
topic: 'test-topic'
acl: sources.AclConfig{
allowed_ips: ['10.0.0.0/8']
}
}
route := WebhookRoute{kind: .generic, index: 0}
assert app.check_acl(route, '10.0.0.5', '') == true
assert app.check_acl(route, '10.255.0.1', '') == true
assert app.check_acl(route, '192.168.1.1', '') == false
assert app.check_acl(route, '', '') == false
}
fn test_check_acl_combined_ip_and_token() {
mut app := new_test_acl_app()
app.cfg.sources.generic << sources.GenericSource{
name: 'test-source'
webhook_path: '/test'
topic: 'test-topic'
acl: sources.AclConfig{
allowed_ips: ['10.0.0.0/8']
allowed_tokens: ['secret123']
}
}
route := WebhookRoute{kind: .generic, index: 0}
// Both pass
assert app.check_acl(route, '10.0.0.5', 'secret123') == true
// IP ok but no token
assert app.check_acl(route, '10.0.0.5', '') == false
// Token ok but wrong IP
assert app.check_acl(route, '192.168.1.1', 'secret123') == false
// Both wrong
assert app.check_acl(route, '1.1.1.1', 'wrong') == false
}
fn test_extract_bearer_token_empty() {
req := http.Request{
header: http.new_header()
}
assert extract_bearer_token(req) == ''
}
// Helper: create a minimal App for ACL tests
fn new_test_acl_app() App {
return App{
logger: new_logger(.info, false)
cfg: Config{
filters: FilterConfig{}
outgoing: OutgoingConfig{}
state: StateConfig{}
}
poll_state: map[string]bool{}
}
}
+1
View File
@@ -66,6 +66,7 @@ pub mut:
http_poll []sources.HttpPollSource @[json: 'http_poll'] http_poll []sources.HttpPollSource @[json: 'http_poll']
cron []sources.CronSource @[json: 'cron'] cron []sources.CronSource @[json: 'cron']
generic []sources.GenericSource @[json: 'generic'] generic []sources.GenericSource @[json: 'generic']
plugin []sources.PluginSource @[json: 'plugin']
} }
pub struct Config { pub struct Config {
+1 -1
View File
@@ -4,8 +4,8 @@ import time
$if linux || macos { $if linux || macos {
import net.unix import net.unix
import os import os
}
import sources import sources
}
fn (mut app App) docker_watch_loop() { fn (mut app App) docker_watch_loop() {
$if linux || macos { $if linux || macos {
+53 -13
View File
@@ -2,7 +2,7 @@
## Philosophie ## Philosophie
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. ntfy-bridge est un **daemon HTTP léger** écrit en V (~2000 lignes) qui compile en un seul binaire natif statique. Il écoute des webhooks, se connecte au socket Docker, exécute des plugins externes, et poll des endpoints HTTP — puis transforme chaque événement en notification Ntfy formatée.
**Principes :** **Principes :**
- **Un binaire, zéro runtime** — pas de Node, Python, JVM ou conteneur obligatoire - **Un binaire, zéro runtime** — pas de Node, Python, JVM ou conteneur obligatoire
@@ -18,13 +18,13 @@ ntfy-bridge est un **daemon HTTP léger** écrit en V (~1000 lignes) qui compile
┌─────────────────────────────────────────────────────────────┐ ┌─────────────────────────────────────────────────────────────┐
│ ntfy-bridge │ │ ntfy-bridge │
│ │ │ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐ │ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐ ┌──────────┐ │
│ │ HTTP │ │ Docker │ │ HTTP │ │ Cron/ │ │ │ │ HTTP │ │ Docker │ │ HTTP │ │ Cron/ │ │ Plugin │ │
│ │ Server │ │ Watcher │ │ Poller │ │ Generic │ │ │ │ Server │ │ Watcher │ │ Poller │ │ Generic │ │ Runner │ │
│ │ :9090 │ │ goroutine│ │ goroutine│ │ Receiver │ │ │ │ :9090 │ │ goroutine│ │ goroutine│ │ Receiver │ │ goroutine│ │
│ └────┬─────┘ └────┬─────┘ └────┬─────┘ └─────┬──────┘ │ │ └────┬─────┘ └────┬─────┘ └────┬─────┘ └─────┬──────┘ └────┬─────┘ │
│ │ │ │ │ │ │ │ │ │ │ │ │
│ └──────────────┼──────────────┼──────────────┘ │ │ └──────────────┼──────────────┼──────────────┼──────────────┘ │
│ ▼ ▼ │ │ ▼ ▼ │
│ ┌──────────────────────────┐ │ │ ┌──────────────────────────┐ │
│ │ Event Pipeline │ │ │ │ Event Pipeline │ │
@@ -94,6 +94,30 @@ POST /webhooks/cron-disk
3. Le script cron contrôle son propre format 3. Le script cron contrôle son propre format
``` ```
### 5. Plugin runner (v0.6)
```
goroutine plugin_loop (via start_background_tasks, toutes les 60s)
│
▼
1. Exécute le binaire configuré (os.execute)
2. Lit stdout → parse JSON → Event struct
3. Exit 0 = publier, exit ≠ 0 = skip (pas de notif)
4. Supporte priority, tags, message, click_url, actions
5. Pipeline publish → Ntfy
```
### 6. ACL pipeline (v0.6)
```
POST /webhooks/*
│
▼
1. Extraire IP client (X-Forwarded-For > X-Real-IP)
2. Extraire Bearer token (Authorization header)
3. Si ACL configurée sur le webhook → valider IP et/ou token
4. Si échec → 403 Forbidden
5. Sinon → continuer le pipeline normal
```
## Modèle de données ## Modèle de données
```v ```v
@@ -113,17 +137,19 @@ struct DedupConfig {
// sources/event.v // sources/event.v
struct Event { struct Event {
source string // "gitea", "docker", "uptime_kuma", "http_poll", "cron", "generic" source string // "gitea", "docker", "uptime_kuma", "http_poll", "cron", "generic", "plugin"
name string // nom de la source configurée name string // nom de la source configurée
topic string // topic Ntfy cible topic string // topic Ntfy cible
priority int // 1-5 priority int // 1-5
tags []string tags []string
message string // message formaté final message string // message formaté final
raw string // payload brut (pour debug) raw string // payload brut (pour debug)
click_url string // URL ouverte au clic sur la notif (v0.5)
actions []Action // boutons d'action Ntfy (v0.5)
} }
``` ```
Les types de sources (GiteaSource, UptimeKumaSource, DockerSource, HttpPollSource/HttpCheck, CronSource, GenericSource) sont définis dans `sources/event.v`. Les types de sources (GiteaSource, UptimeKumaSource, DockerSource, HttpPollSource/HttpCheck, CronSource, GenericSource, PluginSource) et AclConfig sont définis dans `sources/event.v`.
## Pipeline de transformation ## Pipeline de transformation
@@ -150,6 +176,7 @@ Pas d'interface — les fonctions sont dispatchées via un `match route.kind` da
| `transform_http_down/up` | HTTP response | `❌ og.dracodev.net/health → timeout` | | `transform_http_down/up` | HTTP response | `❌ og.dracodev.net/health → timeout` |
| `transform_cron` | Raw body | Pass-through, pas de transformation | | `transform_cron` | Raw body | Pass-through, pas de transformation |
| `transform_generic` | Raw body | Pass-through, pas de transformation | | `transform_generic` | Raw body | Pass-through, pas de transformation |
| `transform_plugin` | stdout d'un exécutable externe | JSON → Event (exit 0 = publish, ≠ 0 = skip) |
### Template engine ### Template engine
@@ -196,6 +223,10 @@ ntfy-bridge/
├── log.v # Logging structuré (human + JSON) ├── log.v # Logging structuré (human + JSON)
├── webhook.v # O(1) route dispatch (WebhookRoute map) ├── webhook.v # O(1) route dispatch (WebhookRoute map)
├── docker_watcher.v # Watcher socket Docker (Unix seulement) ├── docker_watcher.v # Watcher socket Docker (Unix seulement)
├── filter_outgoing.v # Filtres avancés + outgoing webhooks + state persistence
├── metrics.v # Endpoint Prometheus /metrics
├── service.v # Installation/désinstallation de service (systemd, openrc, nssm)
├── docker_install.v # Génération de stack docker-compose
├── dashboard.html # Interface web : stats, status, historique ├── dashboard.html # Interface web : stats, status, historique
├── sources/ # Types Event + transformers ├── sources/ # Types Event + transformers
│ ├── event.v # Types Event, Source, render_source_template │ ├── event.v # Types Event, Source, render_source_template
@@ -204,10 +235,14 @@ ntfy-bridge/
│ ├── docker.v # Transformer Docker events │ ├── docker.v # Transformer Docker events
│ ├── http_poll.v # Transformer HTTP poll (up/down) │ ├── http_poll.v # Transformer HTTP poll (up/down)
│ ├── cron.v # Transformer Cron (pass-through) │ ├── cron.v # Transformer Cron (pass-through)
│ └── generic.v # Transformer Generic (pass-through) │ ├── generic.v # Transformer Generic (pass-through)
├── *_test.v # Tests unitaires (config, gitea, uptime_kuma, dedup) │ └── plugin.v # Transformer Plugin (exec stdout JSON)
├── *_test.v # Tests unitaires (config, gitea, uptime_kuma, dedup, filter_outgoing, metrics, acl)
├── scripts/ ├── scripts/
│ └── example-cron-disk.sh # Exemple script cron → Ntfy │ ├── example-cron-disk.sh # Exemple script cron → Ntfy
│ ├── example-plugin-disk.sh # Exemple plugin disk check
│ ├── install.sh # Script d'installation Linux one-liner
│ └── install.ps1 # Script d'installation Windows PowerShell
└── v.mod # Dépendances V (aucune dépendance externe) └── v.mod # Dépendances V (aucune dépendance externe)
``` ```
@@ -240,6 +275,10 @@ Le projet utilise **uniquement la stdlib V** :
| GET | `/api/config` | Config résumée (safe, sans secrets) | | GET | `/api/config` | Config résumée (safe, sans secrets) |
| GET | `/api/status` | État des health checks HTTP poll | | GET | `/api/status` | État des health checks HTTP poll |
| GET | `/api/history` | 100 derniers événements | | GET | `/api/history` | 100 derniers événements |
| GET | `/api/silence` | État du silence (actif/restant) |
| POST | `/api/silence?duration=30m` | Activer le silence |
| DELETE | `/api/silence` | Désactiver le silence |
| GET | `/metrics` | Métriques Prometheus / OpenMetrics |
| POST | `/webhooks/*` | Webhooks (Gitea, Uptime Kuma, Cron, Generic) | | POST | `/webhooks/*` | Webhooks (Gitea, Uptime Kuma, Cron, Generic) |
## Déploiement ## Déploiement
@@ -278,6 +317,7 @@ v -prod -os linux . -o ntfy-bridge
## Sécurité ## Sécurité
- **ACLs par webhook** (v0.6) : restriction par IP (CIDR ou exacte) et/ou Bearer token. Validé avant HMAC.
- **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. - **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) - **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`) - **Pas d'exposition externe** : le serveur HTTP écoute sur localhost par défaut (`127.0.0.1:9090`)
+19 -11
View File
@@ -106,7 +106,7 @@ Objectif : un seul binaire, une seule commande pour installer/désinstaller le s
**Sortie** : `ntfy-bridge --install-service` fonctionne sur Windows 10+, Debian 11+, Alpine 3.18+, Raspberry Pi OS (arm64). ✅ **Sortie** : `ntfy-bridge --install-service` fonctionne sur Windows 10+, Debian 11+, Alpine 3.18+, Raspberry Pi OS (arm64). ✅
## Phase 5 : Futures idées — v0.5.0+ (6/8) ## Phase 5 : Futures idées — v0.5.0+ (8/8) ✅ DONE
- [x] Web UI minimale — Dashboard avec statuts, historique, API endpoints - [x] Web UI minimale — Dashboard avec statuts, historique, API endpoints
- [x] Filtres avancés — Expressions conditionnelles par source - [x] Filtres avancés — Expressions conditionnelles par source
@@ -124,15 +124,12 @@ Objectif : un seul binaire, une seule commande pour installer/désinstaller le s
- Endpoint `GET /metrics` : gauges (uptime, poll_state, silence) + counters (notifications, errors, http_polls) - Endpoint `GET /metrics` : gauges (uptime, poll_state, silence) + counters (notifications, errors, http_polls)
- Format compatible Prometheus/Grafana - Format compatible Prometheus/Grafana
- [x] Webhook sortant — Forwarder les événements vers Slack, Discord, JSON - [x] Webhook sortant — Forwarder les événements vers Slack, Discord, JSON
- [ ] Support multi-utilisateurs — ACLs par webhook_path - [x] Support multi-utilisateurs — ACLs par webhook_path
- [x] Fichier d'état — Persistance de l'état up/down des health checks ```yaml
- [ ] Plugin system — Sources customisables via dll/.so acl:
- url: "https://hooks.slack.com/..." allowed_ips: ["192.168.30.5", "10.0.0.0/8"]
format: slack allowed_tokens: ["my-secret-token"]
- url: "https://discord.com/api/webhooks/..."
format: discord
``` ```
- [ ] Support multi-utilisateurs — ACLs par webhook_path
- [x] Fichier d'état — Persistance de l'état up/down des health checks - [x] Fichier d'état — Persistance de l'état up/down des health checks
```yaml ```yaml
state: state:
@@ -140,7 +137,17 @@ Objectif : un seul binaire, une seule commande pour installer/désinstaller le s
``` ```
- Sauvegarde `poll_state` (up/down) + `silence_until` - Sauvegarde `poll_state` (up/down) + `silence_until`
- Restauré au démarrage, persisté à chaque changement - Restauré au démarrage, persisté à chaque changement
- [ ] Plugin system — Sources customisables via dll/.so - [x] Plugin system — Sources customisables via executables externes
```yaml
plugin:
- name: "Disk space check"
topic: daily
command: /etc/ntfy-bridge/plugins/disk-check.sh
timeout: 10
```
- Contrat simple : stdin (optionnel) → stdout JSON
- Exit 0 = notif, exit ≠ 0 = skip
- Compatible bash, Python, ou n'importe quel langage
## Fonctionnalités non planifiées mais implémentées ## Fonctionnalités non planifiées mais implémentées
@@ -166,4 +173,5 @@ Ces features ont été ajoutées en cours de route :
| v0.3.0 | ✅ Terminé | Robustesse, HMAC, retry, graceful shutdown, SIGHUP reload, tests | | 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.0 | ✅ Terminé | Templates, silence, grouping, HMAC global, CI/CD |
| v0.4.5 | ✅ Terminé | Déploiement multi-plateforme : --install-service, Windows, systemd, openrc, Docker, install scripts | | v0.4.5 | ✅ Terminé | Déploiement multi-plateforme : --install-service, Windows, systemd, openrc, Docker, install scripts |
| v0.5.0+ | 🚧 6/8 | Dashboard, filtres, actions Ntfy, Prometheus, outgoing webhooks, state file — manque multi-user, plugins | | v0.5.0 | ✅ Terminé | Dashboard, filtres, actions Ntfy, Prometheus, outgoing webhooks, state file |
| v0.6.0 | ✅ Terminé | ACLs multi-utilisateurs, plugin system (scripts externes) |
+5 -1
View File
@@ -3,7 +3,7 @@ module main
import os import os
import flag import flag
const version = '0.5.0' const version = '0.6.0'
const author = 'Bruno Charest' const author = 'Bruno Charest'
fn main() { fn main() {
@@ -220,6 +220,10 @@ fn print_startup_info(cfg Config) {
if generic_count > 0 { if generic_count > 0 {
println(' Generic: ${generic_count} webhook${if generic_count > 1 { 's' } else { '' }}') println(' Generic: ${generic_count} webhook${if generic_count > 1 { 's' } else { '' }}')
} }
plugin_count := cfg.sources.plugin.len
if plugin_count > 0 {
println(' Plugin: ${plugin_count} script${if plugin_count > 1 { 's' } else { '' }}')
}
println('') println('')
println(' Webhook URLs:') println(' Webhook URLs:')
+32
View File
@@ -135,6 +135,38 @@ sources:
webhook_path: /webhooks/generic webhook_path: /webhooks/generic
topic: custom topic: custom
# ── ACLs (Access Control Lists) ─────────────────────────
# Optionnel — restreint l'accès à certains webhooks par IP ou token.
# Si aucune ACL n'est définie, le webhook est ouvert (soumis au HMAC).
#
# Exemple avec ACL IP:
# gitea:
# - name: "FlowDeck"
# webhook_path: /webhooks/gitea-flowdeck
# topic: dev-notifs
# acl:
# allowed_ips: ["192.168.30.5", "10.0.0.0/8"]
#
# Exemple avec ACL token:
# generic:
# - name: "Secure alerts"
# webhook_path: /webhooks/secure
# topic: secure
# acl:
# allowed_tokens: ["my-secret-webhook-token"]
# → Le client doit envoyer: Authorization: Bearer my-secret-webhook-token
# ── Plugins — scripts externes exécutés périodiquement ──
# Chaque plugin est un exécutable appelé toutes les 60s.
# Il lit (optionnel) sur stdin et doit écrire un objet JSON Event sur stdout.
# Exit 0 = envoyer la notif, exit ≠ 0 = ignorer.
#
# plugin:
# - name: "Disk space check"
# topic: daily
# command: /etc/ntfy-bridge/plugins/disk-check.sh
# timeout: 10
# ── Déduplication ────────────────────────────────────── # ── Déduplication ──────────────────────────────────────
dedup: dedup:
enabled: true enabled: true
+27
View File
@@ -0,0 +1,27 @@
#!/bin/bash
# Example ntfy-bridge plugin — disk space checker
#
# A plugin receives nothing on stdin (when polled) and outputs a JSON
# Event object on stdout. Exit 0 = send, exit 1 = skip/drop.
#
# Fields: source, name, topic, priority (1-5), message, tags[], click_url, actions[]
THRESHOLD=90
USAGE=$(df -h / | tail -1 | awk '{print $5}' | sed 's/%//')
if [ "$USAGE" -lt "$THRESHOLD" ]; then
# No alert needed — exit non-zero to skip
exit 1
fi
# Output a JSON Event to stdout
cat <<EOF
{
"source": "plugin",
"name": "disk-check",
"topic": "daily",
"priority": 4,
"message": "💾 **Disk usage: ${USAGE}%**\nThreshold: ${THRESHOLD}%\nHost: $(hostname)",
"tags": ["warning"]
}
EOF
+200
View File
@@ -161,6 +161,9 @@ fn (mut app App) start_background_tasks() {
if app.cfg.sources.docker.len > 0 { if app.cfg.sources.docker.len > 0 {
spawn app.docker_watch_loop() spawn app.docker_watch_loop()
} }
if app.cfg.sources.plugin.len > 0 {
spawn app.plugin_loop()
}
// Grouping flush goroutine // Grouping flush goroutine
spawn app.group_flush_loop() spawn app.group_flush_loop()
} }
@@ -245,6 +248,7 @@ fn (mut app App) config_api_response() http.Response {
safe['http_poll_sources'] = app.cfg.sources.http_poll.len.str() safe['http_poll_sources'] = app.cfg.sources.http_poll.len.str()
safe['cron_sources'] = app.cfg.sources.cron.len.str() safe['cron_sources'] = app.cfg.sources.cron.len.str()
safe['generic_sources'] = app.cfg.sources.generic.len.str() safe['generic_sources'] = app.cfg.sources.generic.len.str()
safe['plugin_sources'] = app.cfg.sources.plugin.len.str()
return new_json_response(.ok, json.encode(safe)) return new_json_response(.ok, json.encode(safe))
} }
@@ -288,6 +292,14 @@ fn (mut app App) handle_webhook(req http.Request) http.Response {
return new_json_response(.not_found, '{"error":"unknown webhook"}') return new_json_response(.not_found, '{"error":"unknown webhook"}')
} }
// ACL check — validate client IP/token before HMAC
client_ip := extract_client_ip(req)
client_token := extract_bearer_token(req)
if !app.check_acl(route, client_ip, client_token) {
app.logger.warn('server', 'webhook', 'ACL denied for ${path} from ${client_ip}')
return new_json_response(.forbidden, '{"error":"access denied"}')
}
// HMAC validation (if secret is configured) // HMAC validation (if secret is configured)
if app.cfg.server.hmac_secret != '' { if app.cfg.server.hmac_secret != '' {
sig_header := req.header.get_custom('X-Hub-Signature-256') or { '' } sig_header := req.header.get_custom('X-Hub-Signature-256') or { '' }
@@ -315,6 +327,10 @@ fn (mut app App) handle_webhook(req http.Request) http.Response {
src := app.cfg.sources.generic[route.index] src := app.cfg.sources.generic[route.index]
event_opt = sources.transform_generic(req.data, src) event_opt = sources.transform_generic(req.data, src)
} }
.plugin {
// Plugin sources are background tasks, not webhook-triggered
return new_json_response(.bad_request, '{"error":"plugin sources are not webhook-triggered"}')
}
} }
event := event_opt or { event := event_opt or {
@@ -596,7 +612,191 @@ fn (mut app App) flush_group_buffer() {
// v0.5 — Dispatch to outgoing webhooks // v0.5 — Dispatch to outgoing webhooks
app.dispatch_outgoing(e) app.dispatch_outgoing(e)
// v0.6 — Save state after each flush
app.save_state()
} }
app.group_buffer = map[string][]GroupEntry{} app.group_buffer = map[string][]GroupEntry{}
app.group_last_flush = time.now() app.group_last_flush = time.now()
} }
// ── v0.6 — ACL (Access Control Lists) ───────────────────────────────────
// check_acl validates the client IP and/or token against the webhook route's ACL config.
// Returns true if access is allowed (no ACL = open access).
fn (mut app App) check_acl(route WebhookRoute, client_ip string, client_token string) bool {
acl := app.get_acl_for_route(route)
if acl.allowed_ips.len == 0 && acl.allowed_tokens.len == 0 {
return true // no ACL configured → open access
}
// Check IP first
if acl.allowed_ips.len > 0 {
if client_ip == '' {
return false
}
mut ip_ok := false
for allowed in acl.allowed_ips {
if ip_matches(client_ip, allowed) {
ip_ok = true
break
}
}
if !ip_ok {
return false
}
}
// Check token
if acl.allowed_tokens.len > 0 {
if client_token == '' {
return false
}
mut token_ok := false
for allowed in acl.allowed_tokens {
if client_token == allowed {
token_ok = true
break
}
}
if !token_ok {
return false
}
}
return true
}
// get_acl_for_route returns the ACL config for a given webhook route.
// Returns an empty ACL (open access) if the index is out of range.
fn (mut app App) get_acl_for_route(route WebhookRoute) sources.AclConfig {
match route.kind {
.gitea {
if route.index >= app.cfg.sources.gitea.len {
return sources.AclConfig{}
}
src := app.cfg.sources.gitea[route.index]
return src.acl
}
.uptime_kuma {
if route.index >= app.cfg.sources.uptime_kuma.len {
return sources.AclConfig{}
}
src := app.cfg.sources.uptime_kuma[route.index]
return src.acl
}
.cron {
if route.index >= app.cfg.sources.cron.len {
return sources.AclConfig{}
}
src := app.cfg.sources.cron[route.index]
return src.acl
}
.generic {
if route.index >= app.cfg.sources.generic.len {
return sources.AclConfig{}
}
src := app.cfg.sources.generic[route.index]
return src.acl
}
.plugin {
if route.index >= app.cfg.sources.plugin.len {
return sources.AclConfig{}
}
src := app.cfg.sources.plugin[route.index]
return src.acl
}
}
}
// extract_client_ip extracts the real client IP from request headers or remote address.
fn extract_client_ip(req http.Request) string {
// Check common proxy/forwarded headers first
forwarded_val := req.header.get_custom('X-Forwarded-For') or { '' }
if forwarded_val != '' {
return forwarded_val.split(',')[0].trim_space()
}
real_ip := req.header.get_custom('X-Real-IP') or { '' }
if real_ip != '' {
return real_ip.trim_space()
}
// Fall back to remote address from the connection
// The http.Request doesn't expose the remote addr directly, so we check the header
return ''
}
// extract_bearer_token extracts a Bearer token from the Authorization header.
fn extract_bearer_token(req http.Request) string {
auth := req.header.get(.authorization) or { return '' }
if auth.starts_with('Bearer ') {
return auth[7..].trim_space()
}
return ''
}
// ip_matches checks if an IP address matches a pattern (CIDR or exact IP).
// Supports full CIDR notation (e.g. "192.168.0.0/16") and exact IP matching.
fn ip_matches(ip string, pattern string) bool {
if ip == pattern {
return true
}
if !pattern.contains('/') {
return ip == pattern
}
// CIDR matching
parts := pattern.split('/')
if parts.len != 2 {
return false
}
cidr_ip := parts[0]
cidr_bits := strconv.atoi(parts[1]) or { return false }
ip_parts := parse_ip(ip) or { return false }
cidr_parts := parse_ip(cidr_ip) or { return false }
// Convert to 32-bit integer for IPv4
ip_int := u32(ip_parts[0]) << 24 | u32(ip_parts[1]) << 16 | u32(ip_parts[2]) << 8 | u32(ip_parts[3])
cidr_int := u32(cidr_parts[0]) << 24 | u32(cidr_parts[1]) << 16 | u32(cidr_parts[2]) << 8 | u32(cidr_parts[3])
mask := ~u32(0) << (32 - u32(cidr_bits))
return (ip_int & mask) == (cidr_int & mask)
}
// parse_ip parses an IPv4 dotted-quad string into 4 octets.
fn parse_ip(ip string) ![]u8 {
parts := ip.split('.')
if parts.len != 4 {
return error('invalid IP')
}
mut octets := []u8{}
for part in parts {
val := strconv.atoi(part) or { return error('invalid octet: ${part}') }
if val < 0 || val > 255 {
return error('octet out of range: ${val}')
}
octets << u8(val)
}
return octets
}
// ── v0.6 — Plugin System ────────────────────────────────────────────────
// plugin_loop runs each configured plugin on a timer (every 60s by default).
fn (mut app App) plugin_loop() {
for {
if app.shutdown_flag {
break
}
for i, src in app.cfg.sources.plugin {
event := sources.transform_plugin('', src) or {
app.logger.debug('server', 'plugin', '${src.name}: no event (exit non-zero or error)')
continue
}
app.logger.info('server', 'plugin', '${src.name}: received event → ${src.topic}')
app.publish(event)
// Prevent unused variable warning for `i`
_ = i
}
time.sleep(60 * time.second)
}
}
+21
View File
@@ -35,6 +35,13 @@ pub:
clear bool // whether to dismiss notification after action clear bool // whether to dismiss notification after action
} }
// v0.6 — ACL per webhook path
pub struct AclConfig {
pub mut:
allowed_ips []string @[json: 'allowed_ips'] // CIDR ranges or single IPs
allowed_tokens []string @[json: 'allowed_tokens'] // bearer tokens
}
pub struct GiteaSource { pub struct GiteaSource {
pub mut: pub mut:
name string @[json: 'name'] name string @[json: 'name']
@@ -44,6 +51,7 @@ pub mut:
priority_map map[string]int @[json: 'priority_map'] priority_map map[string]int @[json: 'priority_map']
tags []string @[json: 'tags'] tags []string @[json: 'tags']
repo string @[json: 'repo'] repo string @[json: 'repo']
acl AclConfig @[json: 'acl']
} }
pub struct UptimeKumaState { pub struct UptimeKumaState {
@@ -61,6 +69,7 @@ pub mut:
priority_map map[string]int @[json: 'priority_map'] priority_map map[string]int @[json: 'priority_map']
tags []string @[json: 'tags'] tags []string @[json: 'tags']
state_map map[string]UptimeKumaState @[json: 'state_map'] state_map map[string]UptimeKumaState @[json: 'state_map']
acl AclConfig @[json: 'acl']
} }
pub struct DockerSource { pub struct DockerSource {
@@ -99,6 +108,7 @@ pub mut:
template string @[json: 'template'] template string @[json: 'template']
priority_map map[string]int @[json: 'priority_map'] priority_map map[string]int @[json: 'priority_map']
tags []string @[json: 'tags'] tags []string @[json: 'tags']
acl AclConfig @[json: 'acl']
} }
pub struct GenericSource { pub struct GenericSource {
@@ -109,6 +119,17 @@ pub mut:
template string @[json: 'template'] template string @[json: 'template']
priority_map map[string]int @[json: 'priority_map'] priority_map map[string]int @[json: 'priority_map']
tags []string @[json: 'tags'] tags []string @[json: 'tags']
acl AclConfig @[json: 'acl']
}
// v0.6 — Plugin source: external executable that receives JSON on stdin, outputs JSON on stdout
pub struct PluginSource {
pub mut:
name string @[json: 'name']
topic string @[json: 'topic']
command string @[json: 'command'] // path to executable
timeout int @[json: 'timeout'] // seconds (default: 10)
acl AclConfig @[json: 'acl']
} }
// render_source_template applies variable substitution to a template string. // render_source_template applies variable substitution to a template string.
+86
View File
@@ -0,0 +1,86 @@
module sources
import os
import time
import x.json2 as json
// transform_plugin runs an external executable, passes the raw event on stdin,
// and expects a JSON Event on stdout. Exit code 0 = success, non-zero = drop.
// timeout defaults to 10s if not set.
pub fn transform_plugin(raw string, source PluginSource) ?Event {
timeout := if source.timeout > 0 { source.timeout } else { 10 }
// Write raw input to a temp file for the plugin to read (optional)
// We use os.execute with timeout via a wrapper
start := time.now()
// Run the plugin command. The plugin receives nothing on stdin;
// it's a polling model — the plugin is called periodically and
// produces output when it has something to report.
result := os.execute(source.command)
elapsed := time.now().unix() - start.unix()
if elapsed > timeout {
return none
}
// Exit code 0 = event, non-zero = skip
if result.exit_code != 0 {
return none
}
output := result.output.trim_space()
if output == '' {
return none
}
// Parse plugin output as JSON
doc := json.decode[json.Any](output, json.DecoderOptions{}) or { return none }
obj := doc.as_map()
priority := obj['priority'] or { json.Any(3) }.int()
mut tags := []string{}
if tags_any := obj['tags'] {
for tag in tags_any.as_array() {
tags << tag.str()
}
}
message := obj['message'] or { json.Any(output) }.str()
click_url := obj['click_url'] or { json.Any('') }.str()
mut actions := []Action{}
if actions_any := obj['actions'] {
for a in actions_any.as_array() {
am := a.as_map()
actions << Action{
action: am['action'] or { json.Any('view') }.str()
label: am['label'] or { json.Any('') }.str()
url: am['url'] or { json.Any('') }.str()
clear: am['clear'] or { json.Any(false) }.bool()
}
}
}
// Resolve topic: plugin output can override, or fall back to source config
mut topic := source.topic
if plugin_topic := obj['topic'] {
t := plugin_topic.str()
if t != '' {
topic = t
}
}
return Event{
source: 'plugin'
name: source.name
topic: topic
priority: int(priority)
tags: tags
message: message
raw: output
click_url: click_url
actions: actions
}
}
+6
View File
@@ -6,6 +6,7 @@ pub enum SourceKind {
uptime_kuma uptime_kuma
cron cron
generic generic
plugin
} }
pub struct WebhookRoute { pub struct WebhookRoute {
@@ -42,6 +43,11 @@ pub fn build_webhook_routes(cfg Config) map[string]WebhookRoute {
index: i index: i
} }
} }
for i, _ in cfg.sources.plugin {
// Plugin sources don't have webhook_path — they run as background tasks
// Reserved for future webhook-triggered plugin support
_ = i
}
return routes return routes
} }