diff --git a/.gitignore b/.gitignore index 9a78dde..0fc0ec4 100644 --- a/.gitignore +++ b/.gitignore @@ -23,6 +23,9 @@ ntfy-bridge.yaml *.swp *.swo +# Windows build artifacts +*.def + # OS .DS_Store Thumbs.db diff --git a/config.v b/config.v new file mode 100644 index 0000000..0d0c427 --- /dev/null +++ b/config.v @@ -0,0 +1,152 @@ +module main + +import os +import yaml +import sources + +pub struct ServerConfig { +pub mut: + url string @[json: 'url'] + auth_token string @[json: 'auth_token'] + listen string @[json: 'listen'] + hmac_secret string @[json: 'hmac_secret'] // shared secret for webhook HMAC-SHA256 validation +} + +pub struct DefaultConfig { +pub mut: + priority int @[json: 'priority'] + tags []string @[json: 'tags'] +} + +pub struct RateLimitConfig { +pub mut: + max_per_minute int @[json: 'max_per_minute'] + max_per_source int @[json: 'max_per_source'] +} + +pub struct DedupConfig { +pub mut: + enabled bool @[json: 'enabled'] + ttl_seconds int @[json: 'ttl_seconds'] + rate_limit RateLimitConfig @[json: 'rate_limit'] +} + +pub struct SourcesConfig { +pub mut: + gitea []sources.GiteaSource @[json: 'gitea'] + uptime_kuma []sources.UptimeKumaSource @[json: 'uptime_kuma'] + docker []sources.DockerSource @[json: 'docker'] + http_poll []sources.HttpPollSource @[json: 'http_poll'] + cron []sources.CronSource @[json: 'cron'] + generic []sources.GenericSource @[json: 'generic'] +} + +pub struct Config { +pub mut: + server ServerConfig @[json: 'server'] + defaults DefaultConfig @[json: 'defaults'] + sources SourcesConfig @[json: 'sources'] + dedup DedupConfig @[json: 'dedup'] +} + +pub fn load_config(path string) !Config { + if !os.exists(path) { + return error('config file not found: ${path}') + } + mut cfg := yaml.decode_file[Config](path)! + + // Env override for auth token + env_token := os.getenv('NTFY_TOKEN') + if env_token != '' { + cfg.server.auth_token = env_token + } + + // Env override for HMAC secret + env_hmac := os.getenv('NTFY_HMAC_SECRET') + if env_hmac != '' { + cfg.server.hmac_secret = env_hmac + } + + validate_config(cfg)! + return cfg +} + +pub fn validate_config(cfg Config) ! { + if cfg.server.url == '' { + return error('server.url is required') + } + if cfg.server.listen == '' { + return error('server.listen is required') + } + if cfg.defaults.priority < 1 || cfg.defaults.priority > 5 { + return error('defaults.priority must be between 1 and 5') + } + + // Collect all webhook paths and detect duplicates + mut seen_paths := map[string]string{} // path -> source description + + for src in cfg.sources.gitea { + if src.webhook_path == '' { + return error('gitea source "${src.name}" missing webhook_path') + } + if src.topic == '' { + return error('gitea source "${src.name}" missing topic') + } + check_duplicate_path(mut seen_paths, src.webhook_path, 'gitea:${src.name}')! + } + for src in cfg.sources.uptime_kuma { + if src.webhook_path == '' { + return error('uptime_kuma source "${src.name}" missing webhook_path') + } + if src.topic == '' { + return error('uptime_kuma source "${src.name}" missing topic') + } + check_duplicate_path(mut seen_paths, src.webhook_path, 'uptime_kuma:${src.name}')! + } + for src in cfg.sources.docker { + if src.topic == '' { + return error('docker source "${src.name}" missing topic') + } + if src.hosts.len == 0 { + return error('docker source "${src.name}" requires at least one host') + } + } + for src in cfg.sources.http_poll { + if src.interval <= 0 { + return error('http_poll source "${src.name}" interval must be > 0') + } + for chk in src.checks { + if chk.url == '' { + return error('http_poll check missing url') + } + if chk.topic == '' { + return error('http_poll check for "${chk.url}" missing topic') + } + } + } + for src in cfg.sources.cron { + if src.webhook_path == '' { + return error('cron source "${src.name}" missing webhook_path') + } + if src.topic == '' { + return error('cron source "${src.name}" missing topic') + } + check_duplicate_path(mut seen_paths, src.webhook_path, 'cron:${src.name}')! + } + for src in cfg.sources.generic { + if src.webhook_path == '' { + return error('generic source "${src.name}" missing webhook_path') + } + if src.topic == '' { + return error('generic source "${src.name}" missing topic') + } + check_duplicate_path(mut seen_paths, src.webhook_path, 'generic:${src.name}')! + } +} + +fn check_duplicate_path(mut seen map[string]string, path string, source_name string) ! { + if existing := seen[path] { + return error('duplicate webhook path "${path}" used by ${existing} and ${source_name}') + } + seen[path] = source_name +} diff --git a/config_test.v b/config_test.v new file mode 100644 index 0000000..d3091da --- /dev/null +++ b/config_test.v @@ -0,0 +1,63 @@ +module main + +import sources + +fn test_validate_minimal_config() { + cfg := Config{ + server: ServerConfig{ + url: 'https://ntfy.example.com' + listen: ':9090' + } + defaults: DefaultConfig{ + priority: 3 + tags: ['loudspeaker'] + } + sources: SourcesConfig{ + generic: [ + sources.GenericSource{ + name: 'test' + webhook_path: '/webhooks/test' + topic: 'test' + }, + ] + } + dedup: DedupConfig{ + enabled: true + ttl_seconds: 300 + rate_limit: RateLimitConfig{ + max_per_minute: 10 + max_per_source: 30 + } + } + } + validate_config(cfg) or { panic(err) } +} + +fn test_validate_missing_url() { + cfg := Config{ + server: ServerConfig{ + listen: ':9090' + } + defaults: DefaultConfig{ + priority: 3 + } + } + if _ := validate_config(cfg) { + panic('expected validation to fail') + } +} + +fn test_validate_invalid_priority() { + cfg := Config{ + server: ServerConfig{ + url: 'https://ntfy.example.com' + listen: ':9090' + } + defaults: DefaultConfig{ + priority: 10 + } + } + if _ := validate_config(cfg) { + panic('expected validation to fail') + } +} diff --git a/dashboard.html b/dashboard.html new file mode 100644 index 0000000..61871e8 --- /dev/null +++ b/dashboard.html @@ -0,0 +1,172 @@ + + + + + +ntfy-bridge Dashboard + + + +
+

πŸ”” ntfy-bridge

+LIVE +
+
+ +
+
+

πŸ“Š Notifications

+
+
β€”
Total
+
β€”
Errors
+
+
+ +
+

⏱️ Uptime

+
β€”
+
Server running
+
+ +
+

πŸ”— Ntfy Server

+
β€”
+
Target
+
+
+ +
+
+

πŸ“‹ Details

