CI / test (push) Successful in 2m30s
- metrics.v: GET /metrics au format OpenMetrics (gauges + counters) - server.v: route /metrics, Stats.http_polls counter, increment in poll loop - metrics_test.v: test du endpoint metrics - ROADMAP: Phase 5 6/8 (Prometheus done)
603 lines
17 KiB
V
603 lines
17 KiB
V
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()
|
||
}
|
||
// 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()
|
||
|
||
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)
|
||
|
||
// 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)
|
||
}
|
||
app.group_buffer = map[string][]GroupEntry{}
|
||
app.group_last_flush = time.now()
|
||
}
|