feat: add incoming webhook API support with configurable endpoints

- Implemented the incoming webhook API to handle alerts from external applications.
- Added configuration options for webhook listening address, port, and endpoints in config.go.
- Created WebhookReceiverConfig and WebhookEndpointConfig structures to manage webhook settings.
- Developed WebhookHandler to process incoming requests, validate tokens, and deliver alerts to specified channels.
- Enhanced existing controller interfaces to support alert delivery.
- Updated message rendering to respect Markdown settings for different channels.
- Added tests for webhook functionality and ensured proper error handling.
This commit is contained in:
2026-09-24 11:21:38 +08:00
parent 90ca08066b
commit 057d7d584d
17 changed files with 595 additions and 65 deletions
+61 -6
View File
@@ -11,8 +11,9 @@ Nukumizu connects to a Komari Dashboard instance, keeps an in-memory view of eve
- **Remote command execution** — dispatches commands through the Komari task API and polls the result (1s interval, up to 60s timeout). - **Remote command execution** — dispatches commands through the Komari task API and polls the result (1s interval, up to 60s timeout).
- **Interactive bots** — QQ (NapCat / OneBot 11) and Telegram bots for `/list`, `/status`, `/info`, `/run`, `/shutdown`, `/reboot`, and more, protected by an admin / trusted-group permission model. - **Interactive bots** — QQ (NapCat / OneBot 11) and Telegram bots for `/list`, `/status`, `/info`, `/run`, `/shutdown`, `/reboot`, and more, protected by an admin / trusted-group permission model.
- **Notification channels** — server status changes are pushed to every enabled channel: QQ, Telegram, Email (SMTP), [ntfy](https://ntfy.sh), and Webhook. - **Notification channels** — server status changes are pushed to every enabled channel: QQ, Telegram, Email (SMTP), [ntfy](https://ntfy.sh), and Webhook.
- **Incoming webhook API** — external applications can push their own alerts in via `POST /api/webhook/<name>`, and Nukumizu relays them to the channels that endpoint lists. Each endpoint carries its own token and target channels, and the API is served on a **separate listener** so it can be exposed without exposing the admin API.
- **Network proxy** — a global proxy URL can be enabled per controller (`networkUseProxy`) for HTTP, WebSocket, and even SMTP (HTTP CONNECT tunnel). - **Network proxy** — a global proxy URL can be enabled per controller (`networkUseProxy`) for HTTP, WebSocket, and even SMTP (HTTP CONNECT tunnel).
- **Customizable message templates** — every bot/notification message is rendered from a template in `config.json`. - **Customizable message templates** — every bot/notification message is rendered from a template in `config.json`, with Markdown formatting switched on per channel.
- **Storage** — SQLite (pure-Go driver) for `user.db` and `log.db`; safe on network shares (WAL disabled). - **Storage** — SQLite (pure-Go driver) for `user.db` and `log.db`; safe on network shares (WAL disabled).
- **Dashboard API** — token-authenticated REST API plus an admin-only live log-streaming WebSocket. - **Dashboard API** — token-authenticated REST API plus an admin-only live log-streaming WebSocket.
- **Web console** — a Vue 3 admin UI for browsing nodes, editing `config.json`, managing bot trust, and tailing logs. The built bundle is embedded in the binary, so a single executable serves both the API and the console. - **Web console** — a Vue 3 admin UI for browsing nodes, editing `config.json`, managing bot trust, and tailing logs. The built bundle is embedded in the binary, so a single executable serves both the API and the console.
@@ -30,7 +31,7 @@ Nukumizu connects to a Komari Dashboard instance, keeps an in-memory view of eve
``` ```
nukumizu-backend/ nukumizu-backend/
├── main.go # Entry point, startup sequence, graceful shutdown ├── main.go # Entry point, startup sequence, graceful shutdown
├── router.go # HTTP route registration ├── router.go # HTTP route registration (main + webhook API)
├── config/ ├── config/
│ ├── config.go # Load config files, apply defaults │ ├── config.go # Load config files, apply defaults
│ └── variables.go # Config schema structs + globals │ └── variables.go # Config schema structs + globals
@@ -40,6 +41,7 @@ nukumizu-backend/
│ ├── user.go # /api/user/login, /api/user/register │ ├── user.go # /api/user/login, /api/user/register
│ ├── server.go # /api/server/list, getInfo, getStatus, exec │ ├── server.go # /api/server/list, getInfo, getStatus, exec
│ ├── settings.go # /api/settings/get, set │ ├── settings.go # /api/settings/get, set
│ ├── webhook.go # /api/webhook/{name} (incoming webhook API)
│ └── health.go # /health │ └── health.go # /health
├── database/ ├── database/
│ └── user.go # user.db (SQLite) user store │ └── user.go # user.db (SQLite) user store
@@ -66,14 +68,14 @@ nukumizu-backend/
│ ├── template/ │ ├── template/
│ │ └── template.go # Message template renderer ({{ variables }}) │ │ └── template.go # Message template renderer ({{ variables }})
│ └── controller/ │ └── controller/
│ ├── controller.go # Manager, Controller / BotController interfaces │ ├── controller.go # Manager, Controller / BotController interfaces, alerts
│ ├── trigger.go # Command parsing, authorization, routing │ ├── trigger.go # Command parsing, authorization, routing
│ ├── processor.go # Command handlers │ ├── processor.go # Command handlers
│ ├── utils.go │ ├── utils.go
│ └── pipes/ │ └── pipes/
│ ├── email.go # Email notification pipe │ ├── email.go # Email notification pipe
│ ├── ntfy.go # ntfy notification pipe │ ├── ntfy.go # ntfy notification pipe
│ ├── webhook.go # Webhook notification pipe │ ├── webhook.go # Outgoing webhook notification pipe
│ ├── qq_napcat/ │ ├── qq_napcat/
│ │ ├── qq.go # QQ (NapCat / OneBot 11) bot controller │ │ ├── qq.go # QQ (NapCat / OneBot 11) bot controller
│ │ └── napcat.go # NapCat WebSocket + HTTP API client │ │ └── napcat.go # NapCat WebSocket + HTTP API client
@@ -139,9 +141,22 @@ There are two configuration files, both read from the working directory unless o
"password": "CHANGE_ME" "password": "CHANGE_ME"
} }
}, },
"webhook": {
"enabled": false,
"listenAddr": "0.0.0.0",
"listenPort": "8081",
"endpoints": {
"example": {
"enabled": true,
"token": "CHANGE_ME",
"notifyPipes": ["qq(napcat)", "telegram", "email", "ntfy"]
}
}
},
"controllerMethod": { "controllerMethod": {
"qq(napcat)": { "qq(napcat)": {
"enabled": false, "enabled": false,
"markdown": false,
"networkUseProxy": false, "networkUseProxy": false,
"napcatAddr": "127.0.0.1", "napcatAddr": "127.0.0.1",
"napcatPort": "3000", "napcatPort": "3000",
@@ -151,12 +166,14 @@ There are two configuration files, both read from the working directory unless o
}, },
"telegram": { "telegram": {
"enabled": false, "enabled": false,
"markdown": true,
"networkUseProxy": false, "networkUseProxy": false,
"botToken": "", "botToken": "",
"listenMethod": "global" "listenMethod": "global"
}, },
"email": { "email": {
"enabled": false, "enabled": false,
"markdown": false,
"networkUseProxy": false, "networkUseProxy": false,
"smtpHost": "", "smtpHost": "",
"smtpPort": 587, "smtpPort": 587,
@@ -168,6 +185,7 @@ There are two configuration files, both read from the working directory unless o
}, },
"ntfy": { "ntfy": {
"enabled": false, "enabled": false,
"markdown": false,
"networkUseProxy": false, "networkUseProxy": false,
"server": "https://ntfy.sh", "server": "https://ntfy.sh",
"topic": "", "topic": "",
@@ -176,6 +194,7 @@ There are two configuration files, both read from the working directory unless o
}, },
"webhook": { "webhook": {
"enabled": false, "enabled": false,
"markdown": false,
"networkUseProxy": false, "networkUseProxy": false,
"url": "", "url": "",
"method": "POST", "method": "POST",
@@ -199,11 +218,13 @@ There are two configuration files, both read from the working directory unless o
Field notes: Field notes:
- `system.networkProxy` is a **system-wide** proxy URL. A controller only uses it when its own `networkUseProxy` is `true`. Applied to Telegram HTTP polling, NapCat HTTP/WebSocket, ntfy and webhook requests, and Email SMTP (tunneled via HTTP CONNECT). - `system.networkProxy` is a **system-wide** proxy URL. A controller only uses it when its own `networkUseProxy` is `true`. Applied to Telegram HTTP polling, NapCat HTTP/WebSocket, ntfy and webhook requests, and Email SMTP (tunneled via HTTP CONNECT).
- `webhook` configures the **incoming** webhook API (see [Incoming webhook API](#incoming-webhook-api)); `controllerMethod.webhook` configures the outgoing webhook notification channel. They are independent.
- `markdown` is a per-channel switch on all five channels. With it `false` (the default) every rendered value is inserted as plain text; with it `true` the values meant to be read verbatim (UUIDs, event messages, commands, command results, alert source and alert content) are wrapped in Markdown code spans / fenced blocks. Nothing is inferred from the channel name, so a channel only ever gets the formatting you asked for — turn it off for a channel whose platform does not render Markdown. On Telegram it also picks the `parse_mode`: with `markdown` off, messages are sent without one, so text containing `*` or `_` is delivered as-is rather than rejected by the API as malformed Markdown.
- `controllerMethod.qq(napcat).listenMethod` / `telegram.listenMethod` — see [Bot recognition modes](#bot-recognition-modes). - `controllerMethod.qq(napcat).listenMethod` / `telegram.listenMethod` — see [Bot recognition modes](#bot-recognition-modes).
- `debug` toggles verbose per-channel message/action logging; these only matter in debug builds / `debugMode`. - `debug` toggles verbose per-channel message/action logging; these only matter in debug builds / `debugMode`.
- `email.useTLS` is kept for configuration compatibility. - `email.useTLS` is kept for configuration compatibility.
- `dataPath` / `dbPath` default to `./data` and `./db`; `user.db` and `log.db` are created under `dbPath`. - `dataPath` / `dbPath` default to `./data` and `./db`; `user.db` and `log.db` are created under `dbPath`.
- Missing keys fall back to built-in defaults (host `0.0.0.0`, port `8080`, NapCat `127.0.0.1:3000`, ntfy server `https://ntfy.sh`, webhook method `POST`, etc.). Message templates have built-in fallbacks too. - Missing keys fall back to built-in defaults (host `0.0.0.0`, port `8080`, webhook API `0.0.0.0:8081`, no webhook endpoints, NapCat `127.0.0.1:3000`, ntfy server `https://ntfy.sh`, webhook method `POST`, etc.). Message templates have built-in fallbacks too. `markdown` defaults to `false`, so add it explicitly for Telegram (see the sample above) to keep its formatting.
### `bot_user_config.json` ### `bot_user_config.json`
@@ -256,7 +277,7 @@ Admins and trusted groups are defined **per bot channel** and map a member ID to
### Message templates ### Message templates
`controllerMessage` templates are rendered before sending. Available variables (rendered through the Telegram pipe are additionally wrapped in Telegram legacy Markdown): `controllerMessage` templates are rendered before sending. Available variables (channels with `markdown: true` additionally wrap the verbatim values in Markdown — see the field notes above):
| Variable | Meaning | | Variable | Meaning |
|---|---| |---|---|
@@ -311,6 +332,38 @@ Middleware applied to the whole server:
- **Security headers** — `X-XSS-Protection`, `X-Content-Type-Options: nosniff`, `X-Frame-Options: DENY`, `Referrer-Policy`, a restrictive CSP. - **Security headers** — `X-XSS-Protection`, `X-Content-Type-Options: nosniff`, `X-Frame-Options: DENY`, `Referrer-Policy`, a restrictive CSP.
- **WebSocket auth** — `utils.WebSocketAuthMiddleware` is attached to `/api/system/getLogs` (route-level, not global): it authenticates the upgrade request and requires an `admin` token before the connection is handed to the log handler. - **WebSocket auth** — `utils.WebSocketAuthMiddleware` is attached to `/api/system/getLogs` (route-level, not global): it authenticates the upgrade request and requires an `admin` token before the connection is handed to the log handler.
### Incoming webhook API
A listener of its own, so external applications can be pointed at it without being able to reach the admin API. It is switched on with `webhook.enabled` and binds `webhook.listenAddr:webhook.listenPort` (default `0.0.0.0:8081`); that half of the configuration is applied at startup, while `webhook.endpoints` is re-read whenever the config is reloaded. Only the rate limit and CORS middleware apply here — no session token is involved.
| Endpoint | Method | Permission | Description |
|---|---|---|---|
| `/api/webhook/<name>` | POST | Endpoint token | Relay an alert to the channels the endpoint lists in `notifyPipes`. Body `{token, subject, content}`. Returns `data: {endpoint, channels}`. |
Every entry under `webhook.endpoints` is one endpoint, addressed by its key as the last path segment: the key `example` is served at `POST /api/webhook/example`. An endpoint holds:
| Field | Meaning |
|---|---|
| `enabled` | Whether the endpoint accepts requests. A disabled endpoint answers `403`. |
| `token` | Shared secret the caller sends as the `token` body field; compared in constant time. An endpoint with an empty token answers `500` instead of accepting requests from anyone. |
| `notifyPipes` | The channels the alert is delivered to, named as in `controllerMethod`: `qq(napcat)`, `telegram`, `email`, `ntfy`, `webhook`. A channel that is unknown or disabled is skipped and reported. |
The alert is rendered per channel as:
```
{{ subject }}
- Source: {{ source }}
- Content:
{{ content }}
- Time: {{ time }}
Sent by Nukumizu Alert System
```
`{{ source }}` is the endpoint name, so recipients can tell which application triggered the alert. On a channel with `markdown: true` the source is wrapped in inline code and the content in a fenced code block; `{{ subject }}` and `{{ time }}` stay plain.
Status codes: `200` delivered, `400` malformed body or empty `subject`/`content`, `401` wrong token, `403` endpoint disabled, `404` unknown endpoint name, `405` non-POST request, `500` endpoint has no token configured, `502` no channel accepted the alert.
## Bots ## Bots
QQ (NapCat) and Telegram bots share one command engine and authorization pipeline, implemented in `internal/controller/`. NapCat speaks OneBot 11 (WebSocket event stream + HTTP actions); Telegram uses `go-telegram/bot` long polling. QQ (NapCat) and Telegram bots share one command engine and authorization pipeline, implemented in `internal/controller/`. NapCat speaks OneBot 11 (WebSocket event stream + HTTP actions); Telegram uses `go-telegram/bot` long polling.
@@ -344,6 +397,8 @@ QQ (NapCat) and Telegram bots share one command engine and authorization pipelin
QQ and Telegram are *interactive* channels. Email, ntfy, and webhook are **status-only** channels — they receive server status-change alerts but cannot run commands. On startup, the welcome message and initial server list are delivered only to the bot channels (QQ / Telegram), honoring each member's `event_bot_started` preference. QQ and Telegram are *interactive* channels. Email, ntfy, and webhook are **status-only** channels — they receive server status-change alerts but cannot run commands. On startup, the welcome message and initial server list are delivered only to the bot channels (QQ / Telegram), honoring each member's `event_bot_started` preference.
All five channels can also carry an alert submitted by an external application through the [incoming webhook API](#incoming-webhook-api). A bot channel delivers it to the groups and admins configured for that channel; a status-only channel delivers it to its configured destination (mail recipients, ntfy topic, outgoing webhook URL). Markdown formatting is decided per channel by its `markdown` setting, never by the channel's name.
## Building ## Building
Requires Go 1.25+ and — to build the web console — Node.js 22+. Requires Go 1.25+ and — to build the web console — Node.js 22+.
+11
View File
@@ -107,6 +107,17 @@ func LoadGlobalConfig(configPath string) (*Config, error) {
cfg.ControllerMethod.Webhook.Headers = map[string]string{} cfg.ControllerMethod.Webhook.Headers = map[string]string{}
} }
// Apply defaults for the incoming webhook API.
if cfg.Webhook.ListenAddr == "" {
cfg.Webhook.ListenAddr = "0.0.0.0"
}
if cfg.Webhook.ListenPort == "" {
cfg.Webhook.ListenPort = "8081"
}
if cfg.Webhook.Endpoints == nil {
cfg.Webhook.Endpoints = map[string]WebhookEndpointConfig{}
}
// Apply defaults for paths. // Apply defaults for paths.
if cfg.DataPath == "" { if cfg.DataPath == "" {
cfg.DataPath = "./data" cfg.DataPath = "./data"
+52 -1
View File
@@ -34,6 +34,7 @@ type KomariConfig struct {
// QQConfig holds QQ (Napcat) Bot controller configuration. // QQConfig holds QQ (Napcat) Bot controller configuration.
type QQConfig struct { type QQConfig struct {
Markdown bool `json:"markdown"`
Enabled bool `json:"enabled"` Enabled bool `json:"enabled"`
NetworkUseProxy bool `json:"networkUseProxy"` NetworkUseProxy bool `json:"networkUseProxy"`
NapcatAddr string `json:"napcatAddr"` NapcatAddr string `json:"napcatAddr"`
@@ -45,6 +46,7 @@ type QQConfig struct {
// TelegramConfig holds Telegram Bot controller configuration. // TelegramConfig holds Telegram Bot controller configuration.
type TelegramConfig struct { type TelegramConfig struct {
Markdown bool `json:"markdown"`
Enabled bool `json:"enabled"` Enabled bool `json:"enabled"`
NetworkUseProxy bool `json:"networkUseProxy"` NetworkUseProxy bool `json:"networkUseProxy"`
BotToken string `json:"botToken"` BotToken string `json:"botToken"`
@@ -53,6 +55,7 @@ type TelegramConfig struct {
// EmailConfig holds Email notification controller configuration. // EmailConfig holds Email notification controller configuration.
type EmailConfig struct { type EmailConfig struct {
Markdown bool `json:"markdown"`
Enabled bool `json:"enabled"` Enabled bool `json:"enabled"`
NetworkUseProxy bool `json:"networkUseProxy"` NetworkUseProxy bool `json:"networkUseProxy"`
SMTPHost string `json:"smtpHost"` SMTPHost string `json:"smtpHost"`
@@ -66,6 +69,7 @@ type EmailConfig struct {
// NtfyConfig holds Ntfy notification controller configuration. // NtfyConfig holds Ntfy notification controller configuration.
type NtfyConfig struct { type NtfyConfig struct {
Markdown bool `json:"markdown"`
Enabled bool `json:"enabled"` Enabled bool `json:"enabled"`
NetworkUseProxy bool `json:"networkUseProxy"` NetworkUseProxy bool `json:"networkUseProxy"`
Server string `json:"server"` Server string `json:"server"`
@@ -74,8 +78,11 @@ type NtfyConfig struct {
Priority string `json:"priority"` Priority string `json:"priority"`
} }
// WebhookConfig holds Webhook notification controller configuration. // WebhookConfig holds the outgoing Webhook notification controller
// configuration. It is the counterpart of WebhookReceiverConfig, which serves
// the incoming webhook API.
type WebhookConfig struct { type WebhookConfig struct {
Markdown bool `json:"markdown"`
Enabled bool `json:"enabled"` Enabled bool `json:"enabled"`
NetworkUseProxy bool `json:"networkUseProxy"` NetworkUseProxy bool `json:"networkUseProxy"`
URL string `json:"url"` URL string `json:"url"`
@@ -93,6 +100,49 @@ type ControllerMethodConfig struct {
Webhook WebhookConfig `json:"webhook"` Webhook WebhookConfig `json:"webhook"`
} }
// WebhookEndpointConfig holds a single incoming webhook endpoint. Endpoints are
// keyed by name under webhook.endpoints; the name is the last path segment of
// the endpoint's URL, so an endpoint named "example" is served at
// POST /api/webhook/example. One endpoint per external application and target
// channel group keeps their tokens and recipients apart.
type WebhookEndpointConfig struct {
// Enabled controls whether the endpoint accepts requests. A disabled
// endpoint answers with 403.
Enabled bool `json:"enabled"`
// Token is the shared secret the caller must send in the request body. An
// endpoint without a token is rejected: an empty token would make the
// endpoint an open relay, so it is treated as a configuration error.
Token string `json:"token"`
// NotifyPipes lists the notification channels the alert is delivered to, by
// controller name (e.g. "qq(napcat)", "telegram", "email", "ntfy",
// "webhook").
NotifyPipes []string `json:"notifyPipes"`
}
// WebhookReceiverConfig holds the incoming webhook API settings. The API is
// served on its own listener instead of the main one, so external applications
// can be given access to the webhook port without exposing the admin API. Only
// the endpoints map is re-read on a settings update; enabled, listenAddr and
// listenPort are applied at startup.
type WebhookReceiverConfig struct {
Enabled bool `json:"enabled"`
ListenAddr string `json:"listenAddr"`
ListenPort string `json:"listenPort"`
Endpoints map[string]WebhookEndpointConfig `json:"endpoints"`
}
// GetWebhookEndpoint returns the incoming webhook endpoint registered under the
// given name, and whether such an endpoint exists.
func GetWebhookEndpoint(name string) (WebhookEndpointConfig, bool) {
if C_globalConfig == nil {
return WebhookEndpointConfig{}, false
}
endpoint, ok := C_globalConfig.Webhook.Endpoints[name]
return endpoint, ok
}
// ControllerMessageConfig holds message templates for controller responses. // ControllerMessageConfig holds message templates for controller responses.
type ControllerMessageConfig struct { type ControllerMessageConfig struct {
BotStarted string `json:"BOT_STARTED"` BotStarted string `json:"BOT_STARTED"`
@@ -108,6 +158,7 @@ type Config struct {
System SystemConfig `json:"system"` System SystemConfig `json:"system"`
Debug DebugConfig `json:"debug"` Debug DebugConfig `json:"debug"`
Komari KomariConfig `json:"komari"` Komari KomariConfig `json:"komari"`
Webhook WebhookReceiverConfig `json:"webhook"`
ControllerMethod ControllerMethodConfig `json:"controllerMethod"` ControllerMethod ControllerMethodConfig `json:"controllerMethod"`
ControllerMessage ControllerMessageConfig `json:"controllerMessage"` ControllerMessage ControllerMessageConfig `json:"controllerMessage"`
DataPath string `json:"dataPath"` DataPath string `json:"dataPath"`
+1 -1
View File
@@ -63,7 +63,7 @@ func ServerListHandler(w http.ResponseWriter, r *http.Request) {
} }
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
result := template.Render("", params) result := template.Render("", params, false)
utils.SendSuccessResponse(w, "", map[string]interface{}{ utils.SendSuccessResponse(w, "", map[string]interface{}{
"list": result, "list": result,
+104
View File
@@ -0,0 +1,104 @@
package handler
import (
"crypto/subtle"
"encoding/json"
"net/http"
"time"
"nukumizu-backend/config"
"nukumizu-backend/internal/controller"
"nukumizu-backend/postLog"
"nukumizu-backend/utils"
)
// maxWebhookBodyBytes caps the size of an incoming webhook request body. The
// endpoint is reachable without a session token, so the body is bounded before
// it is read.
const maxWebhookBodyBytes = 1 << 20 // 1 MiB
// webhookRequest is the JSON body accepted by the incoming webhook API.
type webhookRequest struct {
Token string `json:"token"`
Subject string `json:"subject"`
Content string `json:"content"`
}
// WebhookHandler handles POST /api/webhook/{name}, the incoming webhook API
// served on its own listener (see webhook in config.json). The path segment
// selects the endpoint, which carries the token to present and the notification
// channels to deliver to:
//
// POST /api/webhook/example
// {"token": "...", "subject": "...", "content": "..."}
//
// The alert is rendered per channel and sent through every channel the endpoint
// lists in notifyPipes. This route is not part of the token-authenticated API:
// it authenticates with the endpoint's own shared token.
func WebhookHandler(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
utils.SendErrorResponse(w, http.StatusMethodNotAllowed, "method not allowed")
return
}
name := r.PathValue("name")
endpoint, exists := config.GetWebhookEndpoint(name)
if !exists {
utils.SendErrorResponse(w, http.StatusNotFound, "unknown webhook endpoint: "+name)
return
}
if !endpoint.Enabled {
utils.SendErrorResponse(w, http.StatusForbidden, "webhook endpoint is disabled: "+name)
return
}
// An endpoint without a token would accept requests from anyone, so it is
// treated as a misconfiguration rather than as an open endpoint.
if endpoint.Token == "" {
postLog.Error("Webhook endpoint " + name + " has no token configured, rejecting request")
utils.SendErrorResponse(w, http.StatusInternalServerError, "webhook endpoint is not configured with a token: "+name)
return
}
var req webhookRequest
if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, maxWebhookBodyBytes)).Decode(&req); err != nil {
utils.SendErrorResponse(w, http.StatusBadRequest, "invalid request body: expected a JSON object with token, subject and content")
return
}
if subtle.ConstantTimeCompare([]byte(req.Token), []byte(endpoint.Token)) != 1 {
postLog.Warning("Webhook request rejected for endpoint " + name + ": invalid token")
utils.SendErrorResponse(w, http.StatusUnauthorized, "invalid token")
return
}
if req.Subject == "" || req.Content == "" {
utils.SendErrorResponse(w, http.StatusBadRequest, "missing required parameter: subject and content must not be empty")
return
}
manager := controller.GetManager()
if manager == nil {
utils.SendErrorResponse(w, http.StatusInternalServerError, "controller manager not initialized")
return
}
alert := controller.Alert{
Subject: req.Subject,
Source: name,
Content: req.Content,
Time: time.Now().Format("2006-01-02T15:04:05.000000000-07:00"),
}
delivered, err := manager.NotifyAlert(endpoint.NotifyPipes, alert)
if err != nil {
postLog.Error("Failed to deliver webhook alert for endpoint " + name + ": " + err.Error())
utils.SendErrorResponse(w, http.StatusBadGateway, "failed to send alert: "+err.Error())
return
}
postLog.Info("Webhook alert delivered for endpoint " + name)
utils.SendSuccessResponse(w, "alert sent", map[string]interface{}{
"endpoint": name,
"channels": delivered,
})
}
+93 -7
View File
@@ -1,7 +1,9 @@
package controller package controller
import ( import (
"errors"
"fmt" "fmt"
"strings"
"sync" "sync"
"nukumizu-backend/config" "nukumizu-backend/config"
@@ -12,7 +14,7 @@ import (
// Command represents a parsed bot command. // Command represents a parsed bot command.
type Command struct { type Command struct {
Source string // The source pipe (e.g., "telegram", "qq", "napcat") Source string // Name of the pipe the command arrived on (see Controller.Name)
RawText string // The raw text of the command message RawText string // The raw text of the command message
Command string // The command word (e.g., "list", "status") Command string // The command word (e.g., "list", "status")
Args []string // Command arguments Args []string // Command arguments
@@ -38,8 +40,33 @@ const (
// MessageTypeReply marks a direct reply to a user command. Reserved for the // MessageTypeReply marks a direct reply to a user command. Reserved for the
// BotUserOptions.EventReply opt-out. // BotUserOptions.EventReply opt-out.
MessageTypeReply = "event_reply" MessageTypeReply = "event_reply"
// MessageTypeAlert marks an alert submitted by an external application
// through the incoming webhook API. Not member-controllable: an alert is
// always delivered to the channel's recipients.
MessageTypeAlert = "alert"
) )
// Alert is a free-form notification submitted by an external application
// through the incoming webhook API. Its target channels are chosen per webhook
// endpoint in config.json, not per alert.
type Alert struct {
Subject string // Short one-line title of the alert
Source string // Name of the webhook endpoint the alert was submitted to
Content string // Free-form alert body
Time string // Submission time
}
// Render renders the alert body for a channel, wrapping the source and content
// in Markdown when that channel has markdown enabled (see template.RenderAlert).
func (a Alert) Render(markdown bool) string {
return template.RenderAlert(template.AlertParams{
Subject: a.Subject,
Source: a.Source,
Content: a.Content,
Time: a.Time,
}, markdown)
}
// MemberReceives reports whether a member whose bot_user_config.json options are // MemberReceives reports whether a member whose bot_user_config.json options are
// opts receives an automatic message of the given type. Only member-controllable // opts receives an automatic message of the given type. Only member-controllable
// types are gated; anything else is always delivered. // types are gated; anything else is always delivered.
@@ -58,9 +85,15 @@ type Controller interface {
Start() error Start() error
Stop() Stop()
IsEnabled() bool IsEnabled() bool
// IsMarkdown reports whether the channel renders Markdown, per its own
// "markdown" setting in config.json.
IsMarkdown() bool
SendStatusChange(change node.StatusChange) error SendStatusChange(change node.StatusChange) error
SendServerList(onlineServers, offlineServers string) error SendServerList(onlineServers, offlineServers string) error
SendExecuteResult(serverName, serverUUID, command, result string) error SendExecuteResult(serverName, serverUUID, command, result string) error
// SendAlert delivers a free-form alert submitted through the incoming
// webhook API to the channel's own recipients.
SendAlert(alert Alert) error
} }
// BotController is implemented by controllers that act as chat bots and can // BotController is implemented by controllers that act as chat bots and can
@@ -105,14 +138,14 @@ func (m *Manager) Register(c Controller) {
// bot controllers (QQ/NapCat and Telegram). Notification-only pipes that do // bot controllers (QQ/NapCat and Telegram). Notification-only pipes that do
// not implement BotController are skipped. The message is typed // not implement BotController are skipped. The message is typed
// MessageTypeBotStarted so each controller can honor its members' per-recipient // MessageTypeBotStarted so each controller can honor its members' per-recipient
// EventBotStarted opt-out. // EventBotStarted opt-out. It is rendered once per controller because the
// Markdown formatting depends on each channel's own markdown setting.
func (m *Manager) ShowBotInitMessage() { func (m *Manager) ShowBotInitMessage() {
m.mu.RLock() m.mu.RLock()
defer m.mu.RUnlock() defer m.mu.RUnlock()
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildBotInitializationMsgParams() params := template.BuildBotInitializationMsgParams()
content := template.Render(cfg.ControllerMessage.BotStarted, params)
for _, ctrl := range m.controllers { for _, ctrl := range m.controllers {
if !ctrl.IsEnabled() { if !ctrl.IsEnabled() {
@@ -124,7 +157,7 @@ func (m *Manager) ShowBotInitMessage() {
} }
message := Message{ message := Message{
Source: bot.Name(), Source: bot.Name(),
Content: content, Content: template.Render(cfg.ControllerMessage.BotStarted, params, ctrl.IsMarkdown()),
Type: MessageTypeBotStarted, Type: MessageTypeBotStarted,
} }
if err := bot.SendMessage(message); err != nil { if err := bot.SendMessage(message); err != nil {
@@ -137,14 +170,14 @@ func (m *Manager) ShowBotInitMessage() {
// controllers. The message content is identical to the /list command (same // controllers. The message content is identical to the /list command (same
// template and parameters). Like the init message it is typed // template and parameters). Like the init message it is typed
// MessageTypeBotStarted so members who opted out of bot-started pushes do not // MessageTypeBotStarted so members who opted out of bot-started pushes do not
// receive it. // receive it, and rendered once per controller so each channel's markdown
// setting is honored.
func (m *Manager) ShowBotServerList() { func (m *Manager) ShowBotServerList() {
m.mu.RLock() m.mu.RLock()
defer m.mu.RUnlock() defer m.mu.RUnlock()
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
content := template.Render(cfg.ControllerMessage.ServerList, params)
for _, ctrl := range m.controllers { for _, ctrl := range m.controllers {
if !ctrl.IsEnabled() { if !ctrl.IsEnabled() {
@@ -156,7 +189,7 @@ func (m *Manager) ShowBotServerList() {
} }
message := Message{ message := Message{
Source: bot.Name(), Source: bot.Name(),
Content: content, Content: template.Render(cfg.ControllerMessage.ServerList, params, ctrl.IsMarkdown()),
Type: MessageTypeBotStarted, Type: MessageTypeBotStarted,
} }
if err := bot.SendMessage(message); err != nil { if err := bot.SendMessage(message); err != nil {
@@ -188,6 +221,59 @@ func (m *Manager) NotifyStatusChange(change node.StatusChange) {
} }
} }
// NotifyAlert delivers an alert to the named pipes only, and returns the names
// of the pipes it was handed to. A pipe that is unknown, disabled or fails to
// send is reported through the returned error instead of stopping the delivery
// to the remaining pipes; if no pipe accepted the alert, the error describes
// every failure.
func (m *Manager) NotifyAlert(pipes []string, alert Alert) ([]string, error) {
m.mu.RLock()
defer m.mu.RUnlock()
var delivered, failures []string
for _, name := range pipes {
ctrl, ok := m.controllers[name]
if !ok {
failures = append(failures, fmt.Sprintf("%s: no such channel", name))
continue
}
if !ctrl.IsEnabled() {
failures = append(failures, fmt.Sprintf("%s: channel is disabled", name))
continue
}
if err := ctrl.SendAlert(alert); err != nil {
failures = append(failures, fmt.Sprintf("%s: %v", name, err))
continue
}
delivered = append(delivered, name)
}
if len(failures) > 0 {
postLog.Warning(fmt.Sprintf("Alert %q from %s not delivered by: %s", alert.Subject, alert.Source, strings.Join(failures, "; ")))
}
if len(delivered) == 0 {
if len(failures) == 0 {
return nil, errors.New("no notify channel configured")
}
return nil, errors.New(strings.Join(failures, "; "))
}
return delivered, nil
}
// IsMarkdown reports whether the pipe with the given name renders Markdown, per
// its channel's "markdown" setting in config.json. An unknown pipe renders
// plain text.
func (m *Manager) IsMarkdown(pipeName string) bool {
m.mu.RLock()
defer m.mu.RUnlock()
ctrl, ok := m.controllers[pipeName]
if !ok {
return false
}
return ctrl.IsMarkdown()
}
// StopAll stops all registered controllers. // StopAll stops all registered controllers.
func (m *Manager) StopAll() { func (m *Manager) StopAll() {
m.mu.RLock() m.mu.RLock()
+24 -3
View File
@@ -6,6 +6,7 @@ import (
gomail "gopkg.in/mail.v2" gomail "gopkg.in/mail.v2"
"nukumizu-backend/config" "nukumizu-backend/config"
"nukumizu-backend/internal/controller"
"nukumizu-backend/internal/netproxy" "nukumizu-backend/internal/netproxy"
"nukumizu-backend/internal/node" "nukumizu-backend/internal/node"
"nukumizu-backend/internal/template" "nukumizu-backend/internal/template"
@@ -54,6 +55,12 @@ func (e *EmailController) IsEnabled() bool {
return e.cfg.Enabled return e.cfg.Enabled
} }
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (e *EmailController) IsMarkdown() bool {
return e.cfg.Markdown
}
// SendStatusChange sends a status change notification via Email. // SendStatusChange sends a status change notification via Email.
func (e *EmailController) SendStatusChange(change node.StatusChange) error { func (e *EmailController) SendStatusChange(change node.StatusChange) error {
if !e.cfg.Enabled { if !e.cfg.Enabled {
@@ -66,7 +73,7 @@ func (e *EmailController) SendStatusChange(change node.StatusChange) error {
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
body := template.Render(cfg.ControllerMessage.ServerStatusChanged, params) body := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, e.cfg.Markdown)
subject := fmt.Sprintf("Server Status Change: %s - %s", change.Name, change.Event) subject := fmt.Sprintf("Server Status Change: %s - %s", change.Name, change.Event)
return e.sendEmail(subject, body) return e.sendEmail(subject, body)
@@ -80,7 +87,7 @@ func (e *EmailController) SendServerList(onlineServers, offlineServers string) e
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
body := template.Render(cfg.ControllerMessage.ServerList, params) body := template.Render(cfg.ControllerMessage.ServerList, params, e.cfg.Markdown)
return e.sendEmail("Server List", body) return e.sendEmail("Server List", body)
} }
@@ -93,12 +100,26 @@ func (e *EmailController) SendExecuteResult(serverName, serverUUID, command, res
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
body := template.Render(cfg.ControllerMessage.ServerExecuteResult, params) body := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, e.cfg.Markdown)
subject := fmt.Sprintf("Command Result: %s on %s", command, serverName) subject := fmt.Sprintf("Command Result: %s on %s", command, serverName)
return e.sendEmail(subject, body) return e.sendEmail(subject, body)
} }
// SendAlert sends an alert submitted through the incoming webhook API to the
// configured recipients.
func (e *EmailController) SendAlert(alert controller.Alert) error {
if !e.cfg.Enabled {
return nil
}
if len(e.cfg.To) == 0 {
postLog.Debug("Email controller has no recipients configured")
return nil
}
return e.sendEmail(alert.Subject, alert.Render(e.cfg.Markdown))
}
func (e *EmailController) sendEmail(subject, body string) error { func (e *EmailController) sendEmail(subject, body string) error {
m := gomail.NewMessage() m := gomail.NewMessage()
m.SetHeader("From", e.cfg.From) m.SetHeader("From", e.cfg.From)
+20 -3
View File
@@ -7,6 +7,7 @@ import (
"time" "time"
"nukumizu-backend/config" "nukumizu-backend/config"
"nukumizu-backend/internal/controller"
"nukumizu-backend/internal/netproxy" "nukumizu-backend/internal/netproxy"
"nukumizu-backend/internal/node" "nukumizu-backend/internal/node"
"nukumizu-backend/internal/template" "nukumizu-backend/internal/template"
@@ -52,6 +53,12 @@ func (n *NtfyController) IsEnabled() bool {
return n.cfg.Enabled return n.cfg.Enabled
} }
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (n *NtfyController) IsMarkdown() bool {
return n.cfg.Markdown
}
// SendStatusChange sends a status change notification via Ntfy. // SendStatusChange sends a status change notification via Ntfy.
func (n *NtfyController) SendStatusChange(change node.StatusChange) error { func (n *NtfyController) SendStatusChange(change node.StatusChange) error {
if !n.cfg.Enabled { if !n.cfg.Enabled {
@@ -60,7 +67,7 @@ func (n *NtfyController) SendStatusChange(change node.StatusChange) error {
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params) message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, n.cfg.Markdown)
title := fmt.Sprintf("Server %s: %s", change.Name, change.Event) title := fmt.Sprintf("Server %s: %s", change.Name, change.Event)
return n.publish(title, message) return n.publish(title, message)
@@ -74,7 +81,7 @@ func (n *NtfyController) SendServerList(onlineServers, offlineServers string) er
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params) message := template.Render(cfg.ControllerMessage.ServerList, params, n.cfg.Markdown)
return n.publish("Server List", message) return n.publish("Server List", message)
} }
@@ -87,12 +94,22 @@ func (n *NtfyController) SendExecuteResult(serverName, serverUUID, command, resu
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params) message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, n.cfg.Markdown)
title := fmt.Sprintf("Command Result: %s on %s", command, serverName) title := fmt.Sprintf("Command Result: %s on %s", command, serverName)
return n.publish(title, message) return n.publish(title, message)
} }
// SendAlert sends an alert submitted through the incoming webhook API to the
// configured topic.
func (n *NtfyController) SendAlert(alert controller.Alert) error {
if !n.cfg.Enabled {
return nil
}
return n.publish(alert.Subject, alert.Render(n.cfg.Markdown))
}
func (n *NtfyController) publish(title, message string) error { func (n *NtfyController) publish(title, message string) error {
serverURL := n.cfg.Server serverURL := n.cfg.Server
if serverURL == "" { if serverURL == "" {
+24 -4
View File
@@ -86,6 +86,12 @@ func (q *QQController) IsEnabled() bool {
return q.cfg.Enabled return q.cfg.Enabled
} }
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (q *QQController) IsMarkdown() bool {
return q.cfg.Markdown
}
// handleNapcatEvent processes a raw OneBot event received from the NapCat WebSocket. // handleNapcatEvent processes a raw OneBot event received from the NapCat WebSocket.
func (q *QQController) handleNapcatEvent(raw []byte) { func (q *QQController) handleNapcatEvent(raw []byte) {
var ev oneBotEvent var ev oneBotEvent
@@ -179,7 +185,7 @@ func (q *QQController) processCommand(cmd controller.Command) string {
parsed.ChatID = cmd.ChatID parsed.ChatID = cmd.ChatID
parsed.ChatType = cmd.ChatType parsed.ChatType = cmd.ChatType
parsed.SenderID = cmd.SenderID parsed.SenderID = cmd.SenderID
parsed.Source = "qq_napcat" parsed.Source = q.Name()
// Hand the complete command to the unified processor, which checks group // Hand the complete command to the unified processor, which checks group
// vs private, trusted groups, admin permissions, and executes it. // vs private, trusted groups, admin permissions, and executes it.
@@ -256,7 +262,7 @@ func (q *QQController) SendStatusChange(change node.StatusChange) error {
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params) message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, q.cfg.Markdown)
// Only notify trusted groups and admins whose event_status_notify is true. // Only notify trusted groups and admins whose event_status_notify is true.
if uc := config.C_botUserConfig; uc != nil { if uc := config.C_botUserConfig; uc != nil {
@@ -285,7 +291,7 @@ func (q *QQController) SendServerList(onlineServers, offlineServers string) erro
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params) message := template.Render(cfg.ControllerMessage.ServerList, params, q.cfg.Markdown)
for _, groupID := range q.trustedGroupIDs() { for _, groupID := range q.trustedGroupIDs() {
q.sendGroupMessage(groupID, message) q.sendGroupMessage(groupID, message)
@@ -301,7 +307,7 @@ func (q *QQController) SendExecuteResult(serverName, serverUUID, command, result
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params) message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, q.cfg.Markdown)
for _, groupID := range q.trustedGroupIDs() { for _, groupID := range q.trustedGroupIDs() {
q.sendGroupMessage(groupID, message) q.sendGroupMessage(groupID, message)
@@ -309,6 +315,20 @@ func (q *QQController) SendExecuteResult(serverName, serverUUID, command, result
return nil return nil
} }
// SendAlert sends an alert submitted through the incoming webhook API to all QQ
// trusted groups and admins.
func (q *QQController) SendAlert(alert controller.Alert) error {
if !q.cfg.Enabled {
return nil
}
return q.SendMessage(controller.Message{
Source: q.Name(),
Content: alert.Render(q.cfg.Markdown),
Type: controller.MessageTypeAlert,
})
}
func (q *QQController) sendGroupMessage(groupID string, message string) { func (q *QQController) sendGroupMessage(groupID string, message string) {
if q.napcatClient == nil { if q.napcatClient == nil {
postLog.Warning("Cannot send QQ group message: NapCat client not initialized") postLog.Warning("Cannot send QQ group message: NapCat client not initialized")
+14 -8
View File
@@ -14,11 +14,13 @@ import (
const maxMessageLen = 4000 const maxMessageLen = 4000
// sendMessage sends a text message to a chat, splitting it into chunks that fit // sendMessage sends a text message to a chat, splitting it into chunks that fit
// Telegram's 4096-character limit. All messages are sent with // Telegram's 4096-character limit. When the channel has markdown enabled the
// parse_mode=Markdown so fenced code blocks and inline formatting render as // message is sent with parse_mode=Markdown so fenced code blocks and inline
// rich text. Templates must stay valid under Telegram's legacy Markdown: // formatting render as rich text; templates must then stay valid under
// unpaired '*' or '_' characters (e.g. a lone '*Event: ...' label) make the // Telegram's legacy Markdown, because unpaired '*' or '_' characters (e.g. a
// API reject the whole message. // lone '*Event: ...' label) make the API reject the whole message. With
// markdown disabled the message is sent without a parse mode, so it is
// delivered verbatim whatever it contains.
func (t *TelegramController) sendMessage(message controller.Message) error { func (t *TelegramController) sendMessage(message controller.Message) error {
if t.client == nil { if t.client == nil {
return nil return nil
@@ -40,11 +42,15 @@ func (t *TelegramController) sendMessageChunk(chatID int64, text string) error {
ctx, cancel := context.WithTimeout(context.Background(), apiTimeout) ctx, cancel := context.WithTimeout(context.Background(), apiTimeout)
defer cancel() defer cancel()
_, err := t.client.SendMessage(ctx, &bot.SendMessageParams{ params := &bot.SendMessageParams{
ChatID: chatID, ChatID: chatID,
Text: text, Text: text,
ParseMode: models.ParseModeMarkdownV1, // Telegram legacy Markdown }
}) if t.cfg.Markdown {
params.ParseMode = models.ParseModeMarkdownV1 // Telegram legacy Markdown
}
_, err := t.client.SendMessage(ctx, params)
return err return err
} }
+23 -3
View File
@@ -124,6 +124,12 @@ func (t *TelegramController) IsEnabled() bool {
return t.cfg.Enabled return t.cfg.Enabled
} }
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (t *TelegramController) IsMarkdown() bool {
return t.cfg.Markdown
}
// handleUpdate processes a single Telegram update received via long polling. It // handleUpdate processes a single Telegram update received via long polling. It
// is installed as the framework's default handler (every update with a Message // is installed as the framework's default handler (every update with a Message
// reaches it). Updates are processed sequentially because the bot is created // reaches it). Updates are processed sequentially because the bot is created
@@ -271,7 +277,7 @@ func (t *TelegramController) SendStatusChange(change node.StatusChange) error {
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params) message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, t.cfg.Markdown)
// Only notify trusted groups and admins whose event_status_notify is true. // Only notify trusted groups and admins whose event_status_notify is true.
if uc := config.C_botUserConfig; uc != nil { if uc := config.C_botUserConfig; uc != nil {
@@ -299,7 +305,7 @@ func (t *TelegramController) SendServerList(onlineServers, offlineServers string
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params) message := template.Render(cfg.ControllerMessage.ServerList, params, t.cfg.Markdown)
t.sendToGroups(message) t.sendToGroups(message)
return nil return nil
@@ -313,12 +319,26 @@ func (t *TelegramController) SendExecuteResult(serverName, serverUUID, command,
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params) message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, t.cfg.Markdown)
t.sendToGroups(message) t.sendToGroups(message)
return nil return nil
} }
// SendAlert sends an alert submitted through the incoming webhook API to all
// Telegram trusted groups and admins.
func (t *TelegramController) SendAlert(alert controller.Alert) error {
if !t.cfg.Enabled || t.client == nil {
return nil
}
return t.SendMessage(controller.Message{
Source: t.Name(),
Content: alert.Render(t.cfg.Markdown),
Type: controller.MessageTypeAlert,
})
}
// telegramChatType maps a Telegram chat type to the unified ChatType value used // telegramChatType maps a Telegram chat type to the unified ChatType value used
// by the controller package. Empty means the chat type is unsupported. // by the controller package. Empty means the chat type is unsupported.
func telegramChatType(chatType string) string { func telegramChatType(chatType string) string {
+29 -3
View File
@@ -8,6 +8,7 @@ import (
"time" "time"
"nukumizu-backend/config" "nukumizu-backend/config"
"nukumizu-backend/internal/controller"
"nukumizu-backend/internal/netproxy" "nukumizu-backend/internal/netproxy"
"nukumizu-backend/internal/node" "nukumizu-backend/internal/node"
"nukumizu-backend/internal/template" "nukumizu-backend/internal/template"
@@ -53,6 +54,12 @@ func (w *WebhookController) IsEnabled() bool {
return w.cfg.Enabled return w.cfg.Enabled
} }
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (w *WebhookController) IsMarkdown() bool {
return w.cfg.Markdown
}
// SendStatusChange sends a status change notification via Webhook. // SendStatusChange sends a status change notification via Webhook.
func (w *WebhookController) SendStatusChange(change node.StatusChange) error { func (w *WebhookController) SendStatusChange(change node.StatusChange) error {
if !w.cfg.Enabled { if !w.cfg.Enabled {
@@ -61,7 +68,7 @@ func (w *WebhookController) SendStatusChange(change node.StatusChange) error {
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params) message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, w.cfg.Markdown)
payload := map[string]interface{}{ payload := map[string]interface{}{
"event": change.Event, "event": change.Event,
@@ -82,7 +89,7 @@ func (w *WebhookController) SendServerList(onlineServers, offlineServers string)
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params) message := template.Render(cfg.ControllerMessage.ServerList, params, w.cfg.Markdown)
payload := map[string]interface{}{ payload := map[string]interface{}{
"type": "serverList", "type": "serverList",
@@ -103,7 +110,7 @@ func (w *WebhookController) SendExecuteResult(serverName, serverUUID, command, r
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params) message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, w.cfg.Markdown)
payload := map[string]interface{}{ payload := map[string]interface{}{
"type": "executeResult", "type": "executeResult",
@@ -118,6 +125,25 @@ func (w *WebhookController) SendExecuteResult(serverName, serverUUID, command, r
return w.send(payload) return w.send(payload)
} }
// SendAlert sends an alert submitted through the incoming webhook API to the
// configured URL.
func (w *WebhookController) SendAlert(alert controller.Alert) error {
if !w.cfg.Enabled {
return nil
}
payload := map[string]interface{}{
"type": "alert",
"subject": alert.Subject,
"source": alert.Source,
"content": alert.Content,
"message": alert.Render(w.cfg.Markdown),
"time": alert.Time,
}
return w.send(payload)
}
func (w *WebhookController) send(payload map[string]interface{}) error { func (w *WebhookController) send(payload map[string]interface{}) error {
method := w.cfg.Method method := w.cfg.Method
if method == "" { if method == "" {
+16 -5
View File
@@ -10,16 +10,27 @@ import (
"nukumizu-backend/internal/template" "nukumizu-backend/internal/template"
) )
// commandMarkdown reports whether responses to the given command are rendered
// with Markdown, per the markdown setting of the pipe the command came from
// (see Command.Source).
func commandMarkdown(cmd Command) bool {
mgr := GetManager()
if mgr == nil {
return false
}
return mgr.IsMarkdown(cmd.Source)
}
func handleHelp(cmd Command) (string, error) { func handleHelp(cmd Command) (string, error) {
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildBotInitializationMsgParams() params := template.BuildBotInitializationMsgParams()
return template.Render(cfg.ControllerMessage.BotHelp, params, cmd.Source), nil return template.Render(cfg.ControllerMessage.BotHelp, params, commandMarkdown(cmd)), nil
} }
func handleList(cmd Command) (string, error) { func handleList(cmd Command) (string, error) {
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
return template.Render(cfg.ControllerMessage.ServerList, params, cmd.Source), nil return template.Render(cfg.ControllerMessage.ServerList, params, commandMarkdown(cmd)), nil
} }
func handleStatus(cmd Command) (string, error) { func handleStatus(cmd Command) (string, error) {
@@ -130,7 +141,7 @@ func handleRun(cmd Command) (string, error) {
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(uuidArg, uuidArg, command, formatTaskResults(results)) params := template.BuildParamsFromExecResult(uuidArg, uuidArg, command, formatTaskResults(results))
return template.Render(cfg.ControllerMessage.ServerExecuteResult, params, cmd.Source), nil return template.Render(cfg.ControllerMessage.ServerExecuteResult, params, commandMarkdown(cmd)), nil
} }
func handleInfo(cmd Command) (string, error) { func handleInfo(cmd Command) (string, error) {
@@ -177,10 +188,10 @@ func handleInfo(cmd Command) (string, error) {
return sb.String(), nil return sb.String(), nil
} }
func telegram_handleStart() (string, error) { func telegram_handleStart(cmd Command) (string, error) {
cfg := config.C_globalConfig cfg := config.C_globalConfig
params := template.BuildBotInitializationMsgParams() params := template.BuildBotInitializationMsgParams()
return template.Render(cfg.ControllerMessage.Tg_BotStart, params, "telegram"), nil return template.Render(cfg.ControllerMessage.Tg_BotStart, params, commandMarkdown(cmd)), nil
} }
func handleGetIP(cmd Command) (string, error) { func handleGetIP(cmd Command) (string, error) {
+1 -1
View File
@@ -46,7 +46,7 @@ func (m *Manager) RouteCommand(cmd Command) (string, error) {
if cmd.Source == "telegram" { if cmd.Source == "telegram" {
switch cmd.Command { switch cmd.Command {
case "start": case "start":
return telegram_handleStart() return telegram_handleStart(cmd)
} }
} }
switch cmd.Command { switch cmd.Command {
+55 -4
View File
@@ -19,6 +19,9 @@ type Params struct {
Message string Message string
Command string Command string
Result string Result string
Subject string // Alert subject (see AlertParams)
Source string // Alert source (see AlertParams)
Content string // Alert content (see AlertParams)
OnlineServers string // Pre-formatted multi-line list OnlineServers string // Pre-formatted multi-line list
OfflineServers string // Pre-formatted multi-line list OfflineServers string // Pre-formatted multi-line list
SoftwareVersion string SoftwareVersion string
@@ -30,6 +33,37 @@ type Params struct {
SoftwareDescription string SoftwareDescription string
} }
// AlertTemplate is the body format of an alert submitted by an external
// application through the incoming webhook API.
const AlertTemplate = "{{ subject }}\n- Source: {{ source }}\n- Content:\n{{ content }}\n\n- Time: {{ time }}\nSent by Nukumizu Alert System"
// AlertParams holds the parameters of an alert submitted through the incoming
// webhook API.
type AlertParams struct {
Subject string // Short one-line title of the alert
Source string // Name of the webhook endpoint the alert was submitted to
Content string // Free-form alert body
Time string // Submission time
}
// RenderAlert renders the body of an alert for a channel. The alert source and
// content may be wrapped in Markdown — the source in inline code, the content
// in a fenced code block — when the target channel has markdown enabled
// (markdown); everything else, the timestamp included, stays plain text.
func RenderAlert(alert AlertParams, markdown bool) string {
params := Params{
Time: alert.Time,
Subject: alert.Subject,
Source: alert.Source,
Content: alert.Content,
}
if markdown {
params.Source = "`" + params.Source + "`"
params.Content = "```\n" + params.Content + "\n```"
}
return Render(AlertTemplate, params, false)
}
// BuildBotInitializationMsgParams creates template parameters for the bot initialization message. // BuildBotInitializationMsgParams creates template parameters for the bot initialization message.
func BuildBotInitializationMsgParams() Params { func BuildBotInitializationMsgParams() Params {
return Params{ return Params{
@@ -80,7 +114,13 @@ func BuildParamsFromExecResult(serverName, serverUUID, command, result string) P
} }
} }
// Render substitutes {{ paramName }} placeholders in a template string. // Render substitutes {{ paramName }} placeholders in a template string. The
// markdown argument is the target channel's markdown setting: when true the
// values that are meant to be read verbatim (UUIDs, messages, commands, command
// results) are wrapped in Markdown code spans and blocks, otherwise every value
// is inserted as plain text. Whether a channel renders Markdown comes from the
// configuration alone — the renderer never infers it from the channel name.
//
// Supported placeholders: // Supported placeholders:
// - {{ time }} — current server time // - {{ time }} — current server time
// - {{ serverName }} — server name // - {{ serverName }} — server name
@@ -90,6 +130,9 @@ func BuildParamsFromExecResult(serverName, serverUUID, command, result string) P
// - {{ message }} — event descriptive message // - {{ message }} — event descriptive message
// - {{ command }} — executed command // - {{ command }} — executed command
// - {{ result }} — command execution result // - {{ result }} — command execution result
// - {{ subject }} — alert subject
// - {{ source }} — alert source
// - {{ content }} — alert content
// - {{ list.onlineServers }} — multi-line online server list // - {{ list.onlineServers }} — multi-line online server list
// - {{ list.offlineServers }} — multi-line offline server list // - {{ list.offlineServers }} — multi-line offline server list
// - {{ softwareVersion }} — software version // - {{ softwareVersion }} — software version
@@ -99,11 +142,12 @@ func BuildParamsFromExecResult(serverName, serverUUID, command, result string) P
// - {{ softwareBuildTime }} — software build time // - {{ softwareBuildTime }} — software build time
// - {{ softwareDeveloper }} — software developer // - {{ softwareDeveloper }} — software developer
// - {{ softwareDescription }} — software description // - {{ softwareDescription }} — software description
func Render(tmpl string, params Params, source ...string) string { func Render(tmpl string, params Params, markdown bool) string {
result := tmpl result := tmpl
if len(source) > 0 && source[0] == "telegram" { if markdown {
// Telegram requires special formatting for code blocks and inline code. // Channels that render Markdown get code blocks and inline code for the
// values that are read verbatim.
result = strings.ReplaceAll(result, "{{ time }}", "**"+params.Time+"**") result = strings.ReplaceAll(result, "{{ time }}", "**"+params.Time+"**")
result = strings.ReplaceAll(result, "{{ serverName }}", "**"+params.ServerName+"**") result = strings.ReplaceAll(result, "{{ serverName }}", "**"+params.ServerName+"**")
result = strings.ReplaceAll(result, "{{ serverUUID }}", "`"+params.ServerUUID+"`") result = strings.ReplaceAll(result, "{{ serverUUID }}", "`"+params.ServerUUID+"`")
@@ -140,6 +184,13 @@ func Render(tmpl string, params Params, source ...string) string {
result = strings.ReplaceAll(result, "{{ softwareDeveloper }}", params.SoftwareDeveloper) result = strings.ReplaceAll(result, "{{ softwareDeveloper }}", params.SoftwareDeveloper)
result = strings.ReplaceAll(result, "{{ softwareDescription }}", params.SoftwareDescription) result = strings.ReplaceAll(result, "{{ softwareDescription }}", params.SoftwareDescription)
} }
// Alert values carry their own formatting (see RenderAlert), so they are
// substituted identically in both branches.
result = strings.ReplaceAll(result, "{{ subject }}", params.Subject)
result = strings.ReplaceAll(result, "{{ source }}", params.Source)
result = strings.ReplaceAll(result, "{{ content }}", params.Content)
return result return result
} }
+32
View File
@@ -146,6 +146,9 @@ func main() {
} }
}() }()
// --- Start the incoming webhook listener ---
startWebhookServer(cfg)
// --- Graceful shutdown --- // --- Graceful shutdown ---
quit := make(chan os.Signal, 1) quit := make(chan os.Signal, 1)
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
@@ -167,6 +170,35 @@ func main() {
postLog.Info("Server stopped") postLog.Info("Server stopped")
} }
// startWebhookServer serves the incoming webhook API on its own listener. The
// API is not exposed on the main listener: external applications post alerts to
// this port only, so its rate limiter and CORS policy are configured
// independently. A failure to bind it is logged rather than fatal — the rest of
// the program (bots, status monitoring) keeps running without it.
func startWebhookServer(cfg *config.Config) {
if !cfg.Webhook.Enabled {
postLog.Warning("Incoming webhook API is disabled")
return
}
handler := utils.RateLimitMiddleware(SetupWebhookRouter())
handler = utils.CORSMiddleware(handler)
addr := fmt.Sprintf("%s:%s", cfg.Webhook.ListenAddr, cfg.Webhook.ListenPort)
postLog.Info(fmt.Sprintf("Webhook API listening on %s", addr))
go func() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Webhook server panic: %v", r))
}
}()
if err := http.ListenAndServe(addr, handler); err != nil {
postLog.Error("Webhook server error: " + err.Error())
}
}()
}
// initControllers initializes and starts all configured controllers. // initControllers initializes and starts all configured controllers.
func initControllers() { func initControllers() {
cfg := config.C_globalConfig cfg := config.C_globalConfig
+19
View File
@@ -50,6 +50,25 @@ func SetupRouter() *http.ServeMux {
return mux return mux
} }
// SetupWebhookRouter registers the routes of the incoming webhook API. Unlike
// SetupRouter it is served on its own listener (webhook.listenAddr/listenPort),
// so external applications can be given access to the webhook port without
// reaching the admin API. Every endpoint configured under webhook.endpoints is
// reachable as /api/webhook/<name>.
func SetupWebhookRouter() *http.ServeMux {
postLog.Info("Setting up webhook routers...")
mux := http.NewServeMux()
// The wildcard segment selects the endpoint; requests for a name that is not
// configured fall through to the handler, which answers with a JSON 404.
mux.HandleFunc("/api/webhook/{name}", handler.WebhookHandler)
mux.HandleFunc("/", NotFoundHandler)
postLog.Info("Webhook router setup completed")
return mux
}
// NotFoundHandler returns a 404 JSON response for unknown routes. // NotFoundHandler returns a 404 JSON response for unknown routes.
func NotFoundHandler(w http.ResponseWriter, r *http.Request) { func NotFoundHandler(w http.ResponseWriter, r *http.Request) {
postLog.Debug(fmt.Sprintf("Unknown request: %s %s", r.Method, r.URL.Path)) postLog.Debug(fmt.Sprintf("Unknown request: %s %s", r.Method, r.URL.Path))