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 '
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()
}
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)
}
}