Files
ntfy-bridge/server.v
T
bruno 1a4808f5f7
CI / test (push) Has been cancelled
feat: v0.6.0 — ACLs multi-utilisateurs + Plugin system
**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
2026-08-04 22:25:48 -04:00

803 lines
22 KiB
V
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
module main
import net.http
import time
import x.json2 as json
import crypto.hmac
import crypto.sha256
import encoding.hex
import sources
import os
import strconv
pub struct Stats {
pub mut:
total_notifications int
errors int
http_polls int // v0.5: total HTTP poll checks
by_source map[string]int
start_time time.Time
}
pub struct App {
pub mut:
cfg Config
config_path string // path to config file for SIGHUP reload
logger Logger
ntfy NtfyClient
dedup Dedup
stats Stats
start_time time.Time
poll_state map[string]bool // url -> is_down
webhook_routes map[string]WebhookRoute
recent_events []RecentEvent
shutdown_flag bool
dashboard_html string // cached dashboard HTML
silence_until time.Time // silence all notifications until this time
group_buffer map[string][]GroupEntry // key: source:name:topic
group_last_flush time.Time
}
pub struct RecentEvent {
timestamp time.Time
source string
name string
topic string
priority int
message string
success bool
}
pub struct GroupEntry {
event sources.Event
mut:
received time.Time
count int
}
const recent_events_max = 100
pub fn new_app(cfg Config, config_path string, quiet bool) App {
// json_mode = quiet — when running as a systemd service (quiet mode),
// use JSON structured logging; otherwise use human-readable console output
logger := new_logger(.info, quiet)
dashboard := load_dashboard_html(logger)
return App{
cfg: cfg
config_path: config_path
logger: logger
ntfy: new_ntfy_client(cfg.server, logger)
dedup: new_dedup(cfg.dedup)
stats: Stats{
by_source: map[string]int{}
}
start_time: time.now()
poll_state: map[string]bool{}
webhook_routes: build_webhook_routes(cfg)
recent_events: []RecentEvent{}
shutdown_flag: false
dashboard_html: dashboard
silence_until: time.Time{}
group_buffer: map[string][]GroupEntry{}
group_last_flush: time.now()
}
}
fn load_dashboard_html(logger Logger) string {
// Try dashboard.html from current working directory
if os.exists('dashboard.html') {
mut content := os.read_file('dashboard.html') or {
logger.warn('server', 'dashboard', 'failed to read dashboard.html, using fallback')
return dashboard_fallback_html()
}
content = content.replace('{{VERSION}}', 'v${version}')
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 '<!DOCTYPE html><html><head><title>ntfy-bridge</title><style>body{background:#0f172a;color:#e2e8f0;font-family:system-ui,sans-serif;display:flex;align-items:center;justify-content:center;min-height:100vh;margin:0}.msg{text-align:center}.msg h1{font-size:1.5rem}.msg p{color:#94a3b8;margin-top:8px}code{background:#1e293b;padding:2px 8px;border-radius:4px;font-size:.875rem}</style></head><body><div class=msg><h1>🔔 ntfy-bridge</h1><p>Dashboard file not found.</p><p>Place <code>dashboard.html</code> in the working directory.</p></div></body></html>'
}
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() {
// Register SIGHUP handler for config reload
os.signal_opt(.hup, fn [mut app] (sig os.Signal) {
app.logger.info('server', 'reload', 'SIGHUP received, reloading config...')
new_cfg := load_config(app.config_path) or {
app.logger.error('server', 'reload', 'failed to reload config: ${err}')
return
}
app.cfg = new_cfg
app.webhook_routes = build_webhook_routes(new_cfg)
app.ntfy = new_ntfy_client(new_cfg.server, app.logger)
app.dedup = new_dedup(new_cfg.dedup)
app.logger.info('server', 'reload', 'config reloaded successfully')
}) or { panic(err) }
// Common shutdown handler for SIGINT and SIGTERM
shutdown := fn [mut app] (sig os.Signal) {
app.logger.info('server', 'shutdown', 'received signal, shutting down...')
app.shutdown_flag = true
time.sleep(500 * time.millisecond)
exit(0)
}
os.signal_opt(.int, shutdown) or { panic(err) }
os.signal_opt(.term, shutdown) or { panic(err) }
// Keep goroutine alive
for {
time.sleep(60 * time.second)
}
}
}
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()
}
if app.cfg.sources.plugin.len > 0 {
spawn app.plugin_loop()
}
// Grouping flush goroutine
spawn app.group_flush_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()
}
if req.url == '/api/silence' {
return app.silence_status_response()
}
if req.url == '/metrics' {
return app.metrics_response()
}
}
.post {
if req.url.starts_with('/api/silence') {
return app.silence_create_response(req.url)
}
return app.handle_webhook(req)
}
.delete {
if req.url == '/api/silence' {
return app.silence_clear_response()
}
}
else {}
}
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()
safe['plugin_sources'] = app.cfg.sources.plugin.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"}')
}
// 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)
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)
}
.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 {
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)
// v0.5 — Apply filters
if !app.apply_filters(mut e) {
app.logger.debug('server', 'filter', 'dropped event ${e.source}:${e.name}')
return
}
// Check silence
now := time.now()
if now < app.silence_until {
app.logger.debug('server', 'publish', 'silenced: ${e.source}:${e.name}')
return
}
if !app.dedup.allow(e) {
app.logger.debug('server', 'dedup', 'dropped event ${e.source}:${e.name}')
return
}
// Grouping: buffer the event
key := '${e.source}:${e.name}:${e.topic}'
mut entries := app.group_buffer[key] or { []GroupEntry{} }
// Check if there's a matching event in the buffer (same source/name/topic)
mut found_idx := -1
for i, entry in entries {
if entry.event.message == e.message {
found_idx = i
break
}
}
if found_idx >= 0 {
// Rebuild to mutate count/received
mut updated := entries[found_idx]
updated.count++
updated.received = now
entries[found_idx] = updated
} else {
entries << GroupEntry{
event: e
received: now
count: 1
}
}
app.group_buffer[key] = entries
}
fn (mut app App) add_recent_event(ts time.Time, event sources.Event, success bool) {
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.stats.http_polls++
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
})
}
// ---- Silence Rules ----
fn (mut app App) silence_status_response() http.Response {
is_silenced := time.now() < app.silence_until
remaining := if is_silenced { app.silence_until.unix() - time.now().unix() } else { 0 }
body := '{"silenced":${is_silenced},"remaining_seconds":${remaining}}'
return new_json_response(.ok, body)
}
fn (mut app App) silence_create_response(url string) http.Response {
// Parse duration from query string: /api/silence?duration=30m
duration_str := parse_query_param(url, 'duration')
if duration_str == '' {
return new_json_response(.bad_request, '{"error":"missing duration parameter (e.g. ?duration=30m)"}')
}
seconds := parse_duration(duration_str)
if seconds <= 0 {
return new_json_response(.bad_request, '{"error":"invalid duration: ${duration_str}"}')
}
app.silence_until = time.now().add_seconds(seconds)
app.logger.info('server', 'silence', 'notifications silenced for ${duration_str} (${seconds}s)')
body := '{"silenced":true,"duration":"${duration_str}","seconds":${seconds}}'
return new_json_response(.ok, body)
}
fn (mut app App) silence_clear_response() http.Response {
app.silence_until = time.Time{}
app.logger.info('server', 'silence', 'silence cleared')
return new_json_response(.ok, '{"silenced":false}')
}
// parse_query_param extracts a query parameter value from a URL string.
fn parse_query_param(url string, key string) string {
prefix := '/api/silence?'
if !url.starts_with(prefix) {
return ''
}
qs := url[prefix.len..]
parts := qs.split('&')
for part in parts {
kv := part.split('=')
if kv.len == 2 && kv[0] == key {
return kv[1]
}
}
return ''
}
// parse_duration converts a human duration string (e.g. "30m", "2h", "1h30m") to seconds.
fn parse_duration(s string) int {
mut total := 0
mut current := ''
for c in s {
if c >= `0` && c <= `9` {
current += c.ascii_str()
} else {
val := strconv.atoi(current) or { 0 }
match c {
`s` { total += val }
`m` { total += val * 60 }
`h` { total += val * 3600 }
`d` { total += val * 86400 }
else {}
}
current = ''
}
}
return total
}
// ---- Grouping ----
const group_flush_interval = 10 // seconds
fn (mut app App) group_flush_loop() {
for {
if app.shutdown_flag {
break
}
time.sleep(group_flush_interval * time.second)
app.flush_group_buffer()
}
}
fn (mut app App) flush_group_buffer() {
now := time.now()
for _, entries in app.group_buffer {
if entries.len == 0 {
continue
}
mut total_count := 0
mut latest_entry := entries[0]
for i, entry in entries {
total_count += entry.count
if entry.received > latest_entry.received {
latest_entry = entries[i]
}
}
mut msg := latest_entry.event.message
if total_count > 1 {
msg = '${msg} (×${total_count} in ${group_flush_interval}s)'
}
mut e := latest_entry.event.clone()
e.message = msg
app.ntfy.publish(e) or {
app.logger.error('server', 'publish', 'failed to publish: ${err}')
app.stats.errors++
app.add_recent_event(now, e, false)
continue
}
app.stats.total_notifications++
mut count := app.stats.by_source[e.source] or { 0 }
count++
app.stats.by_source[e.source] = count
app.logger.info('server', 'publish',
'${e.source}:${e.name} → ${e.topic} (grouped ×${total_count})')
app.add_recent_event(now, e, true)
// v0.5 — Dispatch to outgoing webhooks
app.dispatch_outgoing(e)
// v0.6 — Save state after each flush
app.save_state()
}
app.group_buffer = map[string][]GroupEntry{}
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)
}
}