6 Commits
Author SHA1 Message Date
NanamiAdmin c7642ecca5 feat(settings): report which saved keys need a restart
Build / ubuntu-latest (push) Failing after 52s
Build / windows-latest (push) Failing after 41m49s
Every settings update answered with a bare success, so the console could not
tell a change that took effect at once from one written to the file that the
running program would keep ignoring until it was restarted — the listener
address, the storage paths, the Komari dashboard URL.

Classify the patch instead. config.RestartRequiredKeys expands a patch into the
dot-separated paths of its leaves and intersects them with the settings main
reads before it starts serving, and the settings endpoint returns that list as
data.restartRequired. It is a pure function of the patch, so UpdateSettings
keeps its signature and nothing else has to change; the write succeeds either
way, and the field only says which edits are not live.

Matching is by overlap rather than equality, so a patch that names a section —
replacing it, or deleting it with a null — reports the startup-only keys inside
it, while a sibling subtree like webhook.endpoints is not mistaken for
webhook.enabled.

The result is never nil, so an update with nothing to report serializes as []
rather than null and a client can iterate it without a guard.

The console reads the field in ConfigSection and warns instead of confirming,
naming the keys that are pending; Settings.vue's hints for the Komari URL, the
listen address and the storage paths now say which of their fields is affected
rather than labelling the whole card.
2026-09-28 23:28:47 +08:00
NanamiAdmin 457ea42417 fix(netproxy): resolve the proxy address per request
Build / ubuntu-latest (push) Failing after 3m3s
Build / windows-latest (push) Canceled after 3m19s
ProxyFunc parsed system.networkProxy once, when a client was built, and baked
the result into the transport. Editing the proxy address therefore had no effect
on any client that already existed: ntfy, the outgoing webhook, NapCat's HTTP
API, Telegram's polling and the NapCat WebSocket dialer all kept the address
they were built with, and Email's SMTP dialer did the same. Only a controller
rebuild would pick up a new one, and that fires only when controllerMethod
changes — so the address was effectively fixed until a restart.

Resolve it inside the returned function instead. A transport proxy function that
returns a nil URL asks for a direct connection, so this also covers clearing the
setting: a channel built while a proxy was configured now falls back to dialing
directly rather than retrying a dead address.

Opting out is now the only case where ProxyFunc returns nil. That is deliberate:
a function resolved to nothing at construction time is exactly what pinned the
address in the first place.

DialWithTimeout reads the address per dial for the same reason, which is what
lets the gomail NetDialTimeout hook follow a settings change.
2026-09-28 23:25:13 +08:00
NanamiAdmin 7dc6089871 docs(readme): correct how system.networkProxy takes effect
Build / ubuntu-latest (push) Canceled after 3m11s
Build / windows-latest (push) Canceled after 3m23s
The rebuild only fires when the controllerMethod section changed, so editing
networkProxy on its own does not rebuild anything and the channels keep the
proxy they connected with. The previous wording read as though a rebuild would
follow on its own.
2026-09-28 23:22:07 +08:00
NanamiAdmin 2a5a46a984 refactor(controller): apply controller settings by rebuilding the channel set
Build / ubuntu-latest (push) Canceled after 20s
Build / windows-latest (push) Canceled after 24s
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".
2026-09-28 23:21:35 +08:00
NanamiAdmin 0691e85ccf feat(controller): hot-reload the notification pipes
Build / ubuntu-latest (push) Failing after 3m13s
Build / windows-latest (push) Canceled after 6m23s
Email, ntfy and webhook kept their settings in a plain struct field copied at
construction, so editing controllerMethod in config.json wrote the file and
replaced the in-memory configuration while the running channels stayed on their
boot values.

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

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

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

QQ (Napcat) and Telegram implement Reload because the interface requires it, but
neither can apply a change in place: the NapCat client is stopped through a
sync.Once and the Telegram polling context is created with the controller, so
both need to be rebuilt and swapped into the manager. Rather than no-op silently
and let an edit look applied, they log a warning while their settings diverge.
That rebuild is the next step.
2026-09-28 23:15:11 +08:00
NanamiAdmin d9a9918e21 feat(config): notify reload hooks after a settings update
Build / ubuntu-latest (push) Failing after 3m8s
Build / windows-latest (push) Canceled after 4m3s
A settings update replaced the in-memory configuration but nothing told the code
that derives state from it. The logger is the first such consumer: its debug
flag is captured once at startup by postLog.SetDebugMode, so toggling
system.debugMode at runtime changed what config.IsDebugMode() reported while
DEBUG lines stayed filtered by the value from boot.

Add a hook registry to the config package. OnReload registers a callback and
every writer of config.json runs the registered hooks with the configuration now
in effect. It lives in config so packages that already depend on it (the logger,
and later the controller manager) can react without config importing them back,
which would be an import cycle.

Hooks run after settingsLock is released, not inside it. They will rebuild
controllers once those support reloading, which can wait on a network call, and
settingsLock is also held by the node tracker's background save — running a hook
under the lock would stall node registration behind a settings edit.
runSettingsUpdate now owns that lock/notify sequence for all four writer paths
(UpdateSettings and the three webhook endpoint helpers).

