refactor(controller): apply controller settings by rebuilding the channel set
Giving each controller a Reload method worked for the notification pipes, which only hold values, but not for the two bot channels: the NapCat client is stopped through a sync.Once and the Telegram polling context is created with the controller, so neither can be re-pointed once running. It also meant five bespoke implementations of the same idea. Replace the whole set instead. Manager.ReplaceAll swaps in a freshly built set under the registry lock, stops the outgoing controllers outside it, and starts the incoming ones off the calling goroutine so a slow handshake does not hold the settings request open. ReplaceAll is now the only way controllers are installed: Register is gone, because a lone registration would slip a controller in without recording the settings it was built from, which is what NeedsRebuild compares against. The rebuild runs from the reload hook and only fires when the controllerMethod section actually changed, so saving a message template or a webhook endpoint leaves the channels alone. NeedsRebuild compares the whole section by value, and a test pins that reloading a file reproduces it exactly — a default applied inconsistently would make every save look like a controller change and tear down every channel. With controllers immutable after construction and replaced wholesale, the atomic pointers and Reload methods have no writers left, so they are gone and the pipes return to plain fields. notifyReload now serializes hook execution for the same reason: a hook mutates process-wide state, and two overlapping settings updates would otherwise each decide to rebuild from their own view of what was applied. Kept from that work: applyEmailProxy restores gomail's default dialer when networkUseProxy is turned off. The old constructor only ever installed the proxy dialer, so switching the flag back left SMTP tunnelled through a proxy with no way to undo it short of a restart. Documents the resulting behaviour in the README under "What applies without a restart".
This commit is contained in:
@@ -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=<uuid>` (or `all`). Returns `data: {<uuid>: {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: [<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=<same types>` + 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=<same types>` + 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. |
|
||||
|
||||
+10
-2
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user