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 by_source map[string]int } 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') { 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() { // 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() } // 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() } } .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() 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) // 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.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) } app.group_buffer = map[string][]GroupEntry{} app.group_last_flush = time.now() }