+
+ + + +
+
+
Loading...
+
+
+
Loading...
+
+
+
Loading...
+
+
+
+
+ + + + + diff --git a/dedup.v b/dedup.v new file mode 100644 index 0000000..e9d932c --- /dev/null +++ b/dedup.v @@ -0,0 +1,129 @@ +module main + +import time +import sources + +pub struct DedupEntry { +pub mut: + last_seen time.Time + count int +} + +pub struct Dedup { +pub mut: + enabled bool + ttl_seconds int + max_per_minute int + max_per_source int + cache map[string]DedupEntry + topic_history map[string][]time.Time + source_history map[string][]time.Time +} + +pub fn new_dedup(cfg DedupConfig) Dedup { + return Dedup{ + enabled: cfg.enabled + ttl_seconds: cfg.ttl_seconds + max_per_minute: cfg.rate_limit.max_per_minute + max_per_source: cfg.rate_limit.max_per_source + cache: map[string]DedupEntry{} + topic_history: map[string][]time.Time{} + source_history: map[string][]time.Time{} + } +} + +fn simple_hash(s string) u64 { + mut h := u64(14695981039346656037) + for c in s.bytes() { + h ^= u64(c) + h *= u64(1099511628211) + } + return h +} + +fn dedup_key(event sources.Event) string { + h := simple_hash(event.message) + return '${event.source}:${event.name}:${event.topic}:${h}' +} + +fn (mut d Dedup) cleanup(now time.Time) { + now_unix := now.unix() + for key, entry in d.cache { + if now_unix - entry.last_seen.unix() > d.ttl_seconds { + d.cache.delete(key) + } + } + for topic, times in d.topic_history { + mut kept := []time.Time{} + for t in times { + if now_unix - t.unix() <= 60 { + kept << t + } + } + if kept.len == 0 { + d.topic_history.delete(topic) + } else { + d.topic_history[topic] = kept + } + } + for source, times in d.source_history { + mut kept := []time.Time{} + for t in times { + if now_unix - t.unix() <= 60 { + kept << t + } + } + if kept.len == 0 { + d.source_history.delete(source) + } else { + d.source_history[source] = kept + } + } +} + +pub fn (mut d Dedup) allow(event sources.Event) bool { + if !d.enabled { + return true + } + now := time.now() + d.cleanup(now) + + // Rate limit per topic + topic_times := d.topic_history[event.topic] or { []time.Time{} } + if topic_times.len >= d.max_per_minute && d.max_per_minute > 0 { + return false + } + + // Rate limit per source + source_key := '${event.source}:${event.name}' + source_times := d.source_history[source_key] or { []time.Time{} } + if source_times.len >= d.max_per_source && d.max_per_source > 0 { + return false + } + + // Dedup check + key := dedup_key(event) + entry := d.cache[key] or { DedupEntry{} } + if entry.count > 0 { + d.cache[key] = DedupEntry{ + last_seen: now + count: entry.count + 1 + } + return false + } + + d.cache[key] = DedupEntry{ + last_seen: now + count: 1 + } + + mut new_topic_times := topic_times.clone() + new_topic_times << now + d.topic_history[event.topic] = new_topic_times + + mut new_source_times := source_times.clone() + new_source_times << now + d.source_history[source_key] = new_source_times + + return true +} diff --git a/dedup_test.v b/dedup_test.v new file mode 100644 index 0000000..cbc4609 --- /dev/null +++ b/dedup_test.v @@ -0,0 +1,59 @@ +module main + +import sources + +fn test_dedup_allows_first_event() { + mut d := new_dedup(DedupConfig{ + enabled: true + ttl_seconds: 300 + rate_limit: RateLimitConfig{ + max_per_minute: 10 + max_per_source: 30 + } + }) + event := sources.Event{ + source: 'generic' + name: 'test' + topic: 'test' + message: 'hello' + } + assert d.allow(event) == true +} + +fn test_dedup_blocks_duplicate() { + mut d := new_dedup(DedupConfig{ + enabled: true + ttl_seconds: 300 + rate_limit: RateLimitConfig{ + max_per_minute: 10 + max_per_source: 30 + } + }) + event := sources.Event{ + source: 'generic' + name: 'test' + topic: 'test' + message: 'hello' + } + assert d.allow(event) == true + assert d.allow(event) == false +} + +fn test_dedup_disabled() { + mut d := new_dedup(DedupConfig{ + enabled: false + ttl_seconds: 300 + rate_limit: RateLimitConfig{ + max_per_minute: 10 + max_per_source: 30 + } + }) + event := sources.Event{ + source: 'generic' + name: 'test' + topic: 'test' + message: 'hello' + } + assert d.allow(event) == true + assert d.allow(event) == true +} diff --git a/docker_watcher.v b/docker_watcher.v new file mode 100644 index 0000000..602a6f3 --- /dev/null +++ b/docker_watcher.v @@ -0,0 +1,75 @@ +module main + +import time +$if linux || macos { + import net.unix + import os +} +import sources as _ + +fn (mut app App) docker_watch_loop() { + $if linux || macos { + for { + for src in app.cfg.sources.docker { + for host in src.hosts { + app.watch_docker_host(src, host) + } + } + time.sleep(10 * time.second) + } + } $else { + app.logger.warn('docker', 'watch', + 'Docker socket watching is only supported on Linux/macOS') + for { + time.sleep(60 * time.second) + } + } +} + +$if linux || macos { + fn (mut app App) watch_docker_host(src _.DockerSource, host string) { + socket_path := host.replace('unix://', '') + if !os.exists(socket_path) { + app.logger.warn('docker', 'watch', 'socket not found: ${socket_path}') + return + } + + mut conn := unix.connect_stream(socket_path) or { + app.logger.warn('docker', 'watch', 'failed to connect to ${socket_path}: ${err}') + return + } + defer { + conn.close() or {} + } + + filters := build_docker_filters(src.events) + request := 'GET /events?filters=${filters} HTTP/1.1\r\nHost: localhost\r\n\r\n' + conn.write_string(request) or { + app.logger.warn('docker', 'watch', 'write error: ${err}') + return + } + + mut buf := []u8{len: 4096} + for { + read := conn.read(mut buf) or { break } + if read <= 0 { + break + } + lines := buf[..read].bytestr().split_into_lines() + for line in lines { + if line.starts_with('{') { + event := sources.transform_docker_event(line, src, host) or { continue } + app.publish(event) + } + } + } + } + + fn build_docker_filters(events []string) string { + mut parts := []string{} + for event in events { + parts << '"${event}":true' + } + return '{' + parts.join(',') + '}' + } +} diff --git a/ARCHITECTURE.md b/docs/ARCHITECTURE.md similarity index 88% rename from ARCHITECTURE.md rename to docs/ARCHITECTURE.md index 07a238a..890a4fc 100644 --- a/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -92,7 +92,7 @@ POST /webhooks/cron-disk ## ModΓ¨le de donnΓ©es ```v -struct Config { +struct Config { // dans config.v server ServerConfig defaults DefaultConfig sources SourcesConfig @@ -109,7 +109,7 @@ struct DefaultConfig { tags []string // ["loudspeaker"] } -struct Source { +struct Source { // types spΓ©cifiques dans sources/event.v name string webhook_path string topic string @@ -118,7 +118,7 @@ struct Source { tags []string } -struct Event { +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 @@ -174,26 +174,23 @@ ntfy-bridge/ β”œβ”€β”€ ROADMAP.md # Phases de dΓ©veloppement β”œβ”€β”€ CONTRIBUTING.md # Guide de contribution β”œβ”€β”€ ntfy-bridge.example.yaml # Exemple de configuration complet -β”œβ”€β”€ src/ -β”‚ β”œβ”€β”€ main.v # Entrypoint, parsing CLI -β”‚ β”œβ”€β”€ config.v # Chargement + validation YAML -β”‚ β”œβ”€β”€ server.v # HTTP server + routing webhooks -β”‚ β”œβ”€β”€ event.v # Types Event, pipeline -β”‚ β”œβ”€β”€ ntfy.v # Client HTTP Ntfy (POST) -β”‚ β”œβ”€β”€ dedup.v # DΓ©duplication + rate limiting -β”‚ β”œβ”€β”€ template.v # Mini moteur de template -β”‚ β”œβ”€β”€ sources/ -β”‚ β”‚ β”œβ”€β”€ gitea.v # Transformer Gitea -β”‚ β”‚ β”œβ”€β”€ uptime_kuma.v # Transformer Uptime Kuma -β”‚ β”‚ β”œβ”€β”€ docker.v # Watcher Docker socket -β”‚ β”‚ β”œβ”€β”€ http_poll.v # Poller HTTP pΓ©riodique -β”‚ β”‚ └── generic.v # Webhook passe-partout -β”‚ └── log.v # Logging structurΓ© -β”œβ”€β”€ tests/ -β”‚ β”œβ”€β”€ config_test.v # Tests validation config -β”‚ β”œβ”€β”€ gitea_test.v # Tests transformer Gitea -β”‚ β”œβ”€β”€ uptime_kuma_test.v # Tests transformer Kuma -β”‚ └── dedup_test.v # Tests dΓ©duplication +β”œβ”€β”€ 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) +β”œβ”€β”€ sources/ # Types Event + transformers +β”‚ β”œβ”€β”€ event.v # Types Event, Source +β”‚ β”œβ”€β”€ gitea.v # Transformer Gitea +β”‚ β”œβ”€β”€ uptime_kuma.v # Transformer Uptime Kuma +β”‚ β”œβ”€β”€ docker.v # Transformer Docker events +β”‚ β”œβ”€β”€ http_poll.v # Transformer HTTP poll +β”‚ β”œβ”€β”€ generic.v # Webhook passe-partout +β”‚ └── cron.v # Webhook cron +β”œβ”€β”€ *_test.v # Tests unitaires β”œβ”€β”€ scripts/ β”‚ └── example-cron-disk.sh # Exemple script cron β†’ Ntfy └── v.mod # DΓ©pendances V diff --git a/CONTRIBUTING.md b/docs/CONTRIBUTING.md similarity index 59% rename from CONTRIBUTING.md rename to docs/CONTRIBUTING.md index a0feabe..924f4eb 100644 --- a/CONTRIBUTING.md +++ b/docs/CONTRIBUTING.md @@ -30,30 +30,34 @@ v -prod . ## Structure du code ``` -src/ -β”œβ”€β”€ main.v # Entrypoint, flag parsing, orchestration -β”œβ”€β”€ config.v # Config loading + validation -β”œβ”€β”€ server.v # HTTP server, webhook routing -β”œβ”€β”€ event.v # Event types, pipeline -β”œβ”€β”€ ntfy.v # Ntfy HTTP client -β”œβ”€β”€ dedup.v # Deduplication + rate limiting -β”œβ”€β”€ template.v # Mini template engine -β”œβ”€β”€ log.v # Structured logging -└── sources/ # Source-specific transformers - β”œβ”€β”€ gitea.v - β”œβ”€β”€ uptime_kuma.v - β”œβ”€β”€ docker.v - β”œβ”€β”€ http_poll.v - └── generic.v +. +β”œβ”€β”€ main.v # Entrypoint, flag parsing, orchestration +β”œβ”€β”€ config.v # Config loading + validation +β”œβ”€β”€ server.v # HTTP server, webhook routing +β”œβ”€β”€ ntfy.v # Ntfy HTTP client +β”œβ”€β”€ dedup.v # Deduplication + rate limiting +β”œβ”€β”€ template.v # Mini template engine +β”œβ”€β”€ log.v # Structured logging +β”œβ”€β”€ docker_watcher.v # Docker socket watcher (Unix) +β”œβ”€β”€ sources/ # Event types + source-specific transformers +β”‚ β”œβ”€β”€ event.v +β”‚ β”œβ”€β”€ gitea.v +β”‚ β”œβ”€β”€ uptime_kuma.v +β”‚ β”œβ”€β”€ docker.v +β”‚ β”œβ”€β”€ http_poll.v +β”‚ β”œβ”€β”€ generic.v +β”‚ └── cron.v +└── *_test.v # Unit tests ``` ## Ajouter une nouvelle source -1. CrΓ©er `src/sources/ma_source.v` -2. ImplΓ©menter la fonction `transform(raw string, source &Source) ?Event` -3. DΓ©clarer le type dans `sources` du `config.v` -4. Router l'endpoint dans `server.v` -5. Ajouter un test dans `tests/ma_source_test.v` +1. CrΓ©er `sources/ma_source.v` +2. ImplΓ©menter la fonction `transform(raw string, source MaSource) ?Event` +3. DΓ©clarer le type source dans `sources/event.v` +4. RΓ©fΓ©rencer le type dans `config.v` (`SourcesConfig`) +5. Router l'endpoint dans `server.v` +6. Ajouter un test dans `ma_source_test.v` ## Conventions de code @@ -72,7 +76,7 @@ src/ v test . # Test spΓ©cifique -v test tests/gitea_test.v +v test gitea_test.v # Avec verbose v test . -stats diff --git a/ROADMAP.md b/docs/ROADMAP.md similarity index 100% rename from ROADMAP.md rename to docs/ROADMAP.md diff --git a/docs/WEBHOOK_GITEA.md b/docs/WEBHOOK_GITEA.md new file mode 100644 index 0000000..57f3574 --- /dev/null +++ b/docs/WEBHOOK_GITEA.md @@ -0,0 +1,255 @@ +# Configurer les webhooks Gitea avec ntfy-bridge + +Guide tape par tape pour connecter Gitea ntfy-bridge et recevoir des notifications Ntfy pour chaque push, PR, issue ou release. + +--- + +## Prrequis + +- **ntfy-bridge** install et configur (voir [README.md](../README.md)) +- Un serveur **Ntfy** accessible (ex: `https://ntfy.dracodev.net`) +- Un compte **Gitea** avec accs admin au repo (ex: `https://git.dracodev.net`) + +--- + +## Architecture + +``` + Gitea (git.dracodev.net) ntfy-bridge (ton serveur) Ntfy (ntfy.dracodev.net) + β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” + β”‚ bruno/flowdeck β”‚ β”‚ coute sur :9090 β”‚ β”‚ Topic: dev-notifs β”‚ + β”‚ β”‚ POST β”‚ β”‚ POST β”‚ β”‚ + β”‚ vnement (push) │──────────→│ /webhooks/gitea-... │───────→│ Ton tlphone β”‚ + β”‚ JSON + HMAC sig β”‚ β”‚ β”‚ β”‚ β”‚ + β”‚ β”‚ β”‚ Parse, format, forward β”‚ β”‚ β”‚ + β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ +``` + +**Le principe** : Gitea envoie un POST HTTP contenant un JSON avec tous les dtails de l'vnement. ntfy-bridge l'coute, le parse, le formate, et le forward vers Ntfy. + +--- + +## tape 1 β€” Configurer ntfy-bridge + +### 1a. Modifier `ntfy-bridge.yaml` + +```yaml +server: + url: https://ntfy.dracodev.net # Ton serveur Ntfy + listen: "0.0.0.0:9090" # couter sur le rseau (pas juste localhost) + hmac_secret: "un-secret-que-tu-choisis" # Secret partag avec Gitea + +sources: + gitea: + - name: "FlowDeck activity" # Nom affich dans les logs/dashboard + webhook_path: /webhooks/gitea-flowdeck # Chemin d'coute + repo: bruno/flowdeck # Repo concern + topic: dev-notifs # Topic Ntfy cible + priority_map: # Priorits par type d'vent + pull_request: 4 # 4 = high + issue: 3 # 3 = default + push: 2 # 2 = low + template: | # Format du message (optionnel) + {user} {action} sur **[{repo}]**: "{title}" +``` + +**Points cls** : +- `listen: "0.0.0.0:9090"` β€” indispensable si Gitea est sur un autre serveur. `127.0.0.1` = localhost seulement. +- `hmac_secret` β€” le mme secret que tu mettras dans Gitea. Optionnel mais recommand. +- `webhook_path` β€” le chemin dans l'URL. Doit tre unique par source. +- `template` β€” si absent, un format par dfaut est gnr automatiquement avec emoji. + +### 1b. Lancer ntfy-bridge + +```bash +cd /c/dev/git/V/ntfy-bridge +./ntfy-bridge.exe --config ntfy-bridge.yaml +``` + +Sortie attendue : + +``` +────────────────────────────────────────────────────────────────────── + ntfy-bridge v0.1.0 β€” Homelab Notification Hub +────────────────────────────────────────────────────────────────────── + Ntfy server: https://ntfy.dracodev.net + Listening on: http://0.0.0.0:9090 + Dashboard: http://0.0.0.0:9090/ + + Sources loaded: + Gitea: 2 webhooks + Uptime Kuma: 1 webhooks + Docker: 1 watchers + HTTP Poll: 1 pollers (5 checks) + Cron: 2 webhooks + Generic: 1 webhooks +────────────────────────────────────────────────────────────────────── +``` + +### 1c. Vrifier que a marche + +```bash +# Test health endpoint +curl http://localhost:9090/health +# β†’ {"status":"ok","uptime_seconds":12,"total_notifications":0} + +# Test dashboard +curl http://localhost:9090/ +# β†’ HTML du dashboard +``` + +--- + +## tape 2 β€” Configurer Gitea + +### 2a. Aller dans les rglages du repo + +1. Ouvre `https://git.dracodev.net/bruno/flowdeck` +2. Clic sur l'engrenage **Settings** (haut droite) +3. Menu gauche : **Webhooks** +4. Bouton **Add Webhook** β†’ choisir **Gitea** + +### 2b. Remplir le formulaire + +| Champ | Valeur | Notes | +|---|---|---| +| **Target URL** | `http://192.168.30.X:9090/webhooks/gitea-flowdeck` | Remplacer X par l'IP relle du serveur | +| **HTTP Method** | POST | | +| **POST Content Type** | application/json | | +| **Secret** | `un-secret-que-tu-choisis` | Identique `hmac_secret` dans le YAML | +| **Trigger On** | Slectionner les vnements souhauts | Push, Pull Request, Issue, Release... | + +**Trouver la bonne IP** : sur le serveur o tourne ntfy-bridge : + +```bash +# Linux +ip addr show | grep "inet " | grep -v 127.0.0.1 + +# Windows +ipconfig | findstr "IPv4" +``` + +### 2c. Cliquer "Add Webhook" + +Gitea envoie immdiatement un **test ping** pour vrifier la connectivit. + +Rsultat russi : + +``` +βœ“ Delivery successful (HTTP 200) +``` + +En cas d'chec : + +| Erreur | Cause probable | Solution | +|---|---|---| +| Connection refused | ntfy-bridge n'est pas lanc | Lancer `./ntfy-bridge.exe` | +| Connection refused | Firewall bloque le port 9090 | Ouvrir le port dans le firewall | +| HTTP 404 | Mauvais `webhook_path` | Vrifier que `/webhooks/gitea-flowdeck` est exact | +| HTTP 401 | Mauvais `hmac_secret` | Vrifier que le secret est identique des deux cts | +| Timeout | Mauvaise IP ou rseau inaccessible | Vrifier la connectivit rseau | + +--- + +## tape 3 β€” Tester avec un vrai vnement + +### Option A : Test simul (curl) + +Depuis n'importe quel poste : + +```bash +curl -X POST http://192.168.30.X:9090/webhooks/gitea-flowdeck \ + -H "Content-Type: application/json" \ + -H "X-Hub-Signature-256: sha256=..." \ + -d '{ + "commits": [{"id": "abc1234567890abcdef1234567890abcdef1234", "message": "test notification ntfy-bridge"}], + "pusher": {"username": "bruno"}, + "repository": {"full_name": "bruno/flowdeck"} + }' +``` + +> Note : si `hmac_secret` est configur, gnre la signature correcte ou dsactive-le temporairement pour le test. + +### Option B : Test rel + +Fais un push, cre une PR, ou ouvre une issue dans le repo Gitea. Tu devrais recevoir une notification Ntfy dans les secondes qui suivent. + +Exemple de notification reue : + +``` +bruno pushed sur [bruno/flowdeck]: "fix: corrige le bug du dashboard" +``` + +Ou avec le template par dfaut : + +``` + [bruno/flowdeck] bruno pushed: "fix: corrige le bug du dashboard" (abc1234) +``` + +--- + +## Ajouter d'autres repos + +Rpte les mmes tapes pour chaque repo : + +### Dans `ntfy-bridge.yaml` + +```yaml +sources: + gitea: + - name: "FlowDeck activity" + webhook_path: /webhooks/gitea-flowdeck + repo: bruno/flowdeck + topic: dev-notifs + + - name: "ObsiGate activity" + webhook_path: /webhooks/gitea-obsigate + repo: bruno/obsigate + topic: dev-notifs + + - name: "Imago activity" + webhook_path: /webhooks/gitea-imago + repo: bruno/imago + topic: dev-notifs +``` + +### Dans Gitea + +Pour chaque repo : Settings β†’ Webhooks β†’ Add Webhook avec l'URL correspondante. + +```text +http://192.168.30.X:9090/webhooks/gitea-flowdeck β†’ repo bruno/flowdeck +http://192.168.30.X:9090/webhooks/gitea-obsigate β†’ repo bruno/obsigate +http://192.168.30.X:9090/webhooks/gitea-imago β†’ repo bruno/imago +``` + +--- + +## Dpannage rapide + +```bash +# Vrifier que ntfy-bridge tourne +curl http://localhost:9090/health + +# Voir les stats +curl http://localhost:9090/stats + +# Voir la config (sans les secrets) +curl http://localhost:9090/api/config + +# Voir l'historique des derniers vnements +curl http://localhost:9090/api/history + +# Voir le dashboard web +# Ouvre http://localhost:9090/ dans ton navigateur +``` + +--- + +## Rfrences + +- [README.md](../README.md) β€” Vue d'ensemble du projet +- [ARCHITECTURE.md](../ARCHITECTURE.md) β€” Architecture dtaille +- [ntfy-bridge.example.yaml](../ntfy-bridge.example.yaml) β€” Exemple de config complet +- [Documentation webhooks Gitea](https://docs.gitea.com/usage/webhooks) +- [Documentation Ntfy](https://docs.ntfy.sh/) diff --git a/gitea_test.v b/gitea_test.v new file mode 100644 index 0000000..6c57849 --- /dev/null +++ b/gitea_test.v @@ -0,0 +1,37 @@ +module main + +import sources + +fn test_transform_gitea_push() { + raw := '{"repository":{"full_name":"bruno/flowdeck"},"pusher":{"username":"bruno"},"commits":[{"message":"Fix bug","id":"a3f2c1d1234567890"}]}' + src := sources.GiteaSource{ + name: 'FlowDeck' + webhook_path: '/webhooks/gitea-flowdeck' + topic: 'dev-notifs' + priority_map: { + 'push': 2 + } + } + event := sources.transform_gitea(raw, src) or { panic(err) } + assert event.source == 'gitea' + assert event.topic == 'dev-notifs' + assert event.priority == 2 + assert event.message.contains('bruno/flowdeck') + assert event.message.contains('bruno') + assert event.message.contains('Fix bug') +} + +fn test_transform_gitea_pull_request() { + raw := '{"repository":{"full_name":"bruno/flowdeck"},"action":"opened","pull_request":{"title":"New feature","merge_commit_sha":"abc123"}}' + src := sources.GiteaSource{ + name: 'FlowDeck' + webhook_path: '/webhooks/gitea-flowdeck' + topic: 'dev-notifs' + priority_map: { + 'pull_request': 4 + } + } + event := sources.transform_gitea(raw, src) or { panic(err) } + assert event.priority == 4 + assert event.message.contains('New feature') +} diff --git a/log.v b/log.v new file mode 100644 index 0000000..a28b826 --- /dev/null +++ b/log.v @@ -0,0 +1,119 @@ +module main + +import time + +pub enum LogLevel { + debug + info + warn + error +} + +pub struct Logger { + level LogLevel + json_mode bool +} + +pub fn new_logger(level LogLevel, json_mode bool) Logger { + return Logger{ + level: level + json_mode: json_mode + } +} + +fn level_string(level LogLevel) string { + return match level { + .debug { 'DEBUG' } + .info { 'INFO' } + .warn { 'WARN' } + .error { 'ERROR' } + } +} + +fn level_icon(level LogLevel) string { + return match level { + .debug { 'πŸ”' } + .info { 'ℹ️ ' } + .warn { '⚠️ ' } + .error { '❌' } + } +} + +fn (l Logger) should_log(level LogLevel) bool { + return int(level) >= int(l.level) +} + +pub fn (l Logger) log(level LogLevel, source string, event string, message string) { + if !l.should_log(level) { + return + } + + if l.json_mode { + l.log_json(level, source, event, message) + } else { + l.log_human(level, source, event, message) + } +} + +fn (l Logger) log_json(level LogLevel, source string, event string, message string) { + timestamp := time.now().format_rfc3339() + level_str := level_string(level) + // Manual JSON construction β€” simple enough to avoid importing json for each log line + // Escape backslashes and quotes in message + escaped := escape_json(message) + output := '{"timestamp":"${timestamp}","level":"${level_str}","source":"${source}","event":"${event}","message":"${escaped}"}' + if level == .error { + eprintln(output) + } else { + println(output) + } +} + +fn escape_json(s string) string { + mut result := s.clone() + result = result.replace('\\', '\\\\') + result = result.replace('"', '\\"') + result = result.replace('\n', '\\n') + result = result.replace('\r', '\\r') + result = result.replace('\t', '\\t') + return result +} + +fn (l Logger) log_human(level LogLevel, source string, _event string, message string) { + icon := level_icon(level) + ts := time.now() + ts_str := format_time_compact(ts) + output := '${ts_str} ${icon} [${source}] ${message}' + + if level == .error { + eprintln(output) + } else if level == .warn { + eprintln(output) + } else { + println(output) + } +} + +// format_time_compact returns "HH:MM:SS" (compact time for console output) +fn format_time_compact(t time.Time) string { + h := int(t.hour) + m := int(t.minute) + s := int(t.second) + return '${h:02d}:${m:02d}:${s:02d}' +} + +pub fn (l Logger) debug(source string, event string, message string) { + l.log(.debug, source, event, message) +} + +pub fn (l Logger) info(source string, event string, message string) { + l.log(.info, source, event, message) +} + +pub fn (l Logger) warn(source string, event string, message string) { + l.log(.warn, source, event, message) +} + +pub fn (l Logger) error(source string, event string, message string) { + l.log(.error, source, event, message) +} diff --git a/main.v b/main.v new file mode 100644 index 0000000..294cce6 --- /dev/null +++ b/main.v @@ -0,0 +1,202 @@ +module main + +import os +import flag + +const version = '0.1.0' +const author = 'Bruno Charest' + +fn main() { + mut fp := flag.new_flag_parser(os.args) + fp.application('ntfy-bridge') + fp.version(version) + fp.description('Homelab notification hub β€” aggregates Gitea, Docker, Uptime Kuma β†’ Ntfy') + fp.skip_executable() + + config_path := fp.string('config', `c`, 'ntfy-bridge.yaml', 'Path to config file (YAML)') + validate_only := fp.bool('validate', `v`, false, 'Validate config and exit') + quiet := fp.bool('quiet', `q`, false, 'Quiet mode β€” suppress startup banner (for scripts/systemd)') + version_only := fp.bool('version', `V`, false, 'Print version and exit') + + // Show help if no args or --help/-h requested + if os.args.len <= 1 || '--help' in os.args || '-h' in os.args { + print_help() + return + } + + fp.finalize() or { + eprintln(err) + exit(1) + return + } + + if version_only { + println('ntfy-bridge ${version}') + return + } + + cfg := load_config(config_path) or { + eprintln('Failed to load config: ${err}') + exit(1) + return + } + + if validate_only { + println('Config is valid') + return + } + + if !quiet { + print_header() + println('') + print_startup_info(cfg) + } + + mut app := new_app(cfg, quiet) + app.run() or { + eprintln('Server error: ${err}') + exit(1) + } +} + +// ── ASCII Art Header ───────────────────────────────────────────────────── + +const box_w = 60 // inner width between β•‘ borders + +fn build_header() string { + // Tagline β€” centered within the box + tag_text := 'Homelab Notification Hub Β· v${version}' + tag_left := (box_w - tag_text.len) / 2 + tag_right := box_w - tag_left - tag_text.len + tag_line := ' '.repeat(tag_left) + tag_text + ' '.repeat(tag_right) + + // Version line β€” left-aligned with 5-space indent + vinfo := ' Version ${version} β”‚ ${author}' + vinfo_pad := box_w - vinfo.len + + return '╔══════════════════════════════════════════════════════════════╗\n' + + 'β•‘ __ ____ __ _ __ β•‘\n' + + 'β•‘ ____ / /_/ __/_ __ / /_ _____(_)___/ /___ ____ β•‘\n' + + 'β•‘ / __ \\/ __/ /_/ / / /_____/ __ \\/ ___/ / __ / __ `/ _ \\ β•‘\n' + + 'β•‘ / / / / /_/ __/ /_/ /_____/ /_/ / / / / /_/ / /_/ / __/ β•‘\n' + + 'β•‘ /_/ /_/\\__/_/ \\__, / /_.___/_/ /_/\\__,_/\\__, /\\___/ β•‘\n' + + 'β•‘ /____/ /____/ β•‘\n' + + 'β•‘' + tag_line + 'β•‘\n' + + '╠══════════════════════════════════════════════════════════════╣\n' + + 'β•‘' + vinfo + ' '.repeat(vinfo_pad) + 'β•‘\n' + + 'β•šβ•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•β•' +} + +fn print_header() { + println(build_header()) +} + +fn print_help() { + println('') + println(build_header()) + println('') + println(' A simple daemon to aggregate notifications from multiple sources') + println(' (Gitea, Uptime Kuma, Docker, HTTP health checks, Cron scripts)') + println(' and forward them as formatted Ntfy push notifications.') + println('') + println('Usage:') + println(' ntfy-bridge [flags]') + println(' ntfy-bridge --config ntfy-bridge.yaml') + println('') + println('Flags:') + println(' -c, --config Path to config file (default: ntfy-bridge.yaml)') + println(' -v, --validate Validate config and exit') + println(' -q, --quiet Quiet mode β€” suppress startup banner') + println(' -V, --version Print version and exit') + println(' -h, --help Show this help message') + println('') + println('Config:') + println(' Copy ntfy-bridge.example.yaml β†’ ntfy-bridge.yaml and edit it.') + println(' Use NTFY_TOKEN and NTFY_HMAC_SECRET env vars for secrets.') + println('') + println('Docs: https://git.dracodev.net/bruno/ntfy-bridge') + println(line_double_r) + println(' Version: ${version} | Author: ${author}') + println(line_double_r) + println('') +} + +fn print_startup_info(cfg Config) { + println(line_double_r) + println(' Ntfy server: ${cfg.server.url}') + println(' Listening on: http://${cfg.server.listen}') + println(' Dashboard: http://${normalize_listen(cfg.server.listen)}/') + println(' Health check: http://${normalize_listen(cfg.server.listen)}/health') + println('') + println(' Sources configured:') + + gitea_count := cfg.sources.gitea.len + if gitea_count > 0 { + println(' Gitea: ${gitea_count} webhook${if gitea_count > 1 { 's' } else { '' }}') + } + uk_count := cfg.sources.uptime_kuma.len + if uk_count > 0 { + println(' Uptime Kuma: ${uk_count} webhook${if uk_count > 1 { 's' } else { '' }}') + } + docker_count := cfg.sources.docker.len + if docker_count > 0 { + println(' Docker: ${docker_count} watcher${if docker_count > 1 { 's' } else { '' }}') + } + poll_count := cfg.sources.http_poll.len + if poll_count > 0 { + mut total_checks := 0 + for p in cfg.sources.http_poll { + total_checks += p.checks.len + } + println(' HTTP Poll: ${poll_count} poller${if poll_count > 1 { 's' } else { '' }} (${total_checks} check${if total_checks > 1 { 's' } else { '' }})') + } + cron_count := cfg.sources.cron.len + if cron_count > 0 { + println(' Cron: ${cron_count} webhook${if cron_count > 1 { 's' } else { '' }}') + } + generic_count := cfg.sources.generic.len + if generic_count > 0 { + println(' Generic: ${generic_count} webhook${if generic_count > 1 { 's' } else { '' }}') + } + + println('') + println(' Webhook URLs:') + for src in cfg.sources.gitea { + println(' http://${normalize_listen(cfg.server.listen)}${src.webhook_path} β†’ ${src.topic} (Gitea: ${src.name})') + } + for src in cfg.sources.uptime_kuma { + println(' http://${normalize_listen(cfg.server.listen)}${src.webhook_path} β†’ ${src.topic} (Uptime Kuma: ${src.name})') + } + for src in cfg.sources.cron { + println(' http://${normalize_listen(cfg.server.listen)}${src.webhook_path} β†’ ${src.topic} (Cron: ${src.name})') + } + for src in cfg.sources.generic { + println(' http://${normalize_listen(cfg.server.listen)}${src.webhook_path} β†’ ${src.topic} (Generic: ${src.name})') + } + + println('') + if cfg.dedup.enabled { + println(' Dedup: enabled (TTL: ${cfg.dedup.ttl_seconds}s, rate: ${cfg.dedup.rate_limit.max_per_minute}/min)') + } else { + println(' Dedup: disabled') + } + + println(line_double_r) + println('') +} + +// normalize_listen converts "0.0.0.0:9090" to "localhost:9090" for display URLs, +// and leaves specific IPs intact (e.g. "192.168.30.101:9090" stays as-is). +fn normalize_listen(listen string) string { + if listen.starts_with('0.0.0.0') { + return 'localhost' + listen[7..] + } + if listen.starts_with('127.0.0.1') { + return listen // already localhost + } + return listen +} + +// ── Display Constants ───────────────────────────────────────────────────── + +const line_double_r = '══════════════════════════════════════════════════════════════' diff --git a/ntfy-bridge.example.yaml b/ntfy-bridge.example.yaml index a00bcfa..1c773ba 100644 --- a/ntfy-bridge.example.yaml +++ b/ntfy-bridge.example.yaml @@ -14,15 +14,19 @@ server: # Mettre "0.0.0.0:9090" pour exposer rΓ©seau, "127.0.0.1:9090" pour localhost listen: "127.0.0.1:9090" + # Secret partagΓ© pour validation HMAC-SHA256 des webhooks (optionnel) + # Si dΓ©fini, tous les webhooks doivent inclure le header X-Hub-Signature-256 + # Peut aussi Γͺtre passΓ© via variable d'environnement NTFY_HMAC_SECRET + # hmac_secret: "mon-secret-partage" + # ── Defaults ────────────────────────────────────────── # AppliquΓ©s Γ  toutes les sources sauf override explicite defaults: - priority: 3 # 1=min, 2=low, 3=default, 4=high, 5=urgent - tags: ["loudspeaker"] # Tags/emojis Ntfy + priority: 3 # 1=min, 2=low, 3=default, 4=high, 5=urgent + tags: ["loudspeaker"] # Tags/emojis Ntfy # ── Sources ─────────────────────────────────────────── sources: - # ═══════════════════════════════════════════════════ # GITEA β€” webhooks depuis git.dracodev.net # ═══════════════════════════════════════════════════ @@ -54,7 +58,7 @@ sources: topic: alerts state_map: down: { priority: 5, tags: ["rotating_light", "x"] } - up: { priority: 1, tags: ["white_check_mark"] } + up: { priority: 1, tags: ["white_check_mark"] } # ═══════════════════════════════════════════════════ # DOCKER β€” surveillance des conteneurs @@ -79,13 +83,13 @@ sources: # ═══════════════════════════════════════════════════ http_poll: - name: "Health endpoints" - interval: 60 # secondes entre chaque check + interval: 60 # secondes entre chaque check checks: - - url: https://obsigate.dracodev.net/health + - url: https://og.dracodev.net/health topic: alerts - priority: 5 # prioritΓ© si DOWN + priority: 5 # prioritΓ© si DOWN expect_status: 200 - timeout: 10 # secondes + timeout: 10 # secondes - url: https://flowdeck.dracodev.net/health topic: alerts @@ -99,6 +103,18 @@ sources: expect_status: 200 timeout: 10 + - url: http://raspi.8gb.home:3001/ping + topic: alerts + priority: 5 # Uptime Kuma lui-mΓͺme β€” critique + expect_status: 200 + timeout: 10 + + - url: http://raspi.8gb.home:9119/api/status + topic: alerts + priority: 5 + expect_status: 200 + timeout: 10 + # ═══════════════════════════════════════════════════ # CRON β€” reΓ§oit des notifs depuis des scripts shell # ═══════════════════════════════════════════════════ @@ -122,7 +138,7 @@ sources: # ── DΓ©duplication ────────────────────────────────────── dedup: enabled: true - ttl_seconds: 300 # 5 minutes β€” durΓ©e du cache de dΓ©duplication + ttl_seconds: 300 # 5 minutes β€” durΓ©e du cache de dΓ©duplication rate_limit: - max_per_minute: 10 # max notifications/minute par topic - max_per_source: 30 # max notifications/minute toute source confondue + max_per_minute: 10 # max notifications/minute par topic + max_per_source: 30 # max notifications/minute toute source confondue diff --git a/ntfy.v b/ntfy.v new file mode 100644 index 0000000..4fa7c38 --- /dev/null +++ b/ntfy.v @@ -0,0 +1,76 @@ +module main + +import net.http +import time +import sources + +pub struct NtfyClient { +pub mut: + base_url string + auth_token string + logger Logger +} + +pub fn new_ntfy_client(cfg ServerConfig, logger Logger) NtfyClient { + return NtfyClient{ + base_url: cfg.url.trim_right('/') + auth_token: cfg.auth_token + logger: logger + } +} + +fn join_tags(tags []string) string { + return tags.join(',') +} + +// publish sends an event to the Ntfy server with exponential backoff retry (3 attempts). +pub fn (c NtfyClient) publish(event sources.Event) ! { + url := '${c.base_url}/${event.topic}' + mut header := http.Header{} + header.add_custom('Content-Type', 'text/plain')! + header.add_custom('X-Priority', event.priority.str())! + if event.tags.len > 0 { + header.add_custom('X-Tags', join_tags(event.tags))! + } + if c.auth_token != '' { + header.add(.authorization, 'Bearer ${c.auth_token}') + } + + c.logger.debug('ntfy', 'publish', + 'POST ${url} priority=${event.priority} tags=${join_tags(event.tags)}') + + // Retry with exponential backoff: 3 attempts, 1s / 2s / 4s delays + mut last_err := error('unknown') + for attempt := 1; attempt <= 3; attempt++ { + resp := http.fetch(http.FetchConfig{ + url: url + method: .post + header: header + data: event.message + }) or { + last_err = err + if attempt < 3 { + delay := time.second * (1 << (attempt - 1)) + c.logger.warn('ntfy', 'publish', + 'attempt ${attempt}/3 failed: ${err}, retrying in ${delay}') + time.sleep(delay) + } + continue + } + + if resp.status_code >= 200 && resp.status_code < 300 { + c.logger.debug('ntfy', 'publish', 'success ${resp.status_code}') + return + } + + last_err = error('ntfy returned ${resp.status_code}: ${resp.body}') + if attempt < 3 { + delay := time.second * (1 << (attempt - 1)) + c.logger.warn('ntfy', 'publish', + 'attempt ${attempt}/3 failed: ${last_err}, retrying in ${delay}') + time.sleep(delay) + } + } + + return last_err +} diff --git a/server.v b/server.v new file mode 100644 index 0000000..fb6be04 --- /dev/null +++ b/server.v @@ -0,0 +1,390 @@ +module main + +import net.http +import time +import json +import crypto.hmac +import crypto.sha256 +import encoding.hex +import sources +import os + + +pub struct Stats { +pub mut: + total_notifications int + errors int + by_source map[string]int +} + +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 +} + +pub struct RecentEvent { + timestamp time.Time + source string + name string + topic string + priority int + message string + success bool +} + +const recent_events_max = 100 + +pub fn new_app(cfg Config, 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{ + 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 + } +} + +fn load_dashboard_html(logger Logger) string { + // Try dashboard.html from current working directory + if os.exists('dashboard.html') { + content := os.read_file('dashboard.html') or { + logger.warn('server', 'dashboard', 'failed to read dashboard.html, using fallback') + return dashboard_fallback_html() + } + logger.info('server', 'dashboard', 'loaded dashboard.html (${content.len} bytes)') + return content + } + + logger.warn('server', 'dashboard', 'dashboard.html not found, using fallback') + return dashboard_fallback_html() +} + +fn dashboard_fallback_html() string { + return 'ntfy-bridge

πŸ”” ntfy-bridge

Dashboard file not found.

Place dashboard.html in the working directory.

' +} + +pub fn (mut app App) run() ! { + app.start_background_tasks() + + app.logger.info('server', 'start', 'listening on ${app.cfg.server.listen}') + + // Graceful shutdown via signal handling on Unix + $if linux || macos { + spawn app.handle_signals() + } + + mut s := http.Server{ + addr: app.cfg.server.listen + handler: app + } + s.listen_and_serve() +} + +$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) + } +} + +fn (mut app App) start_background_tasks() { + if app.cfg.sources.http_poll.len > 0 { + spawn app.http_poll_loop() + } + if app.cfg.sources.docker.len > 0 { + spawn app.docker_watch_loop() + } +} + +// ---- HTTP Handler ---- + +pub fn (mut app App) handle(req http.Request) http.Response { + match req.method { + .get { + if req.url == '/' || req.url == '/dashboard' { + return app.dashboard_response() + } + if req.url == '/health' { + return app.health_response() + } + if req.url == '/stats' { + return app.stats_response() + } + if req.url == '/api/config' { + return app.config_api_response() + } + if req.url == '/api/status' { + return app.status_api_response() + } + if req.url == '/api/history' { + return app.history_api_response() + } + } + .post { + return app.handle_webhook(req) + } + else {} + } + + return new_json_response(.not_found, '{"error":"not found"}') +} + +// ---- Health & Stats ---- + +fn (mut app App) health_response() http.Response { + uptime := time.now().unix() - app.start_time.unix() + body := '{"status":"ok","uptime_seconds":${uptime},"total_notifications":${app.stats.total_notifications}}' + return new_json_response(.ok, body) +} + +fn (mut app App) stats_response() http.Response { + mut parts := []string{} + for source, count in app.stats.by_source { + parts << '"${source}":${count}' + } + by_source := '{' + parts.join(',') + '}' + body := '{"total":${app.stats.total_notifications},"errors":${app.stats.errors},"by_source":${by_source}}' + return new_json_response(.ok, body) +} + +// ---- API endpoints for dashboard ---- + +fn (mut app App) config_api_response() http.Response { + mut safe := map[string]string{} + safe['server_url'] = app.cfg.server.url + safe['listen'] = app.cfg.server.listen + safe['dedup_enabled'] = app.cfg.dedup.enabled.str() + safe['dedup_ttl'] = app.cfg.dedup.ttl_seconds.str() + safe['default_priority'] = app.cfg.defaults.priority.str() + safe['gitea_sources'] = app.cfg.sources.gitea.len.str() + safe['uptime_kuma_sources'] = app.cfg.sources.uptime_kuma.len.str() + safe['docker_sources'] = app.cfg.sources.docker.len.str() + safe['http_poll_sources'] = app.cfg.sources.http_poll.len.str() + safe['cron_sources'] = app.cfg.sources.cron.len.str() + safe['generic_sources'] = app.cfg.sources.generic.len.str() + + return new_json_response(.ok, json.encode(safe)) +} + +fn (mut app App) status_api_response() http.Response { + mut items := []map[string]string{} + for src in app.cfg.sources.http_poll { + for check in src.checks { + is_down := app.poll_state[check.url] or { false } + status := if is_down { 'up' } else { 'down' } + items << { + 'name': check.url + 'status': status + } + } + } + return new_json_response(.ok, json.encode(items)) +} + +fn (mut app App) history_api_response() http.Response { + mut entries := []map[string]string{} + for ev in app.recent_events { + entries << { + 'timestamp': ev.timestamp.format_rfc3339() + 'source': ev.source + 'name': ev.name + 'topic': ev.topic + 'priority': ev.priority.str() + 'message': ev.message + 'success': ev.success.str() + } + } + return new_json_response(.ok, json.encode(entries)) +} + +// ---- Webhook handling (O(1) dispatch) ---- + +fn (mut app App) handle_webhook(req http.Request) http.Response { + path := req.url + route := app.webhook_routes[path] or { + return new_json_response(.not_found, '{"error":"unknown webhook"}') + } + + // HMAC validation (if secret is configured) + if app.cfg.server.hmac_secret != '' { + sig_header := req.header.get_custom('X-Hub-Signature-256') or { '' } + if sig_header == '' || !app.verify_hmac(req.data, sig_header) { + app.logger.warn('server', 'webhook', 'HMAC validation failed for ${path}') + return new_json_response(.unauthorized, '{"error":"invalid signature"}') + } + } + + mut event_opt := ?sources.Event(none) + match route.kind { + .gitea { + src := app.cfg.sources.gitea[route.index] + event_opt = sources.transform_gitea(req.data, src) + } + .uptime_kuma { + src := app.cfg.sources.uptime_kuma[route.index] + event_opt = sources.transform_uptime_kuma(req.data, src) + } + .cron { + src := app.cfg.sources.cron[route.index] + event_opt = sources.transform_cron(req.data, src) + } + .generic { + src := app.cfg.sources.generic[route.index] + event_opt = sources.transform_generic(req.data, src) + } + } + + event := event_opt or { + app.logger.warn('server', 'webhook', 'failed to transform: ${path}') + return new_json_response(.bad_request, '{"error":"bad request"}') + } + app.publish(event) + return new_json_response(.ok, '{"ok":true}') +} + +fn (mut app App) verify_hmac(payload string, sig_header string) bool { + if !sig_header.starts_with('sha256=') { + return false + } + expected_hex := sig_header[7..] // strip "sha256=" prefix + + key := app.cfg.server.hmac_secret.bytes() + data := payload.bytes() + computed := hmac.new(key, data, sha256.sum, sha256.block_size) + computed_hex := hex.encode(computed) + + return computed_hex == expected_hex +} + +// ---- Publish pipeline ---- + +pub fn (mut app App) publish(event sources.Event) { + mut e := apply_defaults(event, app.cfg.defaults) + 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 + } + 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) +} + +fn (mut app App) add_recent_event(ts time.Time, event sources.Event, success bool) { + app.recent_events << RecentEvent{ + timestamp: ts + source: event.source + name: event.name + topic: event.topic + priority: event.priority + message: event.message + success: success + } + if app.recent_events.len > recent_events_max { + app.recent_events = app.recent_events[app.recent_events.len - recent_events_max..] + } +} + +// ---- HTTP Poll Loop ---- + +fn (mut app App) http_poll_loop() { + for { + if app.shutdown_flag { + break + } + for src in app.cfg.sources.http_poll { + for check in src.checks { + app.run_http_check(check) + } + } + mut min_interval := 60 + for src in app.cfg.sources.http_poll { + if src.interval > 0 && src.interval < min_interval { + min_interval = src.interval + } + } + time.sleep(min_interval * time.second) + } +} + +fn (mut app App) run_http_check(check sources.HttpCheck) { + resp := http.fetch(http.FetchConfig{ + url: check.url + method: .get + read_timeout: check.timeout * time.second + }) or { + app.handle_http_result(check, false, err.msg()) + return + } + is_up := resp.status_code == check.expect_status + app.handle_http_result(check, is_up, 'status ${resp.status_code}') +} + +fn (mut app App) handle_http_result(check sources.HttpCheck, is_up bool, detail string) { + was_up := app.poll_state[check.url] or { true } + if is_up && !was_up { + app.publish(sources.transform_http_up(check)) + app.poll_state[check.url] = true + } else if !is_up && was_up { + app.publish(sources.transform_http_down(check, detail)) + app.poll_state[check.url] = false + } +} + +// ---- Helpers ---- + +fn new_json_response(status http.Status, body string) http.Response { + mut header := http.Header{} + header.add_custom('Content-Type', 'application/json') or { panic(err) } + return http.new_response(http.ResponseConfig{ + status: status + body: body + header: header + }) +} + +fn (mut app App) dashboard_response() http.Response { + mut header := http.Header{} + header.add_custom('Content-Type', 'text/html; charset=utf-8') or { panic(err) } + return http.new_response(http.ResponseConfig{ + status: .ok + body: app.dashboard_html + header: header + }) +} diff --git a/sources/cron.v b/sources/cron.v new file mode 100644 index 0000000..0c81926 --- /dev/null +++ b/sources/cron.v @@ -0,0 +1,13 @@ +module sources + +pub fn transform_cron(raw string, source CronSource) ?Event { + return Event{ + source: 'cron' + name: source.name + topic: source.topic + priority: source.priority_map['default'] or { 0 } + tags: source.tags.clone() + message: raw + raw: raw + } +} diff --git a/sources/docker.v b/sources/docker.v new file mode 100644 index 0000000..6c45fba --- /dev/null +++ b/sources/docker.v @@ -0,0 +1,68 @@ +module sources + +import x.json2 as json + +pub fn transform_docker_event(raw string, source DockerSource, host string) ?Event { + doc := json.decode[json.Any](raw, json.DecoderOptions{}) or { return none } + obj := doc.as_map() + + event_type := get_string(obj, 'status') + actor := obj['Actor'] or { json.Any(map[string]json.Any{}) }.as_map() + attrs := actor['Attributes'] or { json.Any(map[string]json.Any{}) }.as_map() + + container_name := get_string(attrs, 'name') + image := get_string(attrs, 'image') + exit_code := get_string(attrs, 'exitCode') + + priority := source.priority_map[event_type] or { 0 } + + mut tags := source.tags.clone() + if event_type == 'oom' && 'skull' !in tags { + tags << 'skull' + } + + if source.template != '' { + vars := { + 'container_name': container_name + 'status': event_type + 'image': image + 'exit_code': exit_code + 'host': host + } + return Event{ + source: 'docker' + name: source.name + topic: source.topic + priority: priority + tags: tags + message: render_source_template(source.template, vars) + raw: raw + } + } + + emoji := match event_type { + 'die' { 'πŸ’€' } + 'oom' { '☠️' } + 'health_status' { 'πŸ’“' } + else { '🐳' } + } + + mut message := '${emoji} ${container_name} β†’ ${event_type}' + if image != '' { + message += '\nImage: ${image}' + } + if exit_code != '' { + message += '\nExit code: ${exit_code}' + } + message += '\nHost: ${host}' + + return Event{ + source: 'docker' + name: source.name + topic: source.topic + priority: priority + tags: tags + message: message + raw: raw + } +} diff --git a/sources/event.v b/sources/event.v new file mode 100644 index 0000000..e392065 --- /dev/null +++ b/sources/event.v @@ -0,0 +1,111 @@ +module sources + +pub struct Event { +pub mut: + source string + name string + topic string + priority int + tags []string + message string + raw string +} + +pub fn (e Event) clone() Event { + return Event{ + source: e.source + name: e.name + topic: e.topic + priority: e.priority + tags: e.tags.clone() + message: e.message + raw: e.raw + } +} + +pub struct GiteaSource { +pub mut: + name string @[json: 'name'] + webhook_path string @[json: 'webhook_path'] + topic string @[json: 'topic'] + template string @[json: 'template'] + priority_map map[string]int @[json: 'priority_map'] + tags []string @[json: 'tags'] + repo string @[json: 'repo'] +} + +pub struct UptimeKumaState { +pub mut: + priority int @[json: 'priority'] + tags []string @[json: 'tags'] +} + +pub struct UptimeKumaSource { +pub mut: + name string @[json: 'name'] + webhook_path string @[json: 'webhook_path'] + topic string @[json: 'topic'] + template string @[json: 'template'] + priority_map map[string]int @[json: 'priority_map'] + tags []string @[json: 'tags'] + state_map map[string]UptimeKumaState @[json: 'state_map'] +} + +pub struct DockerSource { +pub mut: + name string @[json: 'name'] + topic string @[json: 'topic'] + template string @[json: 'template'] + priority_map map[string]int @[json: 'priority_map'] + tags []string @[json: 'tags'] + hosts []string @[json: 'hosts'] + events []string @[json: 'events'] +} + +pub struct HttpCheck { +pub mut: + url string @[json: 'url'] + topic string @[json: 'topic'] + priority int @[json: 'priority'] + expect_status int @[json: 'expect_status'] + timeout int @[json: 'timeout'] + template string @[json: 'template'] +} + +pub struct HttpPollSource { +pub mut: + name string @[json: 'name'] + interval int @[json: 'interval'] + checks []HttpCheck @[json: 'checks'] +} + +pub struct CronSource { +pub mut: + name string @[json: 'name'] + webhook_path string @[json: 'webhook_path'] + topic string @[json: 'topic'] + template string @[json: 'template'] + priority_map map[string]int @[json: 'priority_map'] + tags []string @[json: 'tags'] +} + +pub struct GenericSource { +pub mut: + name string @[json: 'name'] + webhook_path string @[json: 'webhook_path'] + topic string @[json: 'topic'] + template string @[json: 'template'] + priority_map map[string]int @[json: 'priority_map'] + tags []string @[json: 'tags'] +} + +// render_source_template applies variable substitution to a template string. +// Placeholders are {varname}. Shared by all source transformers in this module. +pub fn render_source_template(template string, vars map[string]string) string { + mut result := template.clone() + for key, value in vars { + placeholder := '{${key}}' + result = result.replace(placeholder, value) + } + return result +} diff --git a/sources/generic.v b/sources/generic.v new file mode 100644 index 0000000..a2622c7 --- /dev/null +++ b/sources/generic.v @@ -0,0 +1,13 @@ +module sources + +pub fn transform_generic(raw string, source GenericSource) ?Event { + return Event{ + source: 'generic' + name: source.name + topic: source.topic + priority: source.priority_map['default'] or { 0 } + tags: source.tags.clone() + message: raw + raw: raw + } +} diff --git a/sources/gitea.v b/sources/gitea.v new file mode 100644 index 0000000..e1c32f1 --- /dev/null +++ b/sources/gitea.v @@ -0,0 +1,114 @@ +module sources + +import x.json2 as json + +pub fn transform_gitea(raw string, source GiteaSource) ?Event { + doc := json.decode[json.Any](raw, json.DecoderOptions{}) or { return none } + obj := doc.as_map() + + event_type := detect_gitea_event(obj) + priority := source.priority_map[event_type] or { 0 } + + repo := get_string(obj, 'repository.full_name') + mut user := get_string(obj, 'pusher.username') + action := get_string(obj, 'action') + mut title := get_string(obj, 'commits.0.message') + mut sha := get_string(obj, 'commits.0.id') + + if event_type == 'pull_request' { + pr := obj['pull_request'] or { json.Any(map[string]json.Any{}) }.as_map() + title = get_string(pr, 'title') + sha = get_string(pr, 'merge_commit_sha') + if user == '' { + user = get_string(pr, 'user.login') + } + } else if event_type == 'issue' { + issue := obj['issue'] or { json.Any(map[string]json.Any{}) }.as_map() + title = get_string(issue, 'title') + sha = '' + if user == '' { + user = get_string(issue, 'user.login') + } + } + + template := source.template + mut message := '' + if template != '' { + vars := { + 'repo': repo + 'user': user + 'action': action + 'title': title + 'sha': sha + } + message = render_source_template(template, vars) + } else { + mut parts := []string{} + if repo != '' { + parts << '[${repo}]' + } + if user != '' { + parts << '${user}' + } + if action != '' { + parts << '${action}:' + } + if title != '' { + parts << '"${title}"' + } + if sha != '' { + short := if sha.len >= 7 { sha[..7] } else { sha } + parts << '(${short})' + } + emoji := match event_type { + 'pull_request' { 'πŸ”€' } + 'issue' { 'πŸ“' } + 'release' { 'πŸš€' } + else { 'πŸ”¨' } + } + + message = '${emoji} ' + parts.join(' ') + } + + return Event{ + source: 'gitea' + name: source.name + topic: source.topic + priority: priority + tags: source.tags.clone() + message: message + raw: raw + } +} + +fn detect_gitea_event(obj map[string]json.Any) string { + if 'commits' in obj { + return 'push' + } + if 'pull_request' in obj { + return 'pull_request' + } + if 'issue' in obj { + return 'issue' + } + if 'release' in obj { + return 'release' + } + return 'unknown' +} + +fn get_string(obj map[string]json.Any, path string) string { + parts := path.split('.') + mut current := obj.clone() + for i, part in parts { + if i == parts.len - 1 { + val := current[part] or { return '' } + return val.str() + } + next := current[part] or { return '' } + current = next.as_map().clone() + } + return '' +} + + diff --git a/sources/http_poll.v b/sources/http_poll.v new file mode 100644 index 0000000..f1edc07 --- /dev/null +++ b/sources/http_poll.v @@ -0,0 +1,37 @@ +module sources + +pub fn transform_http_down(check HttpCheck, err string) Event { + mut message := '' + if check.template != '' { + vars := { + 'url': check.url + 'error': err + 'status': '' + } + message = render_source_template(check.template, vars) + } else { + message = '❌ ${check.url} β†’ ${err}' + } + return Event{ + source: 'http_poll' + name: check.url + topic: check.topic + priority: check.priority + tags: ['x'] + message: message + raw: err + } +} + +pub fn transform_http_up(check HttpCheck) Event { + message := 'βœ… ${check.url} is UP' + return Event{ + source: 'http_poll' + name: check.url + topic: check.topic + priority: 1 + tags: ['white_check_mark'] + message: message + raw: '' + } +} diff --git a/sources/uptime_kuma.v b/sources/uptime_kuma.v new file mode 100644 index 0000000..b389463 --- /dev/null +++ b/sources/uptime_kuma.v @@ -0,0 +1,53 @@ +module sources + +import x.json2 as json + +pub fn transform_uptime_kuma(raw string, source UptimeKumaSource) ?Event { + doc := json.decode[json.Any](raw, json.DecoderOptions{}) or { return none } + obj := doc.as_map() + + monitor := get_string(obj, 'monitor.name') + status := get_string(obj, 'heartbeat.status') + msg := get_string(obj, 'heartbeat.msg') + ping := get_string(obj, 'heartbeat.ping') + time_str := get_string(obj, 'heartbeat.time') + + state_cfg := source.state_map[status] or { UptimeKumaState{} } + priority := state_cfg.priority + mut tags := state_cfg.tags.clone() + + if source.tags.len > 0 && tags.len == 0 { + tags << source.tags + } + + emoji := if status == 'up' { 'βœ…' } else { '🚨' } + mut message := '' + if source.template != '' { + vars := { + 'monitor': monitor + 'status': status + 'msg': msg + 'ping': ping + 'time': time_str + } + message = render_source_template(source.template, vars) + } else { + message = '${emoji} ${monitor} is ${status.to_upper()}' + if msg != '' { + message += ' β€” ${msg}' + } + if time_str != '' { + message += ' β€” ${time_str}' + } + } + + return Event{ + source: 'uptime_kuma' + name: source.name + topic: source.topic + priority: priority + tags: tags + message: message + raw: raw + } +} diff --git a/template.v b/template.v new file mode 100644 index 0000000..c367f81 --- /dev/null +++ b/template.v @@ -0,0 +1,15 @@ +module main + +import sources + +// apply_defaults fills in default priority and tags when the event doesn't set them. +pub fn apply_defaults(event sources.Event, defaults DefaultConfig) sources.Event { + mut e := event.clone() + if e.priority == 0 { + e.priority = defaults.priority + } + if e.tags.len == 0 { + e.tags = defaults.tags.clone() + } + return e +} diff --git a/test-duplicate.yaml b/test-duplicate.yaml new file mode 100644 index 0000000..51439e3 --- /dev/null +++ b/test-duplicate.yaml @@ -0,0 +1,20 @@ +server: + url: https://ntfy.example.com + listen: ":9090" +defaults: + priority: 3 +sources: + gitea: + - name: "test1" + webhook_path: /webhooks/same + topic: test + generic: + - name: "test2" + webhook_path: /webhooks/same + topic: test +dedup: + enabled: false + ttl_seconds: 300 + rate_limit: + max_per_minute: 10 + max_per_source: 30 diff --git a/uptime_kuma_test.v b/uptime_kuma_test.v new file mode 100644 index 0000000..afca277 --- /dev/null +++ b/uptime_kuma_test.v @@ -0,0 +1,41 @@ +module main + +import sources + +fn test_transform_uptime_kuma_down() { + raw := '{"monitor":{"name":"Gitea"},"heartbeat":{"status":"down","msg":"503","time":"14:32"}}' + src := sources.UptimeKumaSource{ + name: 'Critical' + webhook_path: '/webhooks/kuma' + topic: 'alerts' + state_map: { + 'down': sources.UptimeKumaState{ + priority: 5 + tags: ['rotating_light'] + } + } + } + event := sources.transform_uptime_kuma(raw, src) or { panic(err) } + assert event.source == 'uptime_kuma' + assert event.priority == 5 + assert event.message.contains('Gitea') + assert event.message.contains('DOWN') +} + +fn test_transform_uptime_kuma_up() { + raw := '{"monitor":{"name":"Gitea"},"heartbeat":{"status":"up"}}' + src := sources.UptimeKumaSource{ + name: 'Critical' + webhook_path: '/webhooks/kuma' + topic: 'alerts' + state_map: { + 'up': sources.UptimeKumaState{ + priority: 1 + tags: ['white_check_mark'] + } + } + } + event := sources.transform_uptime_kuma(raw, src) or { panic(err) } + assert event.priority == 1 + assert event.message.contains('UP') +} diff --git a/v.mod b/v.mod new file mode 100644 index 0000000..4e1ddb4 --- /dev/null +++ b/v.mod @@ -0,0 +1,8 @@ +Module { + name: 'ntfy-bridge' + description: 'Homelab notification hub β€” aggregates Gitea, Docker, Uptime Kuma β†’ Ntfy' + version: '0.1.0' + license: 'MIT' + dependencies: [] + subdirs: ['sources'] +} diff --git a/webhook.v b/webhook.v new file mode 100644 index 0000000..62742e7 --- /dev/null +++ b/webhook.v @@ -0,0 +1,47 @@ +module main + +// WebhookRoute maps a webhook path to its source type and index for O(1) dispatch. +pub enum SourceKind { + gitea + uptime_kuma + cron + generic +} + +pub struct WebhookRoute { + kind SourceKind + index int // index in the config's source array +} + +// build_webhook_routes constructs a map of webhook_path β†’ WebhookRoute +// for O(1) dispatch instead of scanning all sources on every request. +pub fn build_webhook_routes(cfg Config) map[string]WebhookRoute { + mut routes := map[string]WebhookRoute{} + + for i, src in cfg.sources.gitea { + routes[src.webhook_path] = WebhookRoute{ + kind: .gitea + index: i + } + } + for i, src in cfg.sources.uptime_kuma { + routes[src.webhook_path] = WebhookRoute{ + kind: .uptime_kuma + index: i + } + } + for i, src in cfg.sources.cron { + routes[src.webhook_path] = WebhookRoute{ + kind: .cron + index: i + } + } + for i, src in cfg.sources.generic { + routes[src.webhook_path] = WebhookRoute{ + kind: .generic + index: i + } + } + + return routes +}