diff --git a/internal/controller/controller.go b/internal/controller/controller.go index 3297449..3cd8cee 100644 --- a/internal/controller/controller.go +++ b/internal/controller/controller.go @@ -88,6 +88,15 @@ type Controller interface { // IsMarkdown reports whether the channel renders Markdown, per its own // "markdown" setting in config.json. IsMarkdown() bool + // Reload applies the current configuration to a running controller, so a + // settings update takes effect without a restart. It reads the live + // configuration itself and rebuilds whatever it derived from the old one. + // + // Reload runs on the goroutine serving the settings update while the send + // methods may be running on others, so an implementation must not replace + // state those methods read without synchronizing (see the atomic pointers in + // the notification pipes). + Reload() SendStatusChange(change node.StatusChange) error SendServerList(onlineServers, offlineServers string) error SendExecuteResult(serverName, serverUUID, command, result string) error @@ -134,6 +143,25 @@ func (m *Manager) Register(c Controller) { postLog.Info("Controller registered: " + c.Name()) } +// ReloadAll applies the current configuration to every registered controller, +// so a settings update reaches the channels without a restart. +// +// The controllers are collected under the registry lock and reloaded outside +// it: a reload can rebuild a client and block on the network, and holding m.mu +// across that would stall every notification for its duration. +func (m *Manager) ReloadAll() { + m.mu.RLock() + controllers := make([]Controller, 0, len(m.controllers)) + for _, ctrl := range m.controllers { + controllers = append(controllers, ctrl) + } + m.mu.RUnlock() + + for _, ctrl := range controllers { + ctrl.Reload() + } +} + // ShowBotInitMessage sends the bot initialization message to all enabled // bot controllers (QQ/NapCat and Telegram). Notification-only pipes that do // not implement BotController are skipped. The message is typed diff --git a/internal/controller/pipes/email.go b/internal/controller/pipes/email.go index c95e593..530f660 100644 --- a/internal/controller/pipes/email.go +++ b/internal/controller/pipes/email.go @@ -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 } diff --git a/internal/controller/pipes/ntfy.go b/internal/controller/pipes/ntfy.go index 1da7dbc..fc574ce 100644 --- a/internal/controller/pipes/ntfy.go +++ b/internal/controller/pipes/ntfy.go @@ -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 } diff --git a/internal/controller/pipes/qq_napcat/qq.go b/internal/controller/pipes/qq_napcat/qq.go index 21da416..4b05d98 100644 --- a/internal/controller/pipes/qq_napcat/qq.go +++ b/internal/controller/pipes/qq_napcat/qq.go @@ -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 { diff --git a/internal/controller/pipes/reload_test.go b/internal/controller/pipes/reload_test.go new file mode 100644 index 0000000..e6e5cab --- /dev/null +++ b/internal/controller/pipes/reload_test.go @@ -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") + } +} diff --git a/internal/controller/pipes/telegram/telegram.go b/internal/controller/pipes/telegram/telegram.go index a082db9..ef58935 100644 --- a/internal/controller/pipes/telegram/telegram.go +++ b/internal/controller/pipes/telegram/telegram.go @@ -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 { diff --git a/internal/controller/pipes/webhook.go b/internal/controller/pipes/webhook.go index 8392ec0..8517aaf 100644 --- a/internal/controller/pipes/webhook.go +++ b/internal/controller/pipes/webhook.go @@ -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 } diff --git a/main.go b/main.go index caddf82..5b1e15f 100644 --- a/main.go +++ b/main.go @@ -58,12 +58,17 @@ func main() { postLog.SetDebugMode(cfg.System.DebugMode) postLog.InitLogBroadcaster() - // The logger's debug flag is a process-wide setting rather than something - // read at every log call, so a settings update has to push the new value - // into it. This is the first link in the reload hook chain; the controllers - // join it as they gain reload support. + // A settings update replaces the configuration in memory; these hooks push + // the new values into the state that was derived from the old one. The + // logger's debug flag is process-wide rather than read at every log call, + // and each controller holds its own copy of its channel settings plus the + // clients built from them. config.OnReload(func(updated *config.Config) { postLog.SetDebugMode(updated.System.DebugMode) + + if mgr := controller.GetManager(); mgr != nil { + mgr.ReloadAll() + } }) dbPath := cfg.DBPath