Add YAML config loading with validation, env overrides for secrets, and a web dashboard exposing health, stats, and recent events.
391 lines
11 KiB
V
391 lines
11 KiB
V
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 '<!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() {
|
|
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
|
|
})
|
|
}
|