diff --git a/README.md b/README.md index 0a36be9..0b55ea3 100644 --- a/README.md +++ b/README.md @@ -294,6 +294,18 @@ Admins and trusted groups are defined **per bot channel** and map a member ID to | `{{ list.offlineServers }}` | Formatted list of offline servers | | `{{ softwareVersion }}`, `{{ softwareBuildVer }}`, `{{ softwareCommitHash }}`, `{{ softwareBuildType }}`, `{{ softwareBuildTime }}`, `{{ softwareDeveloper }}`, `{{ softwareDescription }}` | Build metadata (commit hash and build time are injected at compile time) | +### What applies without a restart + +Saving settings applies most of them immediately. How each group takes effect: + +| Settings | How it applies | +|---|---| +| `controllerMethod` (all five channels) | Every channel is **rebuilt**: the running controllers are stopped and a fresh set is built from the new settings. A channel that owns a connection reconnects — Telegram re-runs its `getMe` handshake, NapCat opens a new WebSocket — so notifications sent during the swap are lost. The rebuild only happens when this section actually changed; saving a message template does not disturb the channels. | +| `controllerMessage`, `debug`, `bot_user_config.json`, `bot_node_config.json`, `webhook.endpoints` | Picked up as they are used; nothing is restarted. | +| `system.debugMode` | Applies to both behavior and log filtering. | +| `system.networkProxy` | Applies to connections opened afterwards. Channels that were built with `networkUseProxy` capture the proxy when they connect, so the address is only re-read on the next rebuild; the controller rebuild above does that. | +| `system.listenAddr` / `listenPort`, `webhook.enabled` / `listenAddr` / `listenPort`, `dataPath`, `dbPath`, `komari.dashboardURL` | **Applied at startup only.** Saving them changes the file and the in-memory configuration but not the running listener, database or Komari client — restart to apply. | + ## API Success responses follow the envelope `{"success": true, "message": "...", "data": {...}}`, with the payload nested under a single `data` key. Error responses use `{"success": false, "message": "..."}`. `message` may be omitted on success when there is nothing to report. @@ -322,7 +334,7 @@ Browser WebSocket handshakes cannot carry custom headers, so `/api/system/getLog | `/api/server/getStatus` | GET | admin | Live server status (mirrors the Bot's `/status`). Query `?uuid=` (or `all`). Returns `data: {: {uuid, name, online, report}}`; `report` is `null` when the node has not reported yet. `404` for an unknown single uuid. | | `/api/server/exec` | POST | bot / admin | Execute a command. Body `{uuid: [...], command}`. Dispatches a Komari task and polls until completion (or timeout). Returns `data: {taskID, results}`. | | `/api/settings/get` | GET | admin | `?type=global\|bot_user_config\|bot_node_config` | Returns `data: {config}`, where `config` is the selected config file's content (same layout as the JSON file). | -| `/api/settings/set` | POST | admin | `?type=` + JSON body of partial updates, e.g. `{"system":{"debugMode":true}}` | Deep-merges the body into the selected config file, persists it, and reloads it in memory. Only the given keys change; arrays replace. | +| `/api/settings/set` | POST | admin | `?type=` + JSON body of partial updates, e.g. `{"system":{"debugMode":true}}` | Deep-merges the body into the selected config file, persists it, and reloads it in memory. Only the given keys change; arrays replace. See [What applies without a restart](#what-applies-without-a-restart). | | `/api/webhook/add` | POST | admin | Add an incoming webhook endpoint. Body `{name, enabled?, token?, notifyPipes?}` — only the fields given are stored, the rest start at their defaults. `409` when the name is already configured. | | `/api/webhook/modify` | POST | admin | Change an existing endpoint. Body `{name, ...}` — the fields given are the fields that change (same partial-update rule as `/api/settings/set`, but scoped to one endpoint). `404` for an unknown name, `400` when no other field is given. | | `/api/webhook/delete` | POST | admin | Remove an endpoint. Body `{name}`. `404` for an unknown name. | diff --git a/config/reload.go b/config/reload.go index 61bac53..ba6cd15 100644 --- a/config/reload.go +++ b/config/reload.go @@ -16,6 +16,12 @@ import ( var ( reloadHooksMu sync.Mutex reloadHooks []func(*Config) + + // reloadRunMu serializes hook execution. A hook mutates process-wide state + // — the logger's debug flag, the controller registry — so two overlapping + // settings updates must not run them at the same time, or both would decide + // to rebuild the controllers from their own view of what was applied. + reloadRunMu sync.Mutex ) // OnReload registers a hook to run after every reload of the global @@ -47,13 +53,15 @@ func notifyReload(cfg *Config) { return } - // Copy under the lock, then run outside it: a hook is free to register - // another hook without deadlocking. + // Copy under the registry lock, then release it before running anything: a + // hook is free to register another hook without deadlocking. reloadHooksMu.Lock() hooks := make([]func(*Config), len(reloadHooks)) copy(hooks, reloadHooks) reloadHooksMu.Unlock() + reloadRunMu.Lock() + defer reloadRunMu.Unlock() for _, hook := range hooks { runReloadHook(hook, cfg) } diff --git a/config/reload_test.go b/config/reload_test.go index 1162c17..1f9ca31 100644 --- a/config/reload_test.go +++ b/config/reload_test.go @@ -1,6 +1,7 @@ package config import ( + "reflect" "testing" "nukumizu-backend/global" @@ -90,6 +91,40 @@ func TestReloadHookIsolatesPanic(t *testing.T) { } } +// TestUnrelatedUpdateLeavesControllerMethodAlone guards the trigger for a +// controller rebuild. Whether to rebuild is decided by comparing the whole +// controllerMethod section with the one the running controllers were built +// from, so reloading the file has to reproduce that section byte for byte. A +// default applied inconsistently — a nil recipient slice turned into an empty +// one on the second load, say — would make every settings edit look like a +// controller change and tear down every channel on each save. +func TestUnrelatedUpdateLeavesControllerMethodAlone(t *testing.T) { + writeTempConfig(t, &global.ConfigPath.Global, `{ + "controllerMethod": { + "qq(napcat)": { "enabled": false }, + "email": { "enabled": false, "to": [] }, + "webhook": { "enabled": false, "headers": {} } + } +}`) + + first, err := LoadGlobalConfig(global.ConfigPath.Global) + if err != nil { + t.Fatalf("LoadGlobalConfig: %v", err) + } + before := first.ControllerMethod + + if err := UpdateSettings(SettingGlobal, map[string]interface{}{ + "controllerMessage": map[string]interface{}{"BOT_STARTED": "hello"}, + }); err != nil { + t.Fatalf("UpdateSettings: %v", err) + } + + after := Current().ControllerMethod + if !reflect.DeepEqual(before, after) { + t.Errorf("an unrelated update changed the controllerMethod section:\nbefore: %+v\nafter: %+v", before, after) + } +} + // TestReloadHookSeesEveryWriterPath pins the invariant that each writer of // config.json notifies, not just /api/settings/set: the webhook endpoint // helpers go through the same channel. diff --git a/internal/controller/controller.go b/internal/controller/controller.go index 3cd8cee..e2156c6 100644 --- a/internal/controller/controller.go +++ b/internal/controller/controller.go @@ -3,8 +3,10 @@ package controller import ( "errors" "fmt" + "reflect" "strings" "sync" + "sync/atomic" "nukumizu-backend/config" "nukumizu-backend/internal/node" @@ -88,15 +90,6 @@ 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 @@ -118,6 +111,12 @@ type BotController interface { type Manager struct { mu sync.RWMutex controllers map[string]Controller + + // builtFrom records the controllerMethod section the registered controllers + // were built from, so a settings update that concerns them can be told apart + // from one that does not. It is read on the settings-update goroutine and + // written when the set is replaced. + builtFrom atomic.Pointer[config.ControllerMethodConfig] } var globalManager *Manager @@ -135,30 +134,69 @@ func GetManager() *Manager { return globalManager } -// Register adds a controller to the manager. -func (m *Manager) Register(c Controller) { - m.mu.Lock() - defer m.mu.Unlock() - m.controllers[c.Name()] = c - postLog.Info("Controller registered: " + c.Name()) +// NeedsRebuild reports whether next differs from the controllerMethod section +// the registered controllers were built from. A manager with no controllers yet +// always reports true, so the first call installs the initial set. +// +// Controllers are rebuilt wholesale rather than reconfigured in place: each one +// reads its settings into fields at construction, and two of them own +// connections that cannot be re-pointed (the NapCat WebSocket listener is +// stopped through a sync.Once, the Telegram polling context is created with the +// controller). Replacing the set keeps every channel on the same footing. +func (m *Manager) NeedsRebuild(next config.ControllerMethodConfig) bool { + built := m.builtFrom.Load() + if built == nil { + return true + } + // The section carries a header map and a recipient slice, so it is not + // comparable with ==. + return !reflect.DeepEqual(*built, next) } -// ReloadAll applies the current configuration to every registered controller, -// so a settings update reaches the channels without a restart. +// ReplaceAll stops every registered controller and swaps in next, which the +// caller built from method. It is the only way controllers are installed, at +// startup and after a settings change alike. // -// 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) +// The swap happens under the registry lock so routing flips to the new set +// atomically; stopping and starting happen outside it. Both can block — Stop +// closes sockets, Start performs a handshake — and holding m.mu across them +// would stall every notification for the duration. +func (m *Manager) ReplaceAll(next []Controller, method config.ControllerMethodConfig) { + m.mu.Lock() + previous := m.controllers + m.controllers = make(map[string]Controller, len(next)) + for _, ctrl := range next { + m.controllers[ctrl.Name()] = ctrl } - m.mu.RUnlock() + m.mu.Unlock() - for _, ctrl := range controllers { - ctrl.Reload() + m.builtFrom.Store(&method) + + names := make([]string, 0, len(next)) + for _, ctrl := range next { + names = append(names, ctrl.Name()) + } + postLog.Info("Controller set installed: " + strings.Join(names, ", ")) + + for _, ctrl := range previous { + ctrl.Stop() + } + + // Start off the calling goroutine, the way startup does: Telegram's getMe + // and the NapCat WebSocket handshake would otherwise hold the settings + // request open for as long as they take. The new controllers are already + // routable, and each one can send before Start returns. + for _, ctrl := range next { + go func(ctrl Controller) { + defer func() { + if r := recover(); r != nil { + postLog.Error(fmt.Sprintf("Controller %s panicked on start: %v", ctrl.Name(), r)) + } + }() + if err := ctrl.Start(); err != nil { + postLog.Error(fmt.Sprintf("Failed to start controller %s: %v", ctrl.Name(), err)) + } + }(ctrl) } } diff --git a/internal/controller/controller_test.go b/internal/controller/controller_test.go new file mode 100644 index 0000000..1585f54 --- /dev/null +++ b/internal/controller/controller_test.go @@ -0,0 +1,222 @@ +package controller + +import ( + "sync" + "testing" + "time" + + "nukumizu-backend/config" + "nukumizu-backend/internal/node" +) + +// fakeController records the lifecycle calls it receives so a test can assert +// what ReplaceAll did to it. Start signals on started, because ReplaceAll starts +// controllers off the calling goroutine. +type fakeController struct { + name string + started chan struct{} + + mu sync.Mutex + starts int + stops int +} + +func newFake(name string) *fakeController { + return &fakeController{name: name, started: make(chan struct{}, 4)} +} + +func (f *fakeController) Name() string { return f.name } +func (f *fakeController) IsEnabled() bool { return true } +func (f *fakeController) IsMarkdown() bool { return false } + +func (f *fakeController) Start() error { + f.mu.Lock() + f.starts++ + f.mu.Unlock() + select { + case f.started <- struct{}{}: + default: + } + return nil +} + +func (f *fakeController) Stop() { + f.mu.Lock() + defer f.mu.Unlock() + f.stops++ +} + +func (f *fakeController) lifecycle() (starts, stops int) { + f.mu.Lock() + defer f.mu.Unlock() + return f.starts, f.stops +} + +func (f *fakeController) SendStatusChange(node.StatusChange) error { return nil } +func (f *fakeController) SendServerList(string, string) error { return nil } +func (f *fakeController) SendExecuteResult(string, string, string, string) error { return nil } +func (f *fakeController) SendAlert(Alert) error { return nil } + +// waitStarted blocks until the controller's Start has run. +func (f *fakeController) waitStarted(t *testing.T) { + t.Helper() + select { + case <-f.started: + case <-time.After(5 * time.Second): + t.Fatalf("controller %s was never started", f.name) + } +} + +func newTestManager() *Manager { + return &Manager{controllers: make(map[string]Controller)} +} + +// routedTo returns the controller the manager currently routes the given name +// to. +func (m *Manager) routedTo(name string) Controller { + m.mu.RLock() + defer m.mu.RUnlock() + return m.controllers[name] +} + +func TestReplaceAllInstallsAndStarts(t *testing.T) { + m := newTestManager() + alpha, beta := newFake("alpha"), newFake("beta") + + m.ReplaceAll([]Controller{alpha, beta}, config.ControllerMethodConfig{}) + + alpha.waitStarted(t) + beta.waitStarted(t) + + if got := m.routedTo("alpha"); got != Controller(alpha) { + t.Errorf("alpha is not routable after ReplaceAll: %v", got) + } + if got := m.routedTo("beta"); got != Controller(beta) { + t.Errorf("beta is not routable after ReplaceAll: %v", got) + } +} + +// TestReplaceAllStopsTheOutgoingSet is the invariant the rebuild relies on: the +// old controllers must be shut down, or a rebuilt NapCat or Telegram controller +// would leave its previous connection running. +func TestReplaceAllStopsTheOutgoingSet(t *testing.T) { + m := newTestManager() + outgoing := newFake("alpha") + m.ReplaceAll([]Controller{outgoing}, config.ControllerMethodConfig{}) + outgoing.waitStarted(t) + + incoming := newFake("alpha") + m.ReplaceAll([]Controller{incoming}, config.ControllerMethodConfig{}) + incoming.waitStarted(t) + + if starts, stops := outgoing.lifecycle(); starts != 1 || stops != 1 { + t.Errorf("outgoing controller lifecycle = %d starts / %d stops, want 1/1", starts, stops) + } + if _, stops := incoming.lifecycle(); stops != 0 { + t.Errorf("incoming controller was stopped %d times", stops) + } + if got := m.routedTo("alpha"); got != Controller(incoming) { + t.Error("routing still points at the outgoing controller") + } +} + +func TestReplaceAllDropsChannelsLeftOut(t *testing.T) { + m := newTestManager() + m.ReplaceAll([]Controller{newFake("alpha"), newFake("beta")}, config.ControllerMethodConfig{}) + + m.ReplaceAll([]Controller{newFake("beta")}, config.ControllerMethodConfig{}) + + if got := m.routedTo("alpha"); got != nil { + t.Errorf("a channel missing from the new set is still routable: %v", got) + } + if got := m.routedTo("beta"); got == nil { + t.Error("the surviving channel is not routable") + } +} + +func TestNeedsRebuild(t *testing.T) { + m := newTestManager() + base := config.ControllerMethodConfig{ + Email: config.EmailConfig{Enabled: true, SMTPPort: 587, To: []string{"a@example.com"}}, + } + + if !m.NeedsRebuild(base) { + t.Error("a manager with nothing installed must report that a rebuild is needed") + } + + m.ReplaceAll(nil, base) + if m.NeedsRebuild(base) { + t.Error("the settings the set was built from must not ask for another rebuild") + } + + changed := base + changed.Email.Enabled = false + if !m.NeedsRebuild(changed) { + t.Error("a changed email setting must ask for a rebuild") + } + + // The section carries maps and slices, so it is compared by value rather + // than by identity: equal content must not trigger a rebuild. + withHeaders := config.ControllerMethodConfig{ + Webhook: config.WebhookConfig{Headers: map[string]string{"X-Token": "t"}}, + } + m.ReplaceAll(nil, withHeaders) + equalHeaders := config.ControllerMethodConfig{ + Webhook: config.WebhookConfig{Headers: map[string]string{"X-Token": "t"}}, + } + if m.NeedsRebuild(equalHeaders) { + t.Error("equal header maps must not ask for a rebuild") + } + + differentHeaders := config.ControllerMethodConfig{ + Webhook: config.WebhookConfig{Headers: map[string]string{"X-Token": "other"}}, + } + if !m.NeedsRebuild(differentHeaders) { + t.Error("a changed header must ask for a rebuild") + } +} + +// TestReplaceAllDuringNotification drives ReplaceAll while notifications are +// being routed. Run with -race: the swap replaces the map the routing path +// reads, which is what the registry lock exists to make safe. +func TestReplaceAllDuringNotification(t *testing.T) { + m := newTestManager() + m.ReplaceAll([]Controller{newFake("alpha")}, config.ControllerMethodConfig{}) + + 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: + } + m.NotifyStatusChange(node.StatusChange{UUID: "u1", Name: "alpha", Event: "Online"}) + _ = m.IsMarkdown("alpha") + } + }() + } + + writers.Add(1) + go func() { + defer writers.Done() + for i := 0; i < 20; i++ { + m.ReplaceAll([]Controller{newFake("alpha"), newFake("beta")}, config.ControllerMethodConfig{}) + } + }() + + // Let the writer finish, then release the readers: they only return once + // stop is closed. + writers.Wait() + close(stop) + readers.Wait() + + if got := m.routedTo("beta"); got == nil { + t.Error("the last installed set is not routable") + } +} diff --git a/internal/controller/pipes/email.go b/internal/controller/pipes/email.go index 530f660..7b7133b 100644 --- a/internal/controller/pipes/email.go +++ b/internal/controller/pipes/email.go @@ -3,8 +3,6 @@ package pipes import ( "fmt" "net" - "reflect" - "sync/atomic" gomail "gopkg.in/mail.v2" @@ -18,26 +16,17 @@ 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. +// cfg is written once, by the constructor, and never again: a controller that +// needs different settings is replaced wholesale by the manager rather than +// reconfigured in place, so the send methods can read it without locking. type EmailController struct { - cfg atomic.Pointer[config.EmailConfig] + cfg config.EmailConfig } // NewEmailController creates a new Email controller. func NewEmailController(cfg config.EmailConfig) *EmailController { - 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() + return &EmailController{cfg: cfg} } // applyEmailProxy routes SMTP through the HTTP CONNECT proxy, or restores a @@ -61,7 +50,7 @@ func (e *EmailController) Name() string { // Start initializes the Email controller. func (e *EmailController) Start() error { - if !e.settings().Enabled { + if !e.cfg.Enabled { postLog.Info("Email controller is disabled") return nil } @@ -74,121 +63,90 @@ 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.settings().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.settings().Markdown + return e.cfg.Markdown } // SendStatusChange sends a status change notification via Email. func (e *EmailController) SendStatusChange(change node.StatusChange) error { - s := e.settings() - if !s.Enabled { + if !e.cfg.Enabled { return nil } - if len(s.To) == 0 { + if len(e.cfg.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, s.Markdown) + body := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, e.cfg.Markdown) subject := fmt.Sprintf("Server Status Change: %s - %s", change.Name, change.Event) - return e.sendEmail(s, subject, body) + return e.sendEmail(subject, body) } // SendServerList sends the server list via Email. func (e *EmailController) SendServerList(onlineServers, offlineServers string) error { - s := e.settings() - if !s.Enabled || len(s.To) == 0 { + if !e.cfg.Enabled || len(e.cfg.To) == 0 { return nil } cfg := config.Current() params := template.BuildParamsFromServerList() - body := template.Render(cfg.ControllerMessage.ServerList, params, s.Markdown) + body := template.Render(cfg.ControllerMessage.ServerList, params, e.cfg.Markdown) - return e.sendEmail(s, "Server List", body) + return e.sendEmail("Server List", body) } // SendExecuteResult sends a command execution result via Email. func (e *EmailController) SendExecuteResult(serverName, serverUUID, command, result string) error { - s := e.settings() - if !s.Enabled || len(s.To) == 0 { + if !e.cfg.Enabled || len(e.cfg.To) == 0 { return nil } cfg := config.Current() params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) - body := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, s.Markdown) + body := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, e.cfg.Markdown) subject := fmt.Sprintf("Command Result: %s on %s", command, serverName) - return e.sendEmail(s, 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 { - s := e.settings() - if !s.Enabled { + if !e.cfg.Enabled { return nil } - if len(s.To) == 0 { + if len(e.cfg.To) == 0 { postLog.Debug("Email controller has no recipients configured") return nil } - return e.sendEmail(s, alert.Subject, alert.Render(s.Markdown)) + return e.sendEmail(alert.Subject, alert.Render(e.cfg.Markdown)) } -// 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 { +func (e *EmailController) sendEmail(subject, body string) error { m := gomail.NewMessage() - m.SetHeader("From", s.From) - m.SetHeader("To", s.To...) + m.SetHeader("From", e.cfg.From) + m.SetHeader("To", e.cfg.To...) m.SetHeader("Subject", subject) m.SetBody("text/plain", body) - d := gomail.NewDialer(s.SMTPHost, s.SMTPPort, s.Username, s.Password) + d := gomail.NewDialer(e.cfg.SMTPHost, e.cfg.SMTPPort, e.cfg.Username, e.cfg.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", s.To)) + postLog.Debug("Email sent successfully to " + fmt.Sprintf("%v", e.cfg.To)) return nil } diff --git a/internal/controller/pipes/email_proxy_test.go b/internal/controller/pipes/email_proxy_test.go new file mode 100644 index 0000000..fdf8923 --- /dev/null +++ b/internal/controller/pipes/email_proxy_test.go @@ -0,0 +1,54 @@ +package pipes + +import ( + "net" + "reflect" + "testing" + + gomail "gopkg.in/mail.v2" + + "nukumizu-backend/config" +) + +// dialerIsDefault reports whether gomail still dials directly. Function values +// are not comparable in Go, so the code pointers are compared instead. +func dialerIsDefault() bool { + return reflect.ValueOf(gomail.NetDialTimeout).Pointer() == + reflect.ValueOf(net.DialTimeout).Pointer() +} + +// TestApplyEmailProxyRestoresDefaultDialer covers the case the constructor used +// to get wrong. gomail exposes its dial hook as a package-level variable with no +// unset, and the code only installed the proxy dialer when the flag was set, so +// turning networkUseProxy back off left SMTP tunnelled through a proxy nobody +// had asked for — with no way to undo it short of a restart. +func TestApplyEmailProxyRestoresDefaultDialer(t *testing.T) { + t.Cleanup(func() { gomail.NetDialTimeout = net.DialTimeout }) + + applyEmailProxy(true) + if dialerIsDefault() { + t.Fatal("applyEmailProxy(true) did not install the proxy dialer") + } + + applyEmailProxy(false) + if !dialerIsDefault() { + t.Error("applyEmailProxy(false) left the proxy dialer in place") + } +} + +// TestNewEmailControllerAppliesProxySetting pins that the dialer follows the +// settings a controller is built with, which is what makes a rebuilt controller +// pick up a changed proxy flag. +func TestNewEmailControllerAppliesProxySetting(t *testing.T) { + t.Cleanup(func() { gomail.NetDialTimeout = net.DialTimeout }) + + NewEmailController(config.EmailConfig{NetworkUseProxy: true}) + if dialerIsDefault() { + t.Error("a controller built with networkUseProxy: true did not install the proxy dialer") + } + + NewEmailController(config.EmailConfig{NetworkUseProxy: false}) + if !dialerIsDefault() { + t.Error("a controller rebuilt with networkUseProxy: false left the proxy dialer in place") + } +} diff --git a/internal/controller/pipes/ntfy.go b/internal/controller/pipes/ntfy.go index fc574ce..2b2d6f3 100644 --- a/internal/controller/pipes/ntfy.go +++ b/internal/controller/pipes/ntfy.go @@ -4,7 +4,6 @@ import ( "fmt" "net/http" "strings" - "sync/atomic" "time" "nukumizu-backend/config" @@ -17,32 +16,21 @@ 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. +// cfg and httpClient are written once, by the constructor, and never again: a +// controller that needs different settings is replaced wholesale by the manager +// rather than reconfigured in place, so the send methods can read them without +// locking. type NtfyController struct { - cfg atomic.Pointer[config.NtfyConfig] - httpClient atomic.Pointer[http.Client] + cfg config.NtfyConfig + httpClient *http.Client } // NewNtfyController creates a new Ntfy controller. func NewNtfyController(cfg config.NtfyConfig) *NtfyController { - 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() + return &NtfyController{ + cfg: cfg, + httpClient: netproxy.HTTPClient(cfg.NetworkUseProxy, 10*time.Second), + } } // Name returns the controller name. @@ -52,7 +40,7 @@ func (n *NtfyController) Name() string { // Start initializes the Ntfy controller. func (n *NtfyController) Start() error { - if !n.settings().Enabled { + if !n.cfg.Enabled { postLog.Info("Ntfy controller is disabled") return nil } @@ -65,50 +53,26 @@ 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.settings().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.settings().Markdown + return n.cfg.Markdown } // SendStatusChange sends a status change notification via Ntfy. func (n *NtfyController) SendStatusChange(change node.StatusChange) error { - s := n.settings() - if !s.Enabled { + if !n.cfg.Enabled { return nil } cfg := config.Current() params := template.BuildParamsFromStatusChange(change) - message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, s.Markdown) + message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, n.cfg.Markdown) title := fmt.Sprintf("Server %s: %s", change.Name, change.Event) return n.publish(title, message) @@ -116,28 +80,26 @@ func (n *NtfyController) SendStatusChange(change node.StatusChange) error { // SendServerList sends the server list via Ntfy. func (n *NtfyController) SendServerList(onlineServers, offlineServers string) error { - s := n.settings() - if !s.Enabled { + if !n.cfg.Enabled { return nil } cfg := config.Current() params := template.BuildParamsFromServerList() - message := template.Render(cfg.ControllerMessage.ServerList, params, s.Markdown) + message := template.Render(cfg.ControllerMessage.ServerList, params, n.cfg.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 { - s := n.settings() - if !s.Enabled { + if !n.cfg.Enabled { return nil } cfg := config.Current() params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) - message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, s.Markdown) + message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, n.cfg.Markdown) title := fmt.Sprintf("Command Result: %s on %s", command, serverName) return n.publish(title, message) @@ -146,24 +108,21 @@ 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 { - s := n.settings() - if !s.Enabled { + if !n.cfg.Enabled { return nil } - return n.publish(alert.Subject, alert.Render(s.Markdown)) + return n.publish(alert.Subject, alert.Render(n.cfg.Markdown)) } func (n *NtfyController) publish(title, message string) error { - s := n.settings() - - serverURL := s.Server + serverURL := n.cfg.Server if serverURL == "" { serverURL = "https://ntfy.sh" } serverURL = strings.TrimRight(serverURL, "/") - publishURL := fmt.Sprintf("%s/%s", serverURL, s.Topic) + publishURL := fmt.Sprintf("%s/%s", serverURL, n.cfg.Topic) req, err := http.NewRequest("POST", publishURL, strings.NewReader(message)) if err != nil { @@ -171,14 +130,14 @@ func (n *NtfyController) publish(title, message string) error { } req.Header.Set("Title", title) - if s.Priority != "" && s.Priority != "default" { - req.Header.Set("Priority", s.Priority) + if n.cfg.Priority != "" && n.cfg.Priority != "default" { + req.Header.Set("Priority", n.cfg.Priority) } - if s.Token != "" { - req.Header.Set("Authorization", "Bearer "+s.Token) + if n.cfg.Token != "" { + req.Header.Set("Authorization", "Bearer "+n.cfg.Token) } - resp, err := n.client().Do(req) + resp, err := n.httpClient.Do(req) if err != nil { postLog.Warning("Failed to publish to ntfy: " + err.Error()) return err @@ -189,6 +148,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: " + s.Topic) + postLog.Debug("Ntfy notification sent to topic: " + n.cfg.Topic) return nil } diff --git a/internal/controller/pipes/qq_napcat/qq.go b/internal/controller/pipes/qq_napcat/qq.go index 4b05d98..21da416 100644 --- a/internal/controller/pipes/qq_napcat/qq.go +++ b/internal/controller/pipes/qq_napcat/qq.go @@ -86,23 +86,6 @@ 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 deleted file mode 100644 index e6e5cab..0000000 --- a/internal/controller/pipes/reload_test.go +++ /dev/null @@ -1,257 +0,0 @@ -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 ef58935..a082db9 100644 --- a/internal/controller/pipes/telegram/telegram.go +++ b/internal/controller/pipes/telegram/telegram.go @@ -124,21 +124,6 @@ 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 8517aaf..0b1fb8a 100644 --- a/internal/controller/pipes/webhook.go +++ b/internal/controller/pipes/webhook.go @@ -5,8 +5,6 @@ import ( "encoding/json" "fmt" "net/http" - "reflect" - "sync/atomic" "time" "nukumizu-backend/config" @@ -19,32 +17,21 @@ 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. +// cfg and httpClient are written once, by the constructor, and never again: a +// controller that needs different settings is replaced wholesale by the manager +// rather than reconfigured in place, so the send methods can read them without +// locking. type WebhookController struct { - cfg atomic.Pointer[config.WebhookConfig] - httpClient atomic.Pointer[http.Client] + cfg config.WebhookConfig + httpClient *http.Client } // NewWebhookController creates a new Webhook controller. func NewWebhookController(cfg config.WebhookConfig) *WebhookController { - 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() + return &WebhookController{ + cfg: cfg, + httpClient: netproxy.HTTPClient(cfg.NetworkUseProxy, 10*time.Second), + } } // Name returns the controller name. @@ -54,7 +41,7 @@ func (w *WebhookController) Name() string { // Start initializes the Webhook controller. func (w *WebhookController) Start() error { - if !w.settings().Enabled { + if !w.cfg.Enabled { postLog.Info("Webhook controller is disabled") return nil } @@ -67,51 +54,26 @@ 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.settings().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.settings().Markdown + return w.cfg.Markdown } // SendStatusChange sends a status change notification via Webhook. func (w *WebhookController) SendStatusChange(change node.StatusChange) error { - s := w.settings() - if !s.Enabled { + if !w.cfg.Enabled { return nil } cfg := config.Current() params := template.BuildParamsFromStatusChange(change) - message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, s.Markdown) + message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, w.cfg.Markdown) payload := map[string]interface{}{ "event": change.Event, @@ -126,14 +88,13 @@ func (w *WebhookController) SendStatusChange(change node.StatusChange) error { // SendServerList sends the server list via Webhook. func (w *WebhookController) SendServerList(onlineServers, offlineServers string) error { - s := w.settings() - if !s.Enabled { + if !w.cfg.Enabled { return nil } cfg := config.Current() params := template.BuildParamsFromServerList() - message := template.Render(cfg.ControllerMessage.ServerList, params, s.Markdown) + message := template.Render(cfg.ControllerMessage.ServerList, params, w.cfg.Markdown) payload := map[string]interface{}{ "type": "serverList", @@ -148,14 +109,13 @@ 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 { - s := w.settings() - if !s.Enabled { + if !w.cfg.Enabled { return nil } cfg := config.Current() params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) - message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, s.Markdown) + message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, w.cfg.Markdown) payload := map[string]interface{}{ "type": "executeResult", @@ -173,8 +133,7 @@ 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 { - s := w.settings() - if !s.Enabled { + if !w.cfg.Enabled { return nil } @@ -183,20 +142,15 @@ func (w *WebhookController) SendAlert(alert controller.Alert) error { "subject": alert.Subject, "source": alert.Source, "content": alert.Content, - "message": alert.Render(s.Markdown), + "message": alert.Render(w.cfg.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 { - s := w.settings() - - method := s.Method + method := w.cfg.Method if method == "" { method = "POST" } @@ -206,17 +160,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, s.URL, bytes.NewReader(bodyJSON)) + req, err := http.NewRequest(method, w.cfg.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 s.Headers { + for key, value := range w.cfg.Headers { req.Header.Set(key, value) } - resp, err := w.client().Do(req) + resp, err := w.httpClient.Do(req) if err != nil { postLog.Warning("Failed to send webhook: " + err.Error()) return err @@ -227,6 +181,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 " + s.URL) + postLog.Debug("Webhook notification sent to " + w.cfg.URL) return nil } diff --git a/main.go b/main.go index 5b1e15f..16afeb7 100644 --- a/main.go +++ b/main.go @@ -66,9 +66,17 @@ func main() { config.OnReload(func(updated *config.Config) { postLog.SetDebugMode(updated.System.DebugMode) - if mgr := controller.GetManager(); mgr != nil { - mgr.ReloadAll() + mgr := controller.GetManager() + if mgr == nil || !mgr.NeedsRebuild(updated.ControllerMethod) { + return } + + // Controller settings changed. Each controller reads its settings into + // fields when it is built, and the NapCat and Telegram ones own + // connections that cannot be re-pointed, so the change is applied by + // replacing the whole set rather than reconfiguring it in place. + postLog.Info("Controller settings changed; rebuilding every channel") + mgr.ReplaceAll(buildControllers(updated), updated.ControllerMethod) }) dbPath := cfg.DBPath @@ -212,83 +220,28 @@ func startWebhookServer(cfg *config.Config) { }() } -// initControllers initializes and starts all configured controllers. +// buildControllers constructs one controller per channel from the given +// configuration. The set is built fresh whenever controller settings change, +// because a controller reads its settings once at construction and two of them +// own connections that cannot be re-pointed. +func buildControllers(cfg *config.Config) []controller.Controller { + return []controller.Controller{ + qq_napcat.NewQQController(cfg.ControllerMethod.QQ), + telegram.NewTelegramController(cfg.ControllerMethod.Telegram), + pipes.NewEmailController(cfg.ControllerMethod.Email), + pipes.NewNtfyController(cfg.ControllerMethod.Ntfy), + pipes.NewWebhookController(cfg.ControllerMethod.Webhook), + } +} + +// initControllers installs the initial controller set. func initControllers() { cfg := config.Current() mgr := controller.GetManager() - if mgr == nil { + if mgr == nil || cfg == nil { return } - - // QQ (Napcat) controller. - qqCtrl := qq_napcat.NewQQController(cfg.ControllerMethod.QQ) - mgr.Register(qqCtrl) - go func() { - defer func() { - if r := recover(); r != nil { - postLog.Error(fmt.Sprintf("QQ controller panic: %v", r)) - } - }() - if err := qqCtrl.Start(); err != nil { - postLog.Error("Failed to start QQ controller: " + err.Error()) - } - }() - - // Telegram controller. - tgCtrl := telegram.NewTelegramController(cfg.ControllerMethod.Telegram) - mgr.Register(tgCtrl) - go func() { - defer func() { - if r := recover(); r != nil { - postLog.Error(fmt.Sprintf("Telegram controller panic: %v", r)) - } - }() - if err := tgCtrl.Start(); err != nil { - postLog.Error("Failed to start Telegram controller: " + err.Error()) - } - }() - - // Email controller (status-only). - emailCtrl := pipes.NewEmailController(cfg.ControllerMethod.Email) - mgr.Register(emailCtrl) - go func() { - defer func() { - if r := recover(); r != nil { - postLog.Error(fmt.Sprintf("Email controller panic: %v", r)) - } - }() - if err := emailCtrl.Start(); err != nil { - postLog.Error("Failed to start Email controller: " + err.Error()) - } - }() - - // Ntfy controller (status-only). - ntfyCtrl := pipes.NewNtfyController(cfg.ControllerMethod.Ntfy) - mgr.Register(ntfyCtrl) - go func() { - defer func() { - if r := recover(); r != nil { - postLog.Error(fmt.Sprintf("Ntfy controller panic: %v", r)) - } - }() - if err := ntfyCtrl.Start(); err != nil { - postLog.Error("Failed to start Ntfy controller: " + err.Error()) - } - }() - - // Webhook controller (status-only). - webhookCtrl := pipes.NewWebhookController(cfg.ControllerMethod.Webhook) - mgr.Register(webhookCtrl) - go func() { - defer func() { - if r := recover(); r != nil { - postLog.Error(fmt.Sprintf("Webhook controller panic: %v", r)) - } - }() - if err := webhookCtrl.Start(); err != nil { - postLog.Error("Failed to start Webhook controller: " + err.Error()) - } - }() + mgr.ReplaceAll(buildControllers(cfg), cfg.ControllerMethod) } // startBackgroundTasks starts periodic background goroutines.