A panicking hook is logged and skipped: by then the new configuration is on disk
and published, so reporting the write as failed would be a lie, and the hooks
registered after the broken one still need to run.
2026-09-28 23:10:33 +08:00
20 changed files with 1292 additions and 138 deletions
+13 -1
View File
@@ -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 | | `{{ 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) | | `{{ 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` | Read on every request and every dial, so a new address reaches channels that are already running. Only the per-channel `networkUseProxy` opt-in is fixed when a channel is built, so toggling that still needs the rebuild above. |
| `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 ## 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. 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/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/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/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. Returns `data: {type, restartRequired}`, where `restartRequired` lists the keys the update changed that are only read at startup (see [What applies without a restart](#what-applies-without-a-restart)) — the write succeeds regardless, this only says which edits are not live yet. Always an array, empty when everything took effect. |
| `/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/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/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. | | `/api/webhook/delete` | POST | admin | Remove an endpoint. Body `{name}`. `404` for an unknown name. |
+81
View File
@@ -0,0 +1,81 @@
package config
import (
"fmt"
"sync"
"nukumizu-backend/postLog"
)
// reloadHooks are the callbacks run after the global configuration has been
// reloaded, i.e. after every settings update that touches config.json.
//
// They exist so that packages which already depend on config — the controller
// manager, the logger — can react to an update without config importing them,
// which would be an import cycle. Everything a hook needs is passed in.
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
// configuration, receiving the configuration now in effect. Hooks run in
// registration order.
//
// Register once, at startup, before the first settings update can arrive: a
// hook registered later has already missed the updates that came before it, and
// the configuration it would have seen is not replayed.
//
// A hook runs on the goroutine serving /api/settings/set, so it must not block
// for long. It may be called concurrently by two overlapping updates.
func OnReload(hook func(*Config)) {
reloadHooksMu.Lock()
defer reloadHooksMu.Unlock()
reloadHooks = append(reloadHooks, hook)
}
// notifyReload runs every registered hook with cfg. A nil cfg is the signal
// that the update touched one of the other settings files, which have no
// hook-visible reload, and it is ignored.
//
// A panicking hook is logged and skipped rather than allowed to unwind through
// UpdateSettings: by the time hooks run the new configuration is already on
// disk and published in memory, so reporting the write as failed would be a
// lie, and the hooks registered after the broken one must still run.
func notifyReload(cfg *Config) {
if cfg == nil {
return
}
// 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)
}
}
// runReloadHook runs one hook, isolating a panic to that hook. Recovering in a
// separate function rather than inline is deliberate: a deferred recover placed
// in the loop body would not run until notifyReload itself returned, which
// would abandon the remaining hooks.
func runReloadHook(hook func(*Config), cfg *Config) {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Configuration reload hook panicked: %v", r))
}
}()
hook(cfg)
}
+173
View File
@@ -0,0 +1,173 @@
package config
import (
"reflect"
"testing"
"nukumizu-backend/global"
)
// swapReloadHooks replaces the registered reload hooks for the duration of a
// test and restores the previous set afterwards, so the package-level registry
// cannot leak into another test.
func swapReloadHooks(t *testing.T, hooks ...func(*Config)) {
t.Helper()
reloadHooksMu.Lock()
original := reloadHooks
reloadHooks = hooks
reloadHooksMu.Unlock()
t.Cleanup(func() {
reloadHooksMu.Lock()
reloadHooks = original
reloadHooksMu.Unlock()
})
}
func TestReloadHookRunsOnlyForGlobalSettings(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{"system":{"debugMode":false}}`)
writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{}`)
writeTempConfig(t, &global.ConfigPath.BotNodeConfig, `{}`)
var seen []bool
swapReloadHooks(t, func(cfg *Config) {
seen = append(seen, cfg.System.DebugMode)
})
// A config.json update runs the hook, with the values just written.
if err := UpdateSettings(SettingGlobal, map[string]interface{}{
"system": map[string]interface{}{"debugMode": true},
}); err != nil {
t.Fatalf("UpdateSettings(global): %v", err)
}
if len(seen) != 1 {
t.Fatalf("hook ran %d times for a config.json update, want 1", len(seen))
}
if !seen[0] {
t.Error("hook received a configuration without the updated debugMode")
}
// The other two files reload a singleton that callers read at the point of
// use, so they have nothing to notify.
for _, settingsType := range []string{SettingBotUserConfig, SettingBotNodeConfig} {
if err := UpdateSettings(settingsType, map[string]interface{}{
"unused": map[string]interface{}{"event_reply": true},
}); err != nil {
t.Fatalf("UpdateSettings(%s): %v", settingsType, err)
}
}
if len(seen) != 1 {
t.Errorf("hook ran %d times after updates to the other settings files, want 1", len(seen))
}
}
func TestReloadHookIsolatesPanic(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{"system":{"debugMode":false}}`)
reached := false
swapReloadHooks(t,
func(*Config) { panic("hook under test") },
func(*Config) { reached = true },
)
// The configuration is already on disk and published by the time hooks run,
// so a broken hook must not turn a successful write into a failed request.
if err := UpdateSettings(SettingGlobal, map[string]interface{}{
"system": map[string]interface{}{"debugMode": true},
}); err != nil {
t.Fatalf("a panicking hook must not fail the settings write: %v", err)
}
if !reached {
t.Error("a panicking hook stopped the hooks registered after it")
}
// The write itself must still have landed.
cfg := Current()
if cfg == nil || !cfg.System.DebugMode {
t.Error("the settings write did not take effect")
}
}
// 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.
func TestReloadHookSeesEveryWriterPath(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{
"webhook": { "endpoints": {} }
}`)
var seen int
swapReloadHooks(t, func(*Config) { seen++ })
if err := AddWebhookEndpoint("example", map[string]interface{}{
"enabled": true,
"token": "secret",
"notifyPipes": []interface{}{"ntfy"},
}); err != nil {
t.Fatalf("AddWebhookEndpoint: %v", err)
}
if seen != 1 {
t.Errorf("hook ran %d times after adding an endpoint, want 1", seen)
}
if err := ModifyWebhookEndpoint("example", map[string]interface{}{
"enabled": false,
}); err != nil {
t.Fatalf("ModifyWebhookEndpoint: %v", err)
}
if seen != 2 {
t.Errorf("hook ran %d times after modifying an endpoint, want 2", seen)
}
if err := DeleteWebhookEndpoint("example"); err != nil {
t.Fatalf("DeleteWebhookEndpoint: %v", err)
}
if seen != 3 {
t.Errorf("hook ran %d times after deleting an endpoint, want 3", seen)
}
// A rejected write changes nothing, so it must not notify either.
if err := DeleteWebhookEndpoint("ghost"); err == nil {
t.Error("deleting an unknown endpoint should fail")
}
if seen != 3 {
t.Errorf("hook ran for a rejected write (%d notifications, want 3)", seen)
}
}
+127 -16
View File
@@ -6,6 +6,7 @@ import (
"errors" "errors"
"fmt" "fmt"
"os" "os"
"strings"
"sync" "sync"
"nukumizu-backend/global" "nukumizu-backend/global"
@@ -23,6 +24,89 @@ const (
// accepted constants above. // accepted constants above.
var ErrUnsupportedSettingsType = errors.New("unsupported settings type") var ErrUnsupportedSettingsType = errors.New("unsupported settings type")
// startupOnlySettings are the config.json keys that are read once before the
// program starts serving and never again: the listener addresses, the paths the
// databases are opened from, and the Komari dashboard its client is built
// against. Editing one writes the file and replaces the in-memory
// configuration, but the running program keeps the old value, so an update that
// touches one is reported back to the caller instead of being silently
// accepted.
//
// Keep this in step with main: these are exactly the settings main reads before
// the HTTP server comes up. Everything else — controllerMethod, networkProxy,
// the message templates, the debug switches — is picked up at runtime.
var startupOnlySettings = []string{
"system.listenAddr",
"system.listenPort",
"webhook.enabled",
"webhook.listenAddr",
"webhook.listenPort",
"komari.dashboardURL",
"dataPath",
"dbPath",
}
// RestartRequiredKeys lists the settings in patch that only take effect at
// startup, as dot-separated paths, in the order startupOnlySettings declares
// them. Only config.json carries such settings; an update to one of the other
// files always reports nothing.
//
// The write itself succeeds either way — this is advice for the user, not a
// rejection. The result is never nil, so a caller can put it straight into a
// JSON response and get [] rather than null.
func RestartRequiredKeys(settingsType string, patch map[string]interface{}) []string {
keys := []string{}
if settingsType != SettingGlobal {
return keys
}
patched := patchPaths(patch)
for _, watched := range startupOnlySettings {
for _, path := range patched {
if pathsOverlap(path, watched) {
keys = append(keys, watched)
break
}
}
}
return keys
}
// patchPaths expands a nested settings patch into the dot-separated paths of its
// leaves. An object is descended into rather than reported, so a patch that only
// names sections still resolves to the keys it changes, and a JSON null is a
// leaf because it deletes the key it names.
func patchPaths(patch map[string]interface{}) []string {
paths := []string{}
var walk func(prefix string, node map[string]interface{})
walk = func(prefix string, node map[string]interface{}) {
for key, value := range node {
path := key
if prefix != "" {
path = prefix + "." + key
}
if nested, ok := value.(map[string]interface{}); ok && nested != nil {
walk(path, nested)
continue
}
paths = append(paths, path)
}
}
walk("", patch)
return paths
}
// pathsOverlap reports whether a patched path and a watched setting can affect
// each other: they are the same key, the patch names something inside the
// watched setting, or the patch names a section the watched setting lives in.
// The last case matters because a patch may replace a whole section, which
// changes every key under it.
func pathsOverlap(patched, watched string) bool {
return patched == watched ||
strings.HasPrefix(patched, watched+".") ||
strings.HasPrefix(watched, patched+".")
}
// settingsLock serializes read-modify-write access to the on-disk configuration // settingsLock serializes read-modify-write access to the on-disk configuration
// files so concurrent admin edits (UpdateSettings) and the node tracker's // files so concurrent admin edits (UpdateSettings) and the node tracker's
// background save (SaveBotNodeConfig) cannot lose each other's updates. // background save (SaveBotNodeConfig) cannot lose each other's updates.
@@ -82,19 +166,45 @@ func GetSettings(settingsType string) ([]byte, error) {
// written the matching in-memory singleton is reloaded so runtime code observes // written the matching in-memory singleton is reloaded so runtime code observes
// the new values. // the new values.
func UpdateSettings(settingsType string, patch map[string]interface{}) error { func UpdateSettings(settingsType string, patch map[string]interface{}) error {
settingsLock.Lock() return runSettingsUpdate(func() (*Config, error) {
defer settingsLock.Unlock() return updateSettingsLocked(settingsType, patch)
})
}
return updateSettingsLocked(settingsType, patch) // runSettingsUpdate runs fn under settingsLock and then, once the lock is
// released, runs the reload hooks with whatever configuration fn reports (nil
// when the update did not touch config.json).
//
// The hooks deliberately run outside settingsLock. A hook rebuilds controllers,
// which can wait on a network call, while settingsLock is also held by the node
// tracker's background save (SaveBotNodeConfig); holding it across a hook would
// stall node registration behind an unrelated settings edit.
func runSettingsUpdate(fn func() (*Config, error)) error {
cfg, err := func() (*Config, error) {
settingsLock.Lock()
defer settingsLock.Unlock()
return fn()
}()
if err != nil {
return err
}
notifyReload(cfg)
return nil
} }
// updateSettingsLocked is UpdateSettings without the locking, for callers that // updateSettingsLocked is UpdateSettings without the locking, for callers that
// need to inspect the loaded configuration and write in one critical section // need to inspect the loaded configuration and write in one critical section
// (see the incoming webhook endpoint helpers). Callers must hold settingsLock. // (see the incoming webhook endpoint helpers). Callers must hold settingsLock.
func updateSettingsLocked(settingsType string, patch map[string]interface{}) error { //
// It returns the freshly loaded global configuration, or nil when the settings
// type is one of the other files. The caller is responsible for handing that
// value to notifyReload once settingsLock is released — which runSettingsUpdate
// does for every writer.
func updateSettingsLocked(settingsType string, patch map[string]interface{}) (*Config, error) {
path, err := settingsPath(settingsType) path, err := settingsPath(settingsType)
if err != nil { if err != nil {
return err return nil, err
} }
// Start from whatever is already on disk so nothing is dropped. A missing or // Start from whatever is already on disk so nothing is dropped. A missing or
@@ -104,23 +214,23 @@ func updateSettingsLocked(settingsType string, patch map[string]interface{}) err
if err == nil { if err == nil {
if len(bytes.TrimSpace(data)) > 0 { if len(bytes.TrimSpace(data)) > 0 {
if err := json.Unmarshal(data, &current); err != nil { if err := json.Unmarshal(data, &current); err != nil {
return fmt.Errorf("failed to parse existing %s settings file %s: %w", settingsType, path, err) return nil, fmt.Errorf("failed to parse existing %s settings file %s: %w", settingsType, path, err)
} }
} }
} else if !os.IsNotExist(err) { } else if !os.IsNotExist(err) {
return fmt.Errorf("failed to read existing %s settings file %s: %w", settingsType, path, err) return nil, fmt.Errorf("failed to read existing %s settings file %s: %w", settingsType, path, err)
} }
deepMergeSettings(current, patch) deepMergeSettings(current, patch)
data, err = json.MarshalIndent(current, "", " ") data, err = json.MarshalIndent(current, "", " ")
if err != nil { if err != nil {
return fmt.Errorf("failed to marshal %s settings: %w", settingsType, err) return nil, fmt.Errorf("failed to marshal %s settings: %w", settingsType, err)
} }
data = append(data, '\n') data = append(data, '\n')
if err := os.WriteFile(path, data, 0o644); err != nil { if err := os.WriteFile(path, data, 0o644); err != nil {
return fmt.Errorf("failed to write %s settings file %s: %w", settingsType, path, err) return nil, fmt.Errorf("failed to write %s settings file %s: %w", settingsType, path, err)
} }
return reloadSettings(settingsType, path) return reloadSettings(settingsType, path)
@@ -151,18 +261,19 @@ func deepMergeSettings(dst, src map[string]interface{}) {
} }
// reloadSettings refreshes the in-memory singleton for the given settings type // reloadSettings refreshes the in-memory singleton for the given settings type
// so the running program observes the values just persisted to disk. // so the running program observes the values just persisted to disk. Only
func reloadSettings(settingsType, path string) error { // config.json has a hook-visible reload, so for SettingGlobal it returns the
// configuration now in effect and for the other types it returns nil.
func reloadSettings(settingsType, path string) (*Config, error) {
switch settingsType { switch settingsType {
case SettingGlobal: case SettingGlobal:
_, err := LoadGlobalConfig(path) return LoadGlobalConfig(path)
return err
case SettingBotUserConfig: case SettingBotUserConfig:
_, err := LoadBotUserConfig(path) _, err := LoadBotUserConfig(path)
return err return nil, err
case SettingBotNodeConfig: case SettingBotNodeConfig:
return LoadBotNodeConfig(path) return nil, LoadBotNodeConfig(path)
default: default:
return ErrUnsupportedSettingsType return nil, ErrUnsupportedSettingsType
} }
} }
+157
View File
@@ -4,6 +4,7 @@ import (
"encoding/json" "encoding/json"
"os" "os"
"path/filepath" "path/filepath"
"reflect"
"strings" "strings"
"testing" "testing"
@@ -148,6 +149,162 @@ func TestUpdateSettingsReplacesArraysAndKeepsNumbers(t *testing.T) {
} }
} }
func TestRestartRequiredKeys(t *testing.T) {
cases := []struct {
name string
settingsType string
patch map[string]interface{}
want []string
}{
{
name: "a runtime switch needs no restart",
settingsType: SettingGlobal,
patch: map[string]interface{}{"system": map[string]interface{}{"debugMode": true}},
want: []string{},
},
{
name: "the listen port does",
settingsType: SettingGlobal,
patch: map[string]interface{}{"system": map[string]interface{}{"listenPort": "9090"}},
want: []string{"system.listenPort"},
},
{
name: "only the startup key of a mixed patch is reported",
settingsType: SettingGlobal,
patch: map[string]interface{}{
"system": map[string]interface{}{"listenPort": "9090", "debugMode": true},
},
want: []string{"system.listenPort"},
},
{
name: "a top-level path is reported",
settingsType: SettingGlobal,
patch: map[string]interface{}{"dataPath": "/srv/data"},
want: []string{"dataPath"},
},
{
name: "results follow the declared order, not the patch order",
settingsType: SettingGlobal,
patch: map[string]interface{}{"dbPath": "/srv/db", "dataPath": "/srv/data"},
want: []string{"dataPath", "dbPath"},
},
{
name: "deleting a startup key with null is reported",
settingsType: SettingGlobal,
patch: map[string]interface{}{"system": map[string]interface{}{"listenPort": nil}},
want: []string{"system.listenPort"},
},
{
name: "replacing a whole section reports the startup keys inside it",
settingsType: SettingGlobal,
patch: map[string]interface{}{"webhook": map[string]interface{}{"listenAddr": "127.0.0.1"}},
want: []string{"webhook.listenAddr"},
},
{
// Deleting the section resets the URL to its built-in default.
name: "deleting a section the startup key lives in reports it",
settingsType: SettingGlobal,
patch: map[string]interface{}{"komari": nil},
want: []string{"komari.dashboardURL"},
},
{
name: "deleting a section reports every startup key inside it",
settingsType: SettingGlobal,
patch: map[string]interface{}{"webhook": nil},
want: []string{"webhook.enabled", "webhook.listenAddr", "webhook.listenPort"},
},
{
// An empty object merges nothing, so it changes no key and needs no
// restart — surprising enough to pin.
name: "an empty object changes nothing",
settingsType: SettingGlobal,
patch: map[string]interface{}{"komari": map[string]interface{}{}},
want: []string{},
},
{
// webhook.endpoints must not be mistaken for webhook.enabled.
name: "a sibling subtree is not mistaken for the startup key",
settingsType: SettingGlobal,
patch: map[string]interface{}{
"webhook": map[string]interface{}{
"endpoints": map[string]interface{}{"example": map[string]interface{}{"enabled": true}},
},
},
want: []string{},
},
{
// The Komari credentials are re-read on the next login, so only the
// dashboard URL is startup-only.
name: "komari credentials are not startup-only",
settingsType: SettingGlobal,
patch: map[string]interface{}{
"komari": map[string]interface{}{
"account": map[string]interface{}{"username": "admin", "password": "x"},
},
},
want: []string{},
},
{
name: "controller settings are not startup-only",
settingsType: SettingGlobal,
patch: map[string]interface{}{
"controllerMethod": map[string]interface{}{
"telegram": map[string]interface{}{"enabled": true, "botToken": "t"},
},
},
want: []string{},
},
{
name: "an empty patch reports nothing",
settingsType: SettingGlobal,
patch: map[string]interface{}{},
want: []string{},
},
{
// Only config.json has settings that are read once at startup.
name: "the other settings files never need a restart",
settingsType: SettingBotUserConfig,
patch: map[string]interface{}{"dataPath": "/srv/data"},
want: []string{},
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
got := RestartRequiredKeys(tc.settingsType, tc.patch)
if !reflect.DeepEqual(got, tc.want) {
t.Errorf("RestartRequiredKeys() = %v, want %v", got, tc.want)
}
if got == nil {
t.Error("the result must never be nil, so it serializes as [] rather than null")
}
})
}
}
// TestRestartRequiredKeysIsAdvisory pins that reporting a startup-only key does
// not stop the write: the caller is told, the file is still updated.
func TestRestartRequiredKeysIsAdvisory(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{"system":{"listenPort":"8080"}}`)
keys := RestartRequiredKeys(SettingGlobal, map[string]interface{}{
"system": map[string]interface{}{"listenPort": "9090"},
})
if len(keys) != 1 {
t.Fatalf("expected the listen port to be reported, got %v", keys)
}
if err := UpdateSettings(SettingGlobal, map[string]interface{}{
"system": map[string]interface{}{"listenPort": "9090"},
}); err != nil {
t.Fatalf("UpdateSettings: %v", err)
}
if got := Current().System.ListenPort; got != "9090" {
t.Errorf("the update was not applied: listenPort = %q", got)
}
}
func TestSettingsTypeValidation(t *testing.T) { func TestSettingsTypeValidation(t *testing.T) {
for _, valid := range []string{SettingGlobal, SettingBotUserConfig, SettingBotNodeConfig} { for _, valid := range []string{SettingGlobal, SettingBotUserConfig, SettingBotNodeConfig} {
if !IsValidSettingsType(valid) { if !IsValidSettingsType(valid) {
+18 -22
View File
@@ -60,13 +60,12 @@ func AddWebhookEndpoint(name string, fields map[string]interface{}) error {
return err return err
} }
settingsLock.Lock() return runSettingsUpdate(func() (*Config, error) {
defer settingsLock.Unlock() if _, exists := webhookEndpoint(name); exists {
return nil, fmt.Errorf("%w: %s", ErrWebhookEndpointExists, name)
if _, exists := webhookEndpoint(name); exists { }
return fmt.Errorf("%w: %s", ErrWebhookEndpointExists, name) return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch))
} })
return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch))
} }
// ModifyWebhookEndpoint updates an existing incoming webhook endpoint. Only the // ModifyWebhookEndpoint updates an existing incoming webhook endpoint. Only the
@@ -84,27 +83,24 @@ func ModifyWebhookEndpoint(name string, fields map[string]interface{}) error {
return fmt.Errorf("%w: no fields to update", ErrWebhookEndpointInvalid) return fmt.Errorf("%w: no fields to update", ErrWebhookEndpointInvalid)
} }
settingsLock.Lock() return runSettingsUpdate(func() (*Config, error) {
defer settingsLock.Unlock() if _, exists := webhookEndpoint(name); !exists {
return nil, fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name)
if _, exists := webhookEndpoint(name); !exists { }
return fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name) return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch))
} })
return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch))
} }
// DeleteWebhookEndpoint removes the incoming webhook endpoint registered under // DeleteWebhookEndpoint removes the incoming webhook endpoint registered under
// name. The endpoint stops accepting requests as soon as the configuration is // name. The endpoint stops accepting requests as soon as the configuration is
// reloaded. // reloaded.
func DeleteWebhookEndpoint(name string) error { func DeleteWebhookEndpoint(name string) error {
settingsLock.Lock() return runSettingsUpdate(func() (*Config, error) {
defer settingsLock.Unlock() if _, exists := webhookEndpoint(name); !exists {
return nil, fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name)
if _, exists := webhookEndpoint(name); !exists { }
return fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name) return updateSettingsLocked(SettingGlobal, webhookEndpointDeletePatch(name))
} })
return updateSettingsLocked(SettingGlobal, webhookEndpointDeletePatch(name))
} }
// webhookEndpoint returns the named endpoint held by the loaded configuration. // webhookEndpoint returns the named endpoint held by the loaded configuration.
+4
View File
@@ -13,6 +13,10 @@ export const serverApi = {
// /api/settings/get?type=… / /api/settings/set?type=… // /api/settings/get?type=… / /api/settings/set?type=…
// get → { success, message, data: { config } }. // get → { success, message, data: { config } }.
// set → { success, message, data: { type, restartRequired } }, where
// restartRequired lists the keys the update changed that are only read at
// startup, so the caller can say which edits are not live yet. It is
// always an array, empty when the whole update took effect.
// `type` is one of global | bot_user_config | bot_node_config. // `type` is one of global | bot_user_config | bot_node_config.
// For set, pass a partial object; a JSON null value removes that key. // For set, pass a partial object; a JSON null value removes that key.
export const settingsApi = { export const settingsApi = {
+9 -2
View File
@@ -133,8 +133,15 @@ async function save() {
const patch = wrapRoot(props.section, nest(obj)); const patch = wrapRoot(props.section, nest(obj));
saving.value = true; saving.value = true;
try { try {
await settingsApi.set('global', patch); const res = await settingsApi.set('global', patch);
toast.success(`${props.section.title} saved`); // The backend reports the keys it wrote that are only read at startup.
// A plain "saved" would suggest those are live too.
const pending = (res && res.data && res.data.restartRequired) || [];
if (pending.length) {
toast.warn(`${props.section.title} saved — restart to apply: ${pending.join(', ')}`, 7000);
} else {
toast.success(`${props.section.title} saved`);
}
emit('saved'); emit('saved');
} catch (e) { } catch (e) {
toast.error('Failed to save: ' + e.message); toast.error('Failed to save: ' + e.message);
+2 -2
View File
@@ -11,7 +11,7 @@ const sections = [
{ {
id: 'system', id: 'system',
title: 'System', title: 'System',
hint: 'HTTP listener and global runtime switches.', hint: 'HTTP listener and global runtime switches. A changed listen address or port applies on restart.',
root: ['system'], root: ['system'],
fields: [ fields: [
{ key: 'debugMode', type: 'bool', label: 'Debug mode', help: 'Skipped X-Timestamp checks and verbose debug logging.' }, { key: 'debugMode', type: 'bool', label: 'Debug mode', help: 'Skipped X-Timestamp checks and verbose debug logging.' },
@@ -37,7 +37,7 @@ const sections = [
{ {
id: 'komari', id: 'komari',
title: 'Komari dashboard', title: 'Komari dashboard',
hint: 'Connection the monitor reads node data from. Takes effect on restart.', hint: 'Connection the monitor reads node data from. The URL applies on restart; the account is re-read on the next login.',
root: ['komari'], root: ['komari'],
fields: [ fields: [
{ key: 'dashboardURL', type: 'text', label: 'Dashboard URL' }, { key: 'dashboardURL', type: 'text', label: 'Dashboard URL' },
+8 -1
View File
@@ -47,6 +47,10 @@ func SettingsGetHandler(w http.ResponseWriter, r *http.Request) {
// {"system": {"debugMode": true}} // {"system": {"debugMode": true}}
// //
// Multiple entries may be given at once; only the provided keys are changed. // Multiple entries may be given at once; only the provided keys are changed.
//
// The response carries data.restartRequired: the keys the update changed that
// are only read at startup, so the caller can say which edits are not live yet.
// It is empty for an update that took effect in full.
func SettingsSetHandler(w http.ResponseWriter, r *http.Request) { func SettingsSetHandler(w http.ResponseWriter, r *http.Request) {
if !utils.Auth(w, r, "POST", "admin") { if !utils.Auth(w, r, "POST", "admin") {
return return
@@ -80,7 +84,10 @@ func SettingsSetHandler(w http.ResponseWriter, r *http.Request) {
return return
} }
// The write landed either way. This only tells the caller which of the keys
// it changed will not be live until the program is restarted.
utils.SendSuccessResponse(w, "settings updated successfully", map[string]interface{}{ utils.SendSuccessResponse(w, "settings updated successfully", map[string]interface{}{
"type": settingsType, "type": settingsType,
"restartRequired": config.RestartRequiredKeys(settingsType, patch),
}) })
} }
+128
View File
@@ -0,0 +1,128 @@
package handler
import (
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strconv"
"strings"
"testing"
"time"
"nukumizu-backend/config"
"nukumizu-backend/global"
)
// tempConfigFile points one of the settings paths at a throwaway file, so a
// test can drive the settings API without touching the config files in the
// working directory.
func tempConfigFile(t *testing.T, field *string, content string) {
t.Helper()
path := filepath.Join(t.TempDir(), "config.json")
if err := os.WriteFile(path, []byte(content), 0o644); err != nil {
t.Fatalf("write temp config: %v", err)
}
original := *field
*field = path
t.Cleanup(func() { *field = original })
}
// settingsSetRequest builds an authenticated POST for the settings endpoint.
func settingsSetRequest(t *testing.T, settingsType, body string) *http.Request {
t.Helper()
req := httptest.NewRequest(
http.MethodPost,
"/api/settings/set?type="+settingsType,
strings.NewReader(body),
)
req.Header.Set("X-Token", "test-admin-token")
req.Header.Set("X-Timestamp", strconv.FormatInt(time.Now().Unix(), 10))
return req
}
// settingsSetResponse is the envelope /api/settings/set answers with.
type settingsSetResponse struct {
Success bool `json:"success"`
Data struct {
Type string `json:"type"`
RestartRequired []string `json:"restartRequired"`
} `json:"data"`
}
func decodeSettingsSetResponse(t *testing.T, w *httptest.ResponseRecorder) settingsSetResponse {
t.Helper()
var body settingsSetResponse
if err := json.Unmarshal(w.Body.Bytes(), &body); err != nil {
t.Fatalf("decode response: %v; body=%s", err, w.Body.String())
}
return body
}
// TestSettingsSetReportsStartupOnlyKeys covers the field the console reads to
// tell the user which of their edits are not live yet.
func TestSettingsSetReportsStartupOnlyKeys(t *testing.T) {
setupAdminToken()
tempConfigFile(t, &global.ConfigPath.Global, `{"system":{"debugMode":false,"listenPort":"8080"}}`)
w := httptest.NewRecorder()
SettingsSetHandler(w, settingsSetRequest(t, "global",
`{"system":{"listenPort":"9090","debugMode":true}}`))
if w.Code != http.StatusOK {
t.Fatalf("status = %d, body = %s", w.Code, w.Body.String())
}
body := decodeSettingsSetResponse(t, w)
if !body.Success {
t.Fatalf("not a success envelope: %s", w.Body.String())
}
if body.Data.Type != "global" {
t.Errorf("data.type = %q, want global", body.Data.Type)
}
if len(body.Data.RestartRequired) != 1 || body.Data.RestartRequired[0] != "system.listenPort" {
t.Errorf("restartRequired = %v, want [system.listenPort]", body.Data.RestartRequired)
}
// Being reported as startup-only must not stop the write.
data, err := config.GetSettings(config.SettingGlobal)
if err != nil {
t.Fatalf("GetSettings: %v", err)
}
for _, want := range []string{`"listenPort": "9090"`, `"debugMode": true`} {
if !strings.Contains(string(data), want) {
t.Errorf("the patch was not applied, missing %s:\n%s", want, data)
}
}
}
// TestSettingsSetRestartRequiredIsAlwaysAnArray pins the shape a client
// iterates over: an update with nothing to report must answer [] and not null.
func TestSettingsSetRestartRequiredIsAlwaysAnArray(t *testing.T) {
setupAdminToken()
tempConfigFile(t, &global.ConfigPath.Global, `{"system":{"debugMode":false}}`)
w := httptest.NewRecorder()
SettingsSetHandler(w, settingsSetRequest(t, "global", `{"system":{"debugMode":true}}`))
if w.Code != http.StatusOK {
t.Fatalf("status = %d, body = %s", w.Code, w.Body.String())
}
if !strings.Contains(w.Body.String(), `"restartRequired":[]`) {
t.Errorf("restartRequired should serialize as an empty array: %s", w.Body.String())
}
}
// TestSettingsSetRejectsUnknownType keeps the failure path intact now that the
// success path computes an extra field.
func TestSettingsSetRejectsUnknownType(t *testing.T) {
setupAdminToken()
tempConfigFile(t, &global.ConfigPath.Global, `{}`)
w := httptest.NewRecorder()
SettingsSetHandler(w, settingsSetRequest(t, "nonsense", `{}`))
if w.Code != http.StatusBadRequest {
t.Errorf("status = %d, want 400", w.Code)
}
}
+71 -5
View File
@@ -3,8 +3,10 @@ package controller
import ( import (
"errors" "errors"
"fmt" "fmt"
"reflect"
"strings" "strings"
"sync" "sync"
"sync/atomic"
"nukumizu-backend/config" "nukumizu-backend/config"
"nukumizu-backend/internal/node" "nukumizu-backend/internal/node"
@@ -109,6 +111,12 @@ type BotController interface {
type Manager struct { type Manager struct {
mu sync.RWMutex mu sync.RWMutex
controllers map[string]Controller 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 var globalManager *Manager
@@ -126,12 +134,70 @@ func GetManager() *Manager {
return globalManager return globalManager
} }
// Register adds a controller to the manager. // NeedsRebuild reports whether next differs from the controllerMethod section
func (m *Manager) Register(c Controller) { // 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)
}
// 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 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() m.mu.Lock()
defer m.mu.Unlock() previous := m.controllers
m.controllers[c.Name()] = c m.controllers = make(map[string]Controller, len(next))
postLog.Info("Controller registered: " + c.Name()) for _, ctrl := range next {
m.controllers[ctrl.Name()] = ctrl
}
m.mu.Unlock()
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)
}
} }
// ShowBotInitMessage sends the bot initialization message to all enabled // ShowBotInitMessage sends the bot initialization message to all enabled
+222
View File
@@ -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")
}
}
+20 -7
View File
@@ -2,6 +2,7 @@ package pipes
import ( import (
"fmt" "fmt"
"net"
gomail "gopkg.in/mail.v2" gomail "gopkg.in/mail.v2"
@@ -14,22 +15,34 @@ import (
) )
// EmailController handles email notifications via SMTP. // EmailController handles email notifications via SMTP.
//
// 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 { type EmailController struct {
cfg config.EmailConfig cfg config.EmailConfig
} }
// NewEmailController creates a new Email controller. // NewEmailController creates a new Email controller.
func NewEmailController(cfg config.EmailConfig) *EmailController { func NewEmailController(cfg config.EmailConfig) *EmailController {
if cfg.NetworkUseProxy { applyEmailProxy(cfg.NetworkUseProxy)
// Route SMTP through the HTTP CONNECT proxy. NetDialTimeout is
// gomail's documented hook for overriding how the SMTP connection is
// dialed. There is a single global email channel, so overriding it
// unconditionally when the flag is set is safe.
gomail.NetDialTimeout = netproxy.DialWithTimeout(true)
}
return &EmailController{cfg: cfg} return &EmailController{cfg: cfg}
} }
// applyEmailProxy routes SMTP through the HTTP CONNECT proxy, or restores a
// direct dial. gomail exposes the dial path as the package-level
// NetDialTimeout, whose own default is net.DialTimeout, so turning the proxy
// off has to put that back rather than leave the hook in place. There is a
// single global email channel, so setting a package-level hook here is
// unambiguous.
func applyEmailProxy(useProxy bool) {
if useProxy {
gomail.NetDialTimeout = netproxy.DialWithTimeout(true)
return
}
gomail.NetDialTimeout = net.DialTimeout
}
// Name returns the controller name. // Name returns the controller name.
func (e *EmailController) Name() string { func (e *EmailController) Name() string {
return "email" return "email"
@@ -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")
}
}
+5
View File
@@ -15,6 +15,11 @@ import (
) )
// NtfyController handles notifications via ntfy.sh or a self-hosted ntfy server. // NtfyController handles notifications via ntfy.sh or a self-hosted ntfy server.
//
// 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 { type NtfyController struct {
cfg config.NtfyConfig cfg config.NtfyConfig
httpClient *http.Client httpClient *http.Client
+5
View File
@@ -16,6 +16,11 @@ import (
) )
// WebhookController handles notifications via generic HTTP webhooks. // WebhookController handles notifications via generic HTTP webhooks.
//
// 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 { type WebhookController struct {
cfg config.WebhookConfig cfg config.WebhookConfig
httpClient *http.Client httpClient *http.Client
+24 -10
View File
@@ -2,6 +2,10 @@
// network proxy configured in the system config. Each caller decides whether // network proxy configured in the system config. Each caller decides whether
// to use the proxy by passing its own useProxy flag (the per-channel // to use the proxy by passing its own useProxy flag (the per-channel
// networkUseProxy setting), so proxying is opt-in per channel. // networkUseProxy setting), so proxying is opt-in per channel.
//
// The opt-in is captured when a client is built, but the proxy address is not:
// it is read again on every request and every dial, so editing
// system.networkProxy takes effect on clients that already exist.
package netproxy package netproxy
import ( import (
@@ -39,19 +43,23 @@ func proxyURL() *url.URL {
} }
// ProxyFunc returns a transport proxy function that routes requests through // ProxyFunc returns a transport proxy function that routes requests through
// the configured network proxy when enabled. It returns nil when the caller // the configured network proxy when enabled, and nil when the caller opts out
// opts out or no proxy is configured, meaning direct connection. The returned // of proxying entirely. The returned function is compatible with both
// function is compatible with both http.Transport.Proxy and // http.Transport.Proxy and websocket.Dialer.Proxy.
// websocket.Dialer.Proxy. //
// The proxy address is resolved on every call rather than once here, so a
// settings update that changes system.networkProxy reaches a client that was
// already built. That is also why opting out is the only case that returns nil:
// a function resolved to nothing at construction time would pin its client to
// whatever was configured then. A nil URL from the returned function means no
// proxy is configured and the request goes direct.
func ProxyFunc(useProxy bool) func(*http.Request) (*url.URL, error) { func ProxyFunc(useProxy bool) func(*http.Request) (*url.URL, error) {
if !useProxy { if !useProxy {
return nil return nil
} }
u := proxyURL() return func(*http.Request) (*url.URL, error) {
if u == nil { return proxyURL(), nil
return nil
} }
return http.ProxyURL(u)
} }
// HTTPClient builds an http.Client that sends traffic through the configured // HTTPClient builds an http.Client that sends traffic through the configured
@@ -71,10 +79,16 @@ func HTTPClient(useProxy bool, timeout time.Duration) *http.Client {
// through the configured HTTP CONNECT proxy when enabled. Its signature // through the configured HTTP CONNECT proxy when enabled. Its signature
// matches net.DialTimeout so it can be plugged into gomail's NetDialTimeout // matches net.DialTimeout so it can be plugged into gomail's NetDialTimeout
// to send SMTP over the proxy. // to send SMTP over the proxy.
//
// Like ProxyFunc it reads the proxy address per dial, so clearing or changing
// system.networkProxy reaches a dialer that already exists.
func DialWithTimeout(useProxy bool) func(network, addr string, timeout time.Duration) (net.Conn, error) { func DialWithTimeout(useProxy bool) func(network, addr string, timeout time.Duration) (net.Conn, error) {
u := proxyURL()
return func(network, addr string, timeout time.Duration) (net.Conn, error) { return func(network, addr string, timeout time.Duration) (net.Conn, error) {
if !useProxy || u == nil { if !useProxy {
return net.DialTimeout(network, addr, timeout)
}
u := proxyURL()
if u == nil {
return net.DialTimeout(network, addr, timeout) return net.DialTimeout(network, addr, timeout)
} }
return dialViaProxy(u, addr, timeout) return dialViaProxy(u, addr, timeout)
+133
View File
@@ -0,0 +1,133 @@
package netproxy
import (
"net"
"net/http"
"net/url"
"os"
"path/filepath"
"testing"
"time"
"nukumizu-backend/config"
)
// publishConfig writes a config.json and makes it the configuration in effect,
// which is what a settings update does.
func publishConfig(t *testing.T, body 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)
}
if _, err := config.LoadGlobalConfig(path); err != nil {
t.Fatalf("LoadGlobalConfig: %v", err)
}
}
// resolve runs a proxy function and returns the URL it chose, or "" when it
// chose a direct connection.
func resolve(t *testing.T, proxy func(*http.Request) (*url.URL, error)) string {
t.Helper()
req, err := http.NewRequest(http.MethodGet, "https://example.com/", nil)
if err != nil {
t.Fatalf("NewRequest: %v", err)
}
u, err := proxy(req)
if err != nil {
t.Fatalf("proxy function: %v", err)
}
if u == nil {
return ""
}
return u.String()
}
// TestProxyFuncResolvesPerCall is the property that makes system.networkProxy
// hot-reloadable: the function handed to a transport keeps reading the live
// configuration instead of the address that was configured when it was built.
func TestProxyFuncResolvesPerCall(t *testing.T) {
publishConfig(t, `{"system":{"networkProxy":"http://127.0.0.1:7890"}}`)
proxy := ProxyFunc(true)
if proxy == nil {
t.Fatal("ProxyFunc(true) returned nil, so the channel would never proxy")
}
if got := resolve(t, proxy); got != "http://127.0.0.1:7890" {
t.Errorf("first resolution = %q", got)
}
// The same function must follow a settings update.
publishConfig(t, `{"system":{"networkProxy":"http://127.0.0.1:8888"}}`)
if got := resolve(t, proxy); got != "http://127.0.0.1:8888" {
t.Errorf("after a settings update the same function resolved %q", got)
}
// Clearing the proxy falls back to a direct connection.
publishConfig(t, `{"system":{"networkProxy":""}}`)
if got := resolve(t, proxy); got != "" {
t.Errorf("a cleared proxy still resolved %q", got)
}
}
func TestProxyFuncOptOutReturnsNil(t *testing.T) {
publishConfig(t, `{"system":{"networkProxy":"http://127.0.0.1:7890"}}`)
// A channel with networkUseProxy off must not be handed a function at all,
// so its transport keeps the default direct dialing.
if ProxyFunc(false) != nil {
t.Error("ProxyFunc(false) must return nil")
}
}
func TestProxyFuncNormalizesMissingScheme(t *testing.T) {
publishConfig(t, `{"system":{"networkProxy":"127.0.0.1:7890"}}`)
if got := resolve(t, ProxyFunc(true)); got != "http://127.0.0.1:7890" {
t.Errorf("resolved %q, want the http:// prefix added", got)
}
}
func TestProxyFuncIgnoresUnusableProxy(t *testing.T) {
// A value that cannot be parsed must leave the client dialing directly
// rather than failing every request.
publishConfig(t, `{"system":{"networkProxy":"://missing-scheme"}}`)
if got := resolve(t, ProxyFunc(true)); got != "" {
t.Errorf("an unparseable proxy resolved %q, want a direct connection", got)
}
}
// TestDialWithTimeoutDialsDirectlyWithoutProxy covers the path a cleared
// system.networkProxy takes: the dialer was built while a proxy was configured,
// and must fall back to a direct dial once there is none.
func TestDialWithTimeoutDialsDirectlyWithoutProxy(t *testing.T) {
publishConfig(t, `{"system":{"networkProxy":"http://127.0.0.1:7890"}}`)
dial := DialWithTimeout(true)
if dial == nil {
t.Fatal("DialWithTimeout(true) returned nil")
}
// No proxy is listening on that address, so a dial attempted now would
// fail; clearing the setting is what makes the direct path reachable.
publishConfig(t, `{"system":{"networkProxy":""}}`)
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen: %v", err)
}
defer ln.Close()
go func() {
conn, err := ln.Accept()
if err == nil {
conn.Close()
}
}()
conn, err := dial("tcp", ln.Addr().String(), 5*time.Second)
if err != nil {
t.Fatalf("dial through a cleared proxy: %v", err)
}
conn.Close()
}
+38 -72
View File
@@ -58,6 +58,27 @@ func main() {
postLog.SetDebugMode(cfg.System.DebugMode) postLog.SetDebugMode(cfg.System.DebugMode)
postLog.InitLogBroadcaster() postLog.InitLogBroadcaster()
// A settings update replaces the configuration in memory; these hooks push
// the new values into the state that was derived from the old one. The
// logger's debug flag is process-wide rather than read at every log call,
// and each controller holds its own copy of its channel settings plus the
// clients built from them.
config.OnReload(func(updated *config.Config) {
postLog.SetDebugMode(updated.System.DebugMode)
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 dbPath := cfg.DBPath
if err := postLog.InitLogsDatabase(fmt.Sprintf("%s/log.db", dbPath)); err != nil { if err := postLog.InitLogsDatabase(fmt.Sprintf("%s/log.db", dbPath)); err != nil {
@@ -199,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() { func initControllers() {
cfg := config.Current() cfg := config.Current()
mgr := controller.GetManager() mgr := controller.GetManager()
if mgr == nil { if mgr == nil || cfg == nil {
return return
} }
mgr.ReplaceAll(buildControllers(cfg), cfg.ControllerMethod)
// 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())
}
}()
} }
// startBackgroundTasks starts periodic background goroutines. // startBackgroundTasks starts periodic background goroutines.