Compare commits
6
Commits
ac5587d0c8
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c7642ecca5 | ||
|
|
457ea42417 | ||
|
|
7dc6089871 | ||
|
|
2a5a46a984 | ||
|
|
0691e85ccf | ||
|
|
d9a9918e21 |
@@ -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. |
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
@@ -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
@@ -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, ¤t); err != nil {
|
if err := json.Unmarshal(data, ¤t); 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
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
@@ -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.
|
||||||
|
|||||||
@@ -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 = {
|
||||||
|
|||||||
@@ -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);
|
||||||
|
|||||||
@@ -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
@@ -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),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
|||||||
@@ -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")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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()
|
||||||
|
}
|
||||||
@@ -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.
|
||||||
|
|||||||
Reference in New Issue
Block a user