Files
ntfy-bridge/server.v
T
bruno 26341578fa
CI / test (push) Successful in 2m46s
fix: os.signal_opt API callback-based (V latest)
La nouvelle API est fn(os.Signal, fn(os.Signal)) au lieu de
chan+variadic. Enregistrement séparé pour .hup, .int, .term.
2026-08-03 15:00:39 -04:00

588 lines
16 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
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 '<!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()
}
}
.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()
}