feat(controller): hot-reload the notification pipes
Build / ubuntu-latest (push) Failing after 3m13s
Build / windows-latest (push) Canceled after 6m23s

Email, ntfy and webhook kept their settings in a plain struct field copied at
construction, so editing controllerMethod in config.json wrote the file and
replaced the in-memory configuration while the running channels stayed on their
boot values.

Add Reload to the Controller interface and implement it for the three
notification pipes; Manager.ReloadAll fans an update out to every registered
controller and is wired into the reload hook from main.

Reload runs on the goroutine serving the settings update while the send methods
run on the status-change and incoming-webhook goroutines, so storing the new
settings in a plain field would be a data race. Each pipe publishes its settings
— and the HTTP client derived from them — through atomic pointers, and every
send takes one snapshot so a reload landing mid-send cannot split it across two
configurations.

Only a change to networkUseProxy rebuilds the client; everything else is read
from the settings at send time, so a reload that changes nothing relevant leaves
the client alone. Email needed one more fix: gomail exposes its dial hook as a
package-level variable with no unset, and the constructor only installed the
proxy dialer when the flag was set, so turning networkUseProxy off left SMTP
tunnelled through a proxy. applyEmailProxy now restores net.DialTimeout, which
is gomail's own default.

QQ (Napcat) and Telegram implement Reload because the interface requires it, but
neither can apply a change in place: the NapCat client is stopped through a
sync.Once and the Telegram polling context is created with the controller, so
both need to be rebuilt and swapped into the manager. Rather than no-op silently
and let an edit look applied, they log a warning while their settings diverge.
That rebuild is the next step.
This commit is contained in:
2026-09-28 23:15:11 +08:00
parent d9a9918e21
commit 0691e85ccf
8 changed files with 553 additions and 79 deletions
+83 -28
View File
@@ -2,6 +2,9 @@ package pipes
import (
"fmt"
"net"
"reflect"
"sync/atomic"
gomail "gopkg.in/mail.v2"
@@ -14,20 +17,41 @@ import (
)
// EmailController handles email notifications via SMTP.
//
// The configuration is held behind an atomic pointer rather than in a plain
// field: Reload replaces it on the settings-update goroutine while the send
// methods read it on the status-change and incoming-webhook goroutines.
type EmailController struct {
cfg config.EmailConfig
cfg atomic.Pointer[config.EmailConfig]
}
// NewEmailController creates a new Email controller.
func NewEmailController(cfg config.EmailConfig) *EmailController {
if cfg.NetworkUseProxy {
// Route SMTP through the HTTP CONNECT proxy. NetDialTimeout is
// gomail's documented hook for overriding how the SMTP connection is
// dialed. There is a single global email channel, so overriding it
// unconditionally when the flag is set is safe.
e := &EmailController{}
e.cfg.Store(&cfg)
applyEmailProxy(cfg.NetworkUseProxy)
return e
}
// settings returns the configuration currently in effect. The value it points
// at is never mutated after being stored, so a caller can hold the pointer for
// one whole send without a concurrent Reload disturbing it.
func (e *EmailController) settings() *config.EmailConfig {
return e.cfg.Load()
}
// applyEmailProxy routes SMTP through the HTTP CONNECT proxy, or restores a
// direct dial. gomail exposes the dial path as the package-level
// NetDialTimeout, whose own default is net.DialTimeout, so turning the proxy
// off has to put that back rather than leave the hook in place. There is a
// single global email channel, so setting a package-level hook here is
// unambiguous.
func applyEmailProxy(useProxy bool) {
if useProxy {
gomail.NetDialTimeout = netproxy.DialWithTimeout(true)
return
}
return &EmailController{cfg: cfg}
gomail.NetDialTimeout = net.DialTimeout
}
// Name returns the controller name.
@@ -37,7 +61,7 @@ func (e *EmailController) Name() string {
// Start initializes the Email controller.
func (e *EmailController) Start() error {
if !e.cfg.Enabled {
if !e.settings().Enabled {
postLog.Info("Email controller is disabled")
return nil
}
@@ -50,90 +74,121 @@ func (e *EmailController) Stop() {
postLog.Info("Email controller stopped")
}
// Reload applies the current configuration. Swapping in the new settings is
// enough for everything read at the point of use; the proxy setting is the
// exception, because it is baked into gomail's package-level dialer when the
// controller is built rather than consulted per send.
func (e *EmailController) Reload() {
global := config.Current()
if global == nil {
return
}
current := e.settings()
updated := global.ControllerMethod.Email
// EmailConfig carries a recipient slice, so it is not comparable with ==.
if reflect.DeepEqual(*current, updated) {
return
}
if current.NetworkUseProxy != updated.NetworkUseProxy {
applyEmailProxy(updated.NetworkUseProxy)
}
e.cfg.Store(&updated)
postLog.Debug("Email controller reloaded")
}
// IsEnabled returns whether the controller is enabled.
func (e *EmailController) IsEnabled() bool {
return e.cfg.Enabled
return e.settings().Enabled
}
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (e *EmailController) IsMarkdown() bool {
return e.cfg.Markdown
return e.settings().Markdown
}
// SendStatusChange sends a status change notification via Email.
func (e *EmailController) SendStatusChange(change node.StatusChange) error {
if !e.cfg.Enabled {
s := e.settings()
if !s.Enabled {
return nil
}
if len(e.cfg.To) == 0 {
if len(s.To) == 0 {
postLog.Debug("Email controller has no recipients configured")
return nil
}
cfg := config.Current()
params := template.BuildParamsFromStatusChange(change)
body := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, e.cfg.Markdown)
body := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, s.Markdown)
subject := fmt.Sprintf("Server Status Change: %s - %s", change.Name, change.Event)
return e.sendEmail(subject, body)
return e.sendEmail(s, subject, body)
}
// SendServerList sends the server list via Email.
func (e *EmailController) SendServerList(onlineServers, offlineServers string) error {
if !e.cfg.Enabled || len(e.cfg.To) == 0 {
s := e.settings()
if !s.Enabled || len(s.To) == 0 {
return nil
}
cfg := config.Current()
params := template.BuildParamsFromServerList()
body := template.Render(cfg.ControllerMessage.ServerList, params, e.cfg.Markdown)
body := template.Render(cfg.ControllerMessage.ServerList, params, s.Markdown)
return e.sendEmail("Server List", body)
return e.sendEmail(s, "Server List", body)
}
// SendExecuteResult sends a command execution result via Email.
func (e *EmailController) SendExecuteResult(serverName, serverUUID, command, result string) error {
if !e.cfg.Enabled || len(e.cfg.To) == 0 {
s := e.settings()
if !s.Enabled || len(s.To) == 0 {
return nil
}
cfg := config.Current()
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
body := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, e.cfg.Markdown)
body := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, s.Markdown)
subject := fmt.Sprintf("Command Result: %s on %s", command, serverName)
return e.sendEmail(subject, body)
return e.sendEmail(s, 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 {
s := e.settings()
if !s.Enabled {
return nil
}
if len(e.cfg.To) == 0 {
if len(s.To) == 0 {
postLog.Debug("Email controller has no recipients configured")
return nil
}
return e.sendEmail(alert.Subject, alert.Render(e.cfg.Markdown))
return e.sendEmail(s, alert.Subject, alert.Render(s.Markdown))
}
func (e *EmailController) sendEmail(subject, body string) error {
// sendEmail delivers one message using the settings the caller already
// snapshotted, so the recipients the guard approved are the recipients that
// receive it even if a reload lands mid-send.
func (e *EmailController) sendEmail(s *config.EmailConfig, subject, body string) error {
m := gomail.NewMessage()
m.SetHeader("From", e.cfg.From)
m.SetHeader("To", e.cfg.To...)
m.SetHeader("From", s.From)
m.SetHeader("To", s.To...)
m.SetHeader("Subject", subject)
m.SetBody("text/plain", body)
d := gomail.NewDialer(e.cfg.SMTPHost, e.cfg.SMTPPort, e.cfg.Username, e.cfg.Password)
d := gomail.NewDialer(s.SMTPHost, s.SMTPPort, s.Username, s.Password)
if err := d.DialAndSend(m); err != nil {
postLog.Warning("Failed to send email: " + err.Error())
return err
}
postLog.Debug("Email sent successfully to " + fmt.Sprintf("%v", e.cfg.To))
postLog.Debug("Email sent successfully to " + fmt.Sprintf("%v", s.To))
return nil
}
+71 -25
View File
@@ -4,6 +4,7 @@ import (
"fmt"
"net/http"
"strings"
"sync/atomic"
"time"
"nukumizu-backend/config"
@@ -15,17 +16,33 @@ import (
)
// NtfyController handles notifications via ntfy.sh or a self-hosted ntfy server.
//
// Both fields are held behind atomic pointers rather than in plain fields:
// Reload replaces them on the settings-update goroutine while the send methods
// read them on the status-change and incoming-webhook goroutines.
type NtfyController struct {
cfg config.NtfyConfig
httpClient *http.Client
cfg atomic.Pointer[config.NtfyConfig]
httpClient atomic.Pointer[http.Client]
}
// NewNtfyController creates a new Ntfy controller.
func NewNtfyController(cfg config.NtfyConfig) *NtfyController {
return &NtfyController{
cfg: cfg,
httpClient: netproxy.HTTPClient(cfg.NetworkUseProxy, 10*time.Second),
}
n := &NtfyController{}
n.cfg.Store(&cfg)
n.httpClient.Store(netproxy.HTTPClient(cfg.NetworkUseProxy, 10*time.Second))
return n
}
// settings returns the configuration currently in effect. The value it points
// at is never mutated after being stored, so a caller can hold the pointer for
// one whole send without a concurrent Reload disturbing it.
func (n *NtfyController) settings() *config.NtfyConfig {
return n.cfg.Load()
}
// client returns the HTTP client built for the current proxy setting.
func (n *NtfyController) client() *http.Client {
return n.httpClient.Load()
}
// Name returns the controller name.
@@ -35,7 +52,7 @@ func (n *NtfyController) Name() string {
// Start initializes the Ntfy controller.
func (n *NtfyController) Start() error {
if !n.cfg.Enabled {
if !n.settings().Enabled {
postLog.Info("Ntfy controller is disabled")
return nil
}
@@ -48,26 +65,50 @@ func (n *NtfyController) Stop() {
postLog.Info("Ntfy controller stopped")
}
// Reload applies the current configuration. Everything the publisher reads is
// taken from the settings at send time, so only a change to the proxy flag
// needs more than the swap: that one is baked into the HTTP client's transport
// when the client is built.
func (n *NtfyController) Reload() {
global := config.Current()
if global == nil {
return
}
current := n.settings()
updated := global.ControllerMethod.Ntfy
if *current == updated {
return
}
if current.NetworkUseProxy != updated.NetworkUseProxy {
n.httpClient.Store(netproxy.HTTPClient(updated.NetworkUseProxy, 10*time.Second))
}
n.cfg.Store(&updated)
postLog.Debug("Ntfy controller reloaded")
}
// IsEnabled returns whether the controller is enabled.
func (n *NtfyController) IsEnabled() bool {
return n.cfg.Enabled
return n.settings().Enabled
}
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (n *NtfyController) IsMarkdown() bool {
return n.cfg.Markdown
return n.settings().Markdown
}
// SendStatusChange sends a status change notification via Ntfy.
func (n *NtfyController) SendStatusChange(change node.StatusChange) error {
if !n.cfg.Enabled {
s := n.settings()
if !s.Enabled {
return nil
}
cfg := config.Current()
params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, n.cfg.Markdown)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, s.Markdown)
title := fmt.Sprintf("Server %s: %s", change.Name, change.Event)
return n.publish(title, message)
@@ -75,26 +116,28 @@ func (n *NtfyController) SendStatusChange(change node.StatusChange) error {
// SendServerList sends the server list via Ntfy.
func (n *NtfyController) SendServerList(onlineServers, offlineServers string) error {
if !n.cfg.Enabled {
s := n.settings()
if !s.Enabled {
return nil
}
cfg := config.Current()
params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params, n.cfg.Markdown)
message := template.Render(cfg.ControllerMessage.ServerList, params, s.Markdown)
return n.publish("Server List", message)
}
// SendExecuteResult sends a command execution result via Ntfy.
func (n *NtfyController) SendExecuteResult(serverName, serverUUID, command, result string) error {
if !n.cfg.Enabled {
s := n.settings()
if !s.Enabled {
return nil
}
cfg := config.Current()
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, n.cfg.Markdown)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, s.Markdown)
title := fmt.Sprintf("Command Result: %s on %s", command, serverName)
return n.publish(title, message)
@@ -103,21 +146,24 @@ func (n *NtfyController) SendExecuteResult(serverName, serverUUID, command, resu
// 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 {
s := n.settings()
if !s.Enabled {
return nil
}
return n.publish(alert.Subject, alert.Render(n.cfg.Markdown))
return n.publish(alert.Subject, alert.Render(s.Markdown))
}
func (n *NtfyController) publish(title, message string) error {
serverURL := n.cfg.Server
s := n.settings()
serverURL := s.Server
if serverURL == "" {
serverURL = "https://ntfy.sh"
}
serverURL = strings.TrimRight(serverURL, "/")
publishURL := fmt.Sprintf("%s/%s", serverURL, n.cfg.Topic)
publishURL := fmt.Sprintf("%s/%s", serverURL, s.Topic)
req, err := http.NewRequest("POST", publishURL, strings.NewReader(message))
if err != nil {
@@ -125,14 +171,14 @@ func (n *NtfyController) publish(title, message string) error {
}
req.Header.Set("Title", title)
if n.cfg.Priority != "" && n.cfg.Priority != "default" {
req.Header.Set("Priority", n.cfg.Priority)
if s.Priority != "" && s.Priority != "default" {
req.Header.Set("Priority", s.Priority)
}
if n.cfg.Token != "" {
req.Header.Set("Authorization", "Bearer "+n.cfg.Token)
if s.Token != "" {
req.Header.Set("Authorization", "Bearer "+s.Token)
}
resp, err := n.httpClient.Do(req)
resp, err := n.client().Do(req)
if err != nil {
postLog.Warning("Failed to publish to ntfy: " + err.Error())
return err
@@ -143,6 +189,6 @@ func (n *NtfyController) publish(title, message string) error {
postLog.Warning(fmt.Sprintf("Ntfy publish returned status %d", resp.StatusCode))
}
postLog.Debug("Ntfy notification sent to topic: " + n.cfg.Topic)
postLog.Debug("Ntfy notification sent to topic: " + s.Topic)
return nil
}
+17
View File
@@ -86,6 +86,23 @@ func (q *QQController) IsEnabled() bool {
return q.cfg.Enabled
}
// Reload is not implemented yet for this channel.
//
// Everything this controller uses — the NapCat address, port and token — is
// captured in the client it connected with, and that client's WebSocket
// listener is stopped through a sync.Once and cannot be re-pointed. Applying a
// change therefore means stopping the client and building a replacement
// controller, which lands together with the Telegram one. Until then the
// divergence is reported rather than ignored, so an edit that needs a restart
// does not look like it took effect.
func (q *QQController) Reload() {
global := config.Current()
if global == nil || q.cfg == global.ControllerMethod.QQ {
return
}
postLog.Warning("QQ (Napcat) settings changed but hot reload is not implemented for this channel yet; restart to apply them")
}
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (q *QQController) IsMarkdown() bool {
+257
View File
@@ -0,0 +1,257 @@
package pipes
import (
"fmt"
"net"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"reflect"
"sync"
"sync/atomic"
"testing"
gomail "gopkg.in/mail.v2"
"nukumizu-backend/config"
"nukumizu-backend/internal/node"
)
// writeConfigFile writes one config.json and returns its path. Reload tests
// need several of these up front, because publishing a configuration from
// inside a goroutine cannot use t.Fatalf.
func writeConfigFile(t *testing.T, body string) string {
t.Helper()
path := filepath.Join(t.TempDir(), "config.json")
if err := os.WriteFile(path, []byte(body), 0o644); err != nil {
t.Fatalf("write temp config: %v", err)
}
return path
}
// publishConfig writes body and makes it the configuration in effect, which is
// what a settings update does before the reload hooks run.
func publishConfig(t *testing.T, body string) *config.Config {
t.Helper()
cfg, err := config.LoadGlobalConfig(writeConfigFile(t, body))
if err != nil {
t.Fatalf("LoadGlobalConfig: %v", err)
}
return cfg
}
func TestNtfyControllerReload(t *testing.T) {
first := publishConfig(t, `{
"controllerMethod": {
"ntfy": { "enabled": false, "markdown": false, "server": "https://ntfy.sh", "topic": "old" }
}
}`)
ctrl := NewNtfyController(first.ControllerMethod.Ntfy)
if ctrl.IsEnabled() {
t.Fatal("controller should start disabled")
}
publishConfig(t, `{
"controllerMethod": {
"ntfy": { "enabled": true, "markdown": true, "server": "https://ntfy.example.com", "topic": "new" }
}
}`)
ctrl.Reload()
if !ctrl.IsEnabled() {
t.Error("Reload did not pick up enabled")
}
if !ctrl.IsMarkdown() {
t.Error("Reload did not pick up markdown")
}
s := ctrl.settings()
if s.Topic != "new" || s.Server != "https://ntfy.example.com" {
t.Errorf("Reload did not pick up the new server/topic: %+v", s)
}
}
// TestReloadKeepsDerivedClientsWhenNothingRelevantChanged pins the guard: a
// reload that does not touch the proxy flag must not throw away a perfectly
// good HTTP client.
func TestReloadKeepsDerivedClientsWhenNothingRelevantChanged(t *testing.T) {
first := publishConfig(t, `{
"controllerMethod": {
"webhook": { "enabled": true, "url": "https://example.com/one", "networkUseProxy": false },
"ntfy": { "enabled": true, "topic": "t", "networkUseProxy": false }
}
}`)
webhook := NewWebhookController(first.ControllerMethod.Webhook)
ntfy := NewNtfyController(first.ControllerMethod.Ntfy)
webhookClient, ntfyClient := webhook.client(), ntfy.client()
// The whole point of an unchanged reload: the settings object is equal, so
// nothing is rebuilt.
publishConfig(t, `{
"controllerMethod": {
"webhook": { "enabled": true, "url": "https://example.com/one", "networkUseProxy": false },
"ntfy": { "enabled": true, "topic": "t", "networkUseProxy": false }
}
}`)
webhook.Reload()
ntfy.Reload()
if webhook.client() != webhookClient {
t.Error("webhook client was rebuilt although its settings did not change")
}
if ntfy.client() != ntfyClient {
t.Error("ntfy client was rebuilt although its settings did not change")
}
// A change that leaves the proxy flag alone still keeps the client.
publishConfig(t, `{
"controllerMethod": {
"webhook": { "enabled": true, "url": "https://example.com/two", "networkUseProxy": false },
"ntfy": { "enabled": true, "topic": "t", "networkUseProxy": false }
}
}`)
webhook.Reload()
if webhook.client() != webhookClient {
t.Error("webhook client was rebuilt for a change that did not touch the proxy setting")
}
if got := webhook.settings().URL; got != "https://example.com/two" {
t.Errorf("Reload did not pick up the new URL: %q", got)
}
// Flipping the proxy flag is the one change that must rebuild it.
publishConfig(t, `{
"controllerMethod": {
"webhook": { "enabled": true, "url": "https://example.com/two", "networkUseProxy": true },
"ntfy": { "enabled": true, "topic": "t", "networkUseProxy": false }
}
}`)
webhook.Reload()
if webhook.client() == webhookClient {
t.Error("webhook client was not rebuilt after the proxy setting changed")
}
}
// dialerIsDefault reports whether gomail still dials directly. Function values
// are not comparable in Go, so the code pointers are compared.
func dialerIsDefault() bool {
return reflect.ValueOf(gomail.NetDialTimeout).Pointer() ==
reflect.ValueOf(net.DialTimeout).Pointer()
}
// TestEmailControllerReloadResetsDialer covers the case the constructor used to
// get wrong: gomail's dial hook is a package-level variable with no "unset", so
// turning networkUseProxy off has to actively restore the default dialer rather
// than leave SMTP tunnelled through a proxy nobody asked for.
func TestEmailControllerReloadResetsDialer(t *testing.T) {
t.Cleanup(func() { gomail.NetDialTimeout = net.DialTimeout })
first := publishConfig(t, `{
"system": { "networkProxy": "http://127.0.0.1:7890" },
"controllerMethod": {
"email": { "enabled": true, "networkUseProxy": true, "smtpHost": "smtp.example.com", "to": ["a@example.com"] }
}
}`)
ctrl := NewEmailController(first.ControllerMethod.Email)
if dialerIsDefault() {
t.Fatal("the proxy dialer was not installed for networkUseProxy: true")
}
publishConfig(t, `{
"system": { "networkProxy": "http://127.0.0.1:7890" },
"controllerMethod": {
"email": { "enabled": true, "networkUseProxy": false, "smtpHost": "smtp.example.com", "to": ["a@example.com"] }
}
}`)
ctrl.Reload()
if !dialerIsDefault() {
t.Error("turning networkUseProxy off left gomail's proxy dialer in place")
}
}
// TestWebhookReloadDuringSend drives Reload while notifications are being
// delivered. Run with -race: the settings and the HTTP client are swapped on
// the settings-update goroutine while the send path reads them on others, which
// is exactly the race the atomic pointers exist to prevent.
func TestWebhookReloadDuringSend(t *testing.T) {
var delivered int64
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt64(&delivered, 1)
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
body := func(proxy bool) string {
return fmt.Sprintf(`{
"controllerMethod": {
"webhook": { "enabled": true, "url": %q, "method": "POST", "networkUseProxy": %t }
},
"controllerMessage": { "SERVER_STATUS_CHANGED": "{{ serverName }} {{ event }}" }
}`, srv.URL, proxy)
}
direct := writeConfigFile(t, body(false))
proxied := writeConfigFile(t, body(true))
if _, err := config.LoadGlobalConfig(direct); err != nil {
t.Fatalf("LoadGlobalConfig: %v", err)
}
ctrl := NewWebhookController(config.Current().ControllerMethod.Webhook)
change := node.StatusChange{Event: "Online", UUID: "u1", Name: "alpha"}
var readers, writers sync.WaitGroup
stop := make(chan struct{})
for i := 0; i < 3; i++ {
readers.Add(1)
go func() {
defer readers.Done()
for {
select {
case <-stop:
return
default:
}
if err := ctrl.SendStatusChange(change); err != nil {
// The test server never fails; a transport error here means
// the swap lost a field mid-flight.
t.Errorf("SendStatusChange: %v", err)
return
}
_ = ctrl.IsEnabled()
_ = ctrl.IsMarkdown()
}
}()
}
writers.Add(1)
go func() {
defer writers.Done()
for i := 0; i < 30; i++ {
path := direct
if i%2 == 0 {
path = proxied
}
if _, err := config.LoadGlobalConfig(path); err != nil {
t.Errorf("LoadGlobalConfig: %v", err)
return
}
ctrl.Reload()
}
}()
// Let the writer finish, then release the readers: they only return once
// stop is closed.
writers.Wait()
close(stop)
readers.Wait()
if atomic.LoadInt64(&delivered) == 0 {
t.Error("no notification reached the test server")
}
}
@@ -124,6 +124,21 @@ func (t *TelegramController) IsEnabled() bool {
return t.cfg.Enabled
}
// Reload is not implemented yet for this channel.
//
// The bot client is built from the token and the polling context is created
// with the controller, so a stopped controller cannot be restarted in place:
// applying a change means building a replacement controller and swapping it
// into the manager. Until then the divergence is reported rather than ignored,
// so an edit that needs a restart does not look like it took effect.
func (t *TelegramController) Reload() {
global := config.Current()
if global == nil || t.cfg == global.ControllerMethod.Telegram {
return
}
postLog.Warning("Telegram settings changed but hot reload is not implemented for this channel yet; restart to apply them")
}
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (t *TelegramController) IsMarkdown() bool {
+73 -22
View File
@@ -5,6 +5,8 @@ import (
"encoding/json"
"fmt"
"net/http"
"reflect"
"sync/atomic"
"time"
"nukumizu-backend/config"
@@ -16,17 +18,33 @@ import (
)
// WebhookController handles notifications via generic HTTP webhooks.
//
// Both fields are held behind atomic pointers rather than in plain fields:
// Reload replaces them on the settings-update goroutine while the send methods
// read them on the status-change and incoming-webhook goroutines.
type WebhookController struct {
cfg config.WebhookConfig
httpClient *http.Client
cfg atomic.Pointer[config.WebhookConfig]
httpClient atomic.Pointer[http.Client]
}
// NewWebhookController creates a new Webhook controller.
func NewWebhookController(cfg config.WebhookConfig) *WebhookController {
return &WebhookController{
cfg: cfg,
httpClient: netproxy.HTTPClient(cfg.NetworkUseProxy, 10*time.Second),
}
w := &WebhookController{}
w.cfg.Store(&cfg)
w.httpClient.Store(netproxy.HTTPClient(cfg.NetworkUseProxy, 10*time.Second))
return w
}
// settings returns the configuration currently in effect. The value it points
// at is never mutated after being stored, so a caller can hold the pointer for
// one whole send without a concurrent Reload disturbing it.
func (w *WebhookController) settings() *config.WebhookConfig {
return w.cfg.Load()
}
// client returns the HTTP client built for the current proxy setting.
func (w *WebhookController) client() *http.Client {
return w.httpClient.Load()
}
// Name returns the controller name.
@@ -36,7 +54,7 @@ func (w *WebhookController) Name() string {
// Start initializes the Webhook controller.
func (w *WebhookController) Start() error {
if !w.cfg.Enabled {
if !w.settings().Enabled {
postLog.Info("Webhook controller is disabled")
return nil
}
@@ -49,26 +67,51 @@ func (w *WebhookController) Stop() {
postLog.Info("Webhook controller stopped")
}
// Reload applies the current configuration. Everything the sender reads is
// taken from the settings at send time, so only a change to the proxy flag
// needs more than the swap: that one is baked into the HTTP client's transport
// when the client is built.
func (w *WebhookController) Reload() {
global := config.Current()
if global == nil {
return
}
current := w.settings()
updated := global.ControllerMethod.Webhook
// WebhookConfig carries a header map, so it is not comparable with ==.
if reflect.DeepEqual(*current, updated) {
return
}
if current.NetworkUseProxy != updated.NetworkUseProxy {
w.httpClient.Store(netproxy.HTTPClient(updated.NetworkUseProxy, 10*time.Second))
}
w.cfg.Store(&updated)
postLog.Debug("Webhook controller reloaded")
}
// IsEnabled returns whether the controller is enabled.
func (w *WebhookController) IsEnabled() bool {
return w.cfg.Enabled
return w.settings().Enabled
}
// IsMarkdown returns whether the channel renders Markdown, per its markdown
// setting in config.json.
func (w *WebhookController) IsMarkdown() bool {
return w.cfg.Markdown
return w.settings().Markdown
}
// SendStatusChange sends a status change notification via Webhook.
func (w *WebhookController) SendStatusChange(change node.StatusChange) error {
if !w.cfg.Enabled {
s := w.settings()
if !s.Enabled {
return nil
}
cfg := config.Current()
params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, w.cfg.Markdown)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, s.Markdown)
payload := map[string]interface{}{
"event": change.Event,
@@ -83,13 +126,14 @@ func (w *WebhookController) SendStatusChange(change node.StatusChange) error {
// SendServerList sends the server list via Webhook.
func (w *WebhookController) SendServerList(onlineServers, offlineServers string) error {
if !w.cfg.Enabled {
s := w.settings()
if !s.Enabled {
return nil
}
cfg := config.Current()
params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params, w.cfg.Markdown)
message := template.Render(cfg.ControllerMessage.ServerList, params, s.Markdown)
payload := map[string]interface{}{
"type": "serverList",
@@ -104,13 +148,14 @@ func (w *WebhookController) SendServerList(onlineServers, offlineServers string)
// SendExecuteResult sends a command execution result via Webhook.
func (w *WebhookController) SendExecuteResult(serverName, serverUUID, command, result string) error {
if !w.cfg.Enabled {
s := w.settings()
if !s.Enabled {
return nil
}
cfg := config.Current()
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, w.cfg.Markdown)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, s.Markdown)
payload := map[string]interface{}{
"type": "executeResult",
@@ -128,7 +173,8 @@ func (w *WebhookController) SendExecuteResult(serverName, serverUUID, command, r
// 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 {
s := w.settings()
if !s.Enabled {
return nil
}
@@ -137,15 +183,20 @@ func (w *WebhookController) SendAlert(alert controller.Alert) error {
"subject": alert.Subject,
"source": alert.Source,
"content": alert.Content,
"message": alert.Render(w.cfg.Markdown),
"message": alert.Render(s.Markdown),
"time": alert.Time,
}
return w.send(payload)
}
// send posts one payload using the settings in effect at the moment it is
// called, so the URL, method and headers all come from the same configuration
// even if a reload lands mid-send.
func (w *WebhookController) send(payload map[string]interface{}) error {
method := w.cfg.Method
s := w.settings()
method := s.Method
if method == "" {
method = "POST"
}
@@ -155,17 +206,17 @@ func (w *WebhookController) send(payload map[string]interface{}) error {
return fmt.Errorf("failed to marshal webhook payload: %w", err)
}
req, err := http.NewRequest(method, w.cfg.URL, bytes.NewReader(bodyJSON))
req, err := http.NewRequest(method, s.URL, bytes.NewReader(bodyJSON))
if err != nil {
return fmt.Errorf("failed to create webhook request: %w", err)
}
req.Header.Set("Content-Type", "application/json")
for key, value := range w.cfg.Headers {
for key, value := range s.Headers {
req.Header.Set(key, value)
}
resp, err := w.httpClient.Do(req)
resp, err := w.client().Do(req)
if err != nil {
postLog.Warning("Failed to send webhook: " + err.Error())
return err
@@ -176,6 +227,6 @@ func (w *WebhookController) send(payload map[string]interface{}) error {
postLog.Warning(fmt.Sprintf("Webhook returned status %d", resp.StatusCode))
}
postLog.Debug("Webhook notification sent to " + w.cfg.URL)
postLog.Debug("Webhook notification sent to " + s.URL)
return nil
}