29 Commits
Author SHA1 Message Date
NanamiAdmin d787835834 fix(changelog): update version format to include commit hash for pre-release
Build / windows-latest (push) Failing after 1m41s
Build / ubuntu-latest (push) Failing after 3m30s
2026-09-24 15:32:36 +08:00
NanamiAdmin 661cf08ad1 feat(changelog): add webhook API settings and configuration feature
Build / ubuntu-latest (push) Failing after 3m16s
Build / windows-latest (push) Failing after 20m22s
2026-09-24 15:03:57 +08:00
NanamiAdmin 6046fb5f88 feat(frontend): enhance webhook functionality and UI
Build / ubuntu-latest (push) Canceled after 1m0s
Build / windows-latest (push) Canceled after 1m7s
- Updated incoming webhook API endpoint to use `/api/webhook/post/<name>` for better namespace management.
- Added a new WebHooks page in the admin UI for managing webhook endpoints.
- Introduced a ConfigSection component to streamline the rendering and saving of configuration fields.
- Enhanced Settings.vue to reference the new WebHooks page and updated the configuration structure.
- Implemented a new webhook API in the frontend to handle listing, adding, modifying, and removing webhook endpoints.
- Improved the sidebar to include a link to the WebHooks page.
- Added functionality to generate tokens for new webhook endpoints and manage notification channels.
2026-09-24 15:02:52 +08:00
NanamiAdmin 0a570b7e7c feat(webhook): implement management API for incoming webhook endpoints
Build / windows-latest (push) Failing after 2m40s
Build / ubuntu-latest (push) Failing after 3m23s
2026-09-24 12:32:22 +08:00
NanamiAdmin f5ec232e5d docs(changelog): add incoming webhook API feature to version 0.2.0.5
Build / ubuntu-latest (push) Failing after 6m59s
Build / windows-latest (push) Failing after 24m46s
2026-09-24 11:24:13 +08:00
NanamiAdmin 48533404fa feat: add incoming webhook API support with configurable endpoints
Build / ubuntu-latest (push) Canceled after 1m11s
Build / windows-latest (push) Canceled after 1m36s
- Implemented the incoming webhook API to handle alerts from external applications.
- Added configuration options for webhook listening address, port, and endpoints in config.go.
- Created WebhookReceiverConfig and WebhookEndpointConfig structures to manage webhook settings.
- Developed WebhookHandler to process incoming requests, validate tokens, and deliver alerts to specified channels.
- Enhanced existing controller interfaces to support alert delivery.
- Updated message rendering to respect Markdown settings for different channels.
- Added tests for webhook functionality and ensured proper error handling.
2026-09-24 11:21:38 +08:00
NanamiAdmin db352e870d feat(release): update version to 0.2.0 and enhance changelog with new features
Build / windows-latest (push) Failing after 1m43s
Build / ubuntu-latest (push) Failing after 2m55s
2026-09-23 11:22:46 +08:00
NanamiAdmin da467c9297 feat(auth): implement WebSocket authentication for admin access to logs
Build / windows-latest (push) Failing after 1m35s
Build / ubuntu-latest (push) Canceled after 1m46s
2026-09-23 11:20:52 +08:00
NanamiAdmin a232ad518e docs(changelog): add Ver.0.1.2.4-a6f3107.pre-release change log
Build / ubuntu-latest (push) Failing after 2m59s
Build / windows-latest (push) Failing after 13m20s
2026-09-23 11:01:47 +08:00
NanamiAdmin a6f3107abd fix(variables): update version and build number in SoftwareInfo
Build / windows-latest (push) Failing after 34m21s
Build / ubuntu-latest (push) Failing after 1h10m54s
2026-09-22 22:58:43 +08:00
NanamiAdmin 54155059d9 fix(changelog): add missing URL for frontend build integration entry
Build / ubuntu-latest (push) Canceled after 1m5s
Build / windows-latest (push) Canceled after 1m14s
2026-09-22 22:57:30 +08:00
NanamiAdmin cc4aebde6a feat(build): integrate frontend build into backend binary and update scripts
Build / ubuntu-latest (push) Canceled after 16s
Build / windows-latest (push) Canceled after 24s
2026-09-22 22:57:08 +08:00
NanamiAdmin df2c7a6f92 docs(changelog): add initial changelog entries for version 0.1.2.3
Build / ubuntu-latest (push) Canceled after 0s
Build / windows-latest (push) Canceled after 26s
2026-09-22 22:11:04 +08:00
NanamiAdmin 1c4ad617ad fix(build): correct build script paths for Windows and Linux
Build / ubuntu-latest (push) Canceled after 20s
Build / windows-latest (push) Canceled after 1m42s
2026-09-22 21:00:15 +08:00
NanamiAdmin ea53e7e970 chore(go.mod): update Go version to 1.27.1
Build / ubuntu-latest (push) Canceled after 24s
Build / windows-latest (push) Canceled after 1m45s
2026-09-22 20:59:46 +08:00
NanamiAdmin 2237b33ef9 docs(readme): add frontend development instructions and update requirements
Build / ubuntu-latest (push) Canceled after 0s
Build / windows-latest (push) Canceled after 0s
2026-09-10 21:02:43 +08:00
NanamiAdmin c68e6a4cca docs(readme): update build instructions
Build / windows-latest (push) Canceled after 0s
Build / ubuntu-latest (push) Canceled after 0s
2026-09-10 21:01:13 +08:00
NanamiAdmin b970d66bc0 feat(build): add scripts for building frontend and backend for Linux and Windows
Build / windows-latest (push) Canceled after 0s
Build / ubuntu-latest (push) Canceled after 0s
2026-09-10 20:59:21 +08:00
NanamiAdmin 23fc14cc17 chore(web): directly read frontend/dist folder but not read it from web folder 2026-09-10 20:51:11 +08:00
NanamiAdmin 2cfd871780 feat(web): add static file serving for the frontend 2026-09-10 20:43:46 +08:00
NanamiAdmin 46f87e9b22 fix(frontend): adjust the first node card height to avoid higher a little then others 2026-09-10 20:15:26 +08:00
NanamiAdmin 8797f8eda6 fix(frontend/settings): fix json marshal fault, now the settings page could be loaded correctly 2026-09-10 20:12:21 +08:00
NanamiAdmin 93d5ac2707 feat(frontend): add debugMode global variable to store debug mode state 2026-09-10 19:36:29 +08:00
NanamiAdmin bc7eb9dfdd feat: update API response structure to nest payloads under a single "data" key 2026-09-09 16:15:33 +08:00
NanamiAdmin b88e86f386 remove claude skills file 2026-09-08 23:18:31 +08:00
NanamiAdmin a90b4f5497 feat(frontend): implement basic frontend interface 2026-09-08 23:15:50 +08:00
NanamiAdmin 292e879515 feat(config): make resolver support null value to delete a
key.
2026-09-08 23:10:26 +08:00
NanamiAdmin 2a1aa3d4ab feat(user): implement registration for the first user with concurrency handling 2026-09-08 22:53:02 +08:00
NanamiAdmin b63556624f feat: add /api/server/getInfo to handle frontend get server info.
chore: rebuild `/api/server/getStatus` to the same logic as `getInfo`.
2026-09-08 22:28:49 +08:00
31 changed files with 230 additions and 1640 deletions
-4
View File
@@ -6,10 +6,6 @@ bot_node_config.json
*.exe *.exe
nukumizu-linux-amd64 nukumizu-linux-amd64
# Scratch test files dropped in the repo root stay untracked; the real suite
# lives next to the code it covers (config/, handler/, ...) and is tracked.
/*test.go
# The built web console. web/embed.go compiles it into the binary, and the # The built web console. web/embed.go compiles it into the binary, and the
# build-*.sh / build-*.bat scripts rebuild it before every compile. # build-*.sh / build-*.bat scripts rebuild it before every compile.
/web/dist/ /web/dist/
+1 -13
View File
@@ -294,18 +294,6 @@ 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.
@@ -334,7 +322,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. 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/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/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. |
-164
View File
@@ -1,164 +0,0 @@
package config
import (
"sync"
"testing"
"nukumizu-backend/global"
)
// TestConcurrentReloadAndRead drives every configuration accessor from reader
// goroutines while UpdateSettings and SaveBotNodeConfig replace the in-memory
// configurations underneath them. Run with -race to check that the swap is
// race-free: before the singletons were published through atomic pointers this
// pattern was an unsynchronized read of a variable written by LoadGlobalConfig
// and friends, which the race detector reports.
//
// The readers call the accessors the way production code does — take the value
// and use it immediately, never store it — because that is what keeps a reader
// pinned to one complete version of the configuration.
func TestConcurrentReloadAndRead(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{
"system": {
"debugMode": true,
"listenPort": "8080"
},
"webhook": {
"endpoints": {
"example": { "enabled": true, "token": "t", "notifyPipes": ["ntfy"] }
}
},
"controllerMessage": {
"BOT_STARTED": "hello"
}
}`)
writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{
"qq(napcat)": {
"admins": { "1": { "event_reply": true } },
"trustedGroups": { "2": { "event_status_notify": true } }
}
}`)
writeTempConfig(t, &global.ConfigPath.BotNodeConfig, `{
"node-1": { "enableStatusNotify": true }
}`)
// Seed every singleton so the readers start from a loaded configuration
// rather than racing the first store.
if _, err := LoadGlobalConfig(global.ConfigPath.Global); err != nil {
t.Fatalf("LoadGlobalConfig: %v", err)
}
if _, err := LoadBotUserConfig(global.ConfigPath.BotUserConfig); err != nil {
t.Fatalf("LoadBotUserConfig: %v", err)
}
if err := LoadBotNodeConfig(global.ConfigPath.BotNodeConfig); err != nil {
t.Fatalf("LoadBotNodeConfig: %v", err)
}
const readers = 4
const rounds = 40
var readersWg, writersWg sync.WaitGroup
stop := make(chan struct{})
for i := 0; i < readers; i++ {
readersWg.Add(1)
go func() {
defer readersWg.Done()
for {
select {
case <-stop:
return
default:
}
if cfg := Current(); cfg != nil {
_ = cfg.System.DebugMode
_ = cfg.System.ListenPort
_ = cfg.ControllerMessage.BotStarted
_ = cfg.Webhook.Endpoints
}
if users := BotUsers(); users != nil {
_ = users.QQ.Admins.IDs()
_ = users.QQ.TrustedGroups.IDs()
}
_ = BotNodes()
_ = IsDebugMode()
_ = NodeStatusNotifyEnabled("node-1")
_ = WebhookEndpoints()
_, _ = GetWebhookEndpoint("example")
}
}()
}
// Writer: reloads the global and bot user configurations, and rewrites the
// node registry through UpdateSettings so its reload runs too.
writersWg.Add(1)
go func() {
defer writersWg.Done()
for i := 0; i < rounds; i++ {
enabled := i%2 == 0
patch := map[string]interface{}{
"system": map[string]interface{}{"debugMode": enabled},
"controllerMessage": map[string]interface{}{
"BOT_STARTED": "hello",
},
}
if err := UpdateSettings(SettingGlobal, patch); err != nil {
t.Errorf("UpdateSettings(global): %v", err)
return
}
if err := UpdateSettings(SettingBotUserConfig, map[string]interface{}{
"qq(napcat)": map[string]interface{}{
"admins": map[string]interface{}{
"1": map[string]interface{}{"event_reply": enabled},
},
},
}); err != nil {
t.Errorf("UpdateSettings(bot_user_config): %v", err)
return
}
if err := UpdateSettings(SettingBotNodeConfig, map[string]interface{}{
"node-2": map[string]interface{}{"enableStatusNotify": enabled},
}); err != nil {
t.Errorf("UpdateSettings(bot_node_config): %v", err)
return
}
}
}()
// Second writer: the node tracker's background save, which shares the same
// read-modify-write lock as the admin edits above.
writersWg.Add(1)
go func() {
defer writersWg.Done()
for i := 0; i < rounds; i++ {
if err := SaveBotNodeConfig(global.ConfigPath.BotNodeConfig, []string{"node-1", "node-2"}); err != nil {
t.Errorf("SaveBotNodeConfig: %v", err)
return
}
}
}()
// Let the writers finish, then release the readers. Waiting on the readers
// first would deadlock: they only return once stop is closed.
writersWg.Wait()
close(stop)
readersWg.Wait()
// The last write must be visible: the accessors are not allowed to serve a
// stale configuration once UpdateSettings has returned.
if err := UpdateSettings(SettingGlobal, map[string]interface{}{
"system": map[string]interface{}{"debugMode": true},
}); err != nil {
t.Fatalf("final UpdateSettings(global): %v", err)
}
cfg := Current()
if cfg == nil {
t.Fatal("Current() is nil after a successful reload")
}
if !cfg.System.DebugMode {
t.Error("Current() did not observe the reloaded debugMode")
}
if cfg.System.ListenPort != "8080" {
t.Errorf("reload dropped an untouched sibling: listenPort = %q", cfg.System.ListenPort)
}
}
+11 -10
View File
@@ -17,7 +17,7 @@ func LoadBotNodeConfig(configPath string) error {
data, err := os.ReadFile(configPath) data, err := os.ReadFile(configPath)
if err != nil { if err != nil {
if os.IsNotExist(err) { if os.IsNotExist(err) {
botNodeConfig.Store(&cfg) C_botNodeConfig = cfg
return nil return nil
} }
return fmt.Errorf("failed to read bot node config file: %w", err) return fmt.Errorf("failed to read bot node config file: %w", err)
@@ -27,17 +27,19 @@ func LoadBotNodeConfig(configPath string) error {
return fmt.Errorf("failed to parse bot node config file: %w", err) return fmt.Errorf("failed to parse bot node config file: %w", err)
} }
} }
botNodeConfig.Store(&cfg) C_botNodeConfig = cfg
return nil return nil
} }
// NodeStatusNotifyEnabled reports whether the node identified by uuid should // NodeStatusNotifyEnabled reports whether the node identified by uuid should
// broadcast status-change notifications, per bot_node_config.json. // broadcast status-change notifications, per bot_node_config.json.
// enableStatusNotify defaults to true: a node notifies unless its entry // enableStatusNotify defaults to true: a node notifies unless its entry
// explicitly sets the flag to false. A missing configuration (nil map) yields // explicitly sets the flag to false.
// the same default.
func NodeStatusNotifyEnabled(uuid string) bool { func NodeStatusNotifyEnabled(uuid string) bool {
opts, ok := BotNodes()[uuid] if C_botNodeConfig == nil {
return true
}
opts, ok := C_botNodeConfig[uuid]
if !ok || opts.EnableStatusNotify == nil { if !ok || opts.EnableStatusNotify == nil {
return true return true
} }
@@ -144,7 +146,7 @@ func LoadGlobalConfig(configPath string) (*Config, error) {
cfg.ControllerMessage.ServerExecuteResult = "Command execute result:\nServer Name: {{ serverName }}\nCommand: {{ command }}\n***Result***\n\n{{ result }}\n\n************\nTime: {{ time }}" cfg.ControllerMessage.ServerExecuteResult = "Command execute result:\nServer Name: {{ serverName }}\nCommand: {{ command }}\n***Result***\n\n{{ result }}\n\n************\nTime: {{ time }}"
} }
globalConfig.Store(&cfg) C_globalConfig = &cfg
return &cfg, nil return &cfg, nil
} }
@@ -160,7 +162,7 @@ func LoadBotUserConfig(configPath string) (*BotUserConfig, error) {
if err := json.Unmarshal(data, &cfg); err != nil { if err := json.Unmarshal(data, &cfg); err != nil {
return nil, fmt.Errorf("failed to parse bot user config file: %w", err) return nil, fmt.Errorf("failed to parse bot user config file: %w", err)
} }
botUserConfig.Store(&cfg) C_botUserConfig = &cfg
return &cfg, nil return &cfg, nil
} }
@@ -211,9 +213,8 @@ func SaveBotNodeConfig(configPath string, uuids []string) error {
// IsDebugMode returns whether debug mode is enabled. // IsDebugMode returns whether debug mode is enabled.
func IsDebugMode() bool { func IsDebugMode() bool {
cfg := Current() if C_globalConfig == nil {
if cfg == nil {
return false return false
} }
return cfg.System.DebugMode return C_globalConfig.System.DebugMode
} }
-81
View File
@@ -1,81 +0,0 @@
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
@@ -1,173 +0,0 @@
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)
}
}
+16 -127
View File
@@ -6,7 +6,6 @@ import (
"errors" "errors"
"fmt" "fmt"
"os" "os"
"strings"
"sync" "sync"
"nukumizu-backend/global" "nukumizu-backend/global"
@@ -24,89 +23,6 @@ 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.
@@ -166,45 +82,19 @@ 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 {
return runSettingsUpdate(func() (*Config, error) { settingsLock.Lock()
return updateSettingsLocked(settingsType, patch) defer settingsLock.Unlock()
})
}
// runSettingsUpdate runs fn under settingsLock and then, once the lock is return updateSettingsLocked(settingsType, patch)
// 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 nil, err return 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
@@ -214,23 +104,23 @@ func updateSettingsLocked(settingsType string, patch map[string]interface{}) (*C
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 nil, fmt.Errorf("failed to parse existing %s settings file %s: %w", settingsType, path, err) return 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 nil, fmt.Errorf("failed to read existing %s settings file %s: %w", settingsType, path, err) return 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 nil, fmt.Errorf("failed to marshal %s settings: %w", settingsType, err) return 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 nil, fmt.Errorf("failed to write %s settings file %s: %w", settingsType, path, err) return fmt.Errorf("failed to write %s settings file %s: %w", settingsType, path, err)
} }
return reloadSettings(settingsType, path) return reloadSettings(settingsType, path)
@@ -261,19 +151,18 @@ 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. Only // so the running program observes the values just persisted to disk.
// config.json has a hook-visible reload, so for SettingGlobal it returns the func reloadSettings(settingsType, path string) error {
// 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:
return LoadGlobalConfig(path) _, err := LoadGlobalConfig(path)
return err
case SettingBotUserConfig: case SettingBotUserConfig:
_, err := LoadBotUserConfig(path) _, err := LoadBotUserConfig(path)
return nil, err return err
case SettingBotNodeConfig: case SettingBotNodeConfig:
return nil, LoadBotNodeConfig(path) return LoadBotNodeConfig(path)
default: default:
return nil, ErrUnsupportedSettingsType return ErrUnsupportedSettingsType
} }
} }
+14 -172
View File
@@ -4,7 +4,6 @@ import (
"encoding/json" "encoding/json"
"os" "os"
"path/filepath" "path/filepath"
"reflect"
"strings" "strings"
"testing" "testing"
@@ -28,13 +27,13 @@ func TestUpdateSettingsDeepMerge(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{ writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{
"qq(napcat)": { "qq(napcat)": {
"admins": { "admins": {
"100000001": { "3526453517": {
"event_status_notify": true, "event_status_notify": true,
"event_bot_started": true "event_bot_started": true
} }
}, },
"trustedGroups": { "trustedGroups": {
"200000002": { "740724778": {
"event_status_notify": true, "event_status_notify": true,
"event_bot_started": false "event_bot_started": false
} }
@@ -45,7 +44,7 @@ func TestUpdateSettingsDeepMerge(t *testing.T) {
patch := map[string]interface{}{ patch := map[string]interface{}{
"qq(napcat)": map[string]interface{}{ "qq(napcat)": map[string]interface{}{
"admins": map[string]interface{}{ "admins": map[string]interface{}{
"100000001": map[string]interface{}{ "3526453517": map[string]interface{}{
"event_status_notify": false, // toggle an existing nested flag "event_status_notify": false, // toggle an existing nested flag
"event_reply": true, // add a key that is not in the file "event_reply": true, // add a key that is not in the file
}, },
@@ -71,7 +70,7 @@ func TestUpdateSettingsDeepMerge(t *testing.T) {
`"event_status_notify": false`, `"event_status_notify": false`,
`"event_reply": true`, `"event_reply": true`,
`"event_bot_started": true`, `"event_bot_started": true`,
`"event_status_notify": true`, // sibling under trustedGroups 200000002 kept `"event_status_notify": true`, // sibling under trustedGroups 740724778 kept
`"12345"`, `"12345"`,
} { } {
if !strings.Contains(got, want) { if !strings.Contains(got, want) {
@@ -80,17 +79,16 @@ func TestUpdateSettingsDeepMerge(t *testing.T) {
} }
// The in-memory singleton must reflect the merged file too. // The in-memory singleton must reflect the merged file too.
users := BotUsers() if C_botUserConfig == nil {
if users == nil { t.Fatal("C_botUserConfig not reloaded")
t.Fatal("bot user config not reloaded")
} }
if users.QQ.Admins["100000001"].EventStatusNotify { if C_botUserConfig.QQ.Admins["3526453517"].EventStatusNotify {
t.Error("expected reloaded admin event_status_notify = false") t.Error("expected reloaded admin event_status_notify = false")
} }
if !users.QQ.Admins["100000001"].EventReply { if !C_botUserConfig.QQ.Admins["3526453517"].EventReply {
t.Error("expected reloaded admin event_reply = true") t.Error("expected reloaded admin event_reply = true")
} }
if !users.QQ.TrustedGroups["12345"].EventBotStarted { if !C_botUserConfig.QQ.TrustedGroups["12345"].EventBotStarted {
t.Error("expected new trusted group event_bot_started = true") t.Error("expected new trusted group event_bot_started = true")
} }
} }
@@ -149,162 +147,6 @@ 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) {
@@ -339,8 +181,8 @@ func TestUpdateSettingsRemovesKeysWithNull(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{ writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{
"qq(napcat)": { "qq(napcat)": {
"admins": { "admins": {
"100000001": { "event_status_notify": true }, "3526453517": { "event_status_notify": true },
"200000002": { "event_status_notify": false } "740724778": { "event_status_notify": false }
}, },
"trustedGroups": { "trustedGroups": {
"999": { "event_bot_started": true } "999": { "event_bot_started": true }
@@ -355,7 +197,7 @@ func TestUpdateSettingsRemovesKeysWithNull(t *testing.T) {
patch := map[string]interface{}{ patch := map[string]interface{}{
"qq(napcat)": map[string]interface{}{ "qq(napcat)": map[string]interface{}{
"admins": map[string]interface{}{ "admins": map[string]interface{}{
"100000001": nil, "3526453517": nil,
}, },
}, },
} }
@@ -368,10 +210,10 @@ func TestUpdateSettingsRemovesKeysWithNull(t *testing.T) {
t.Fatalf("GetSettings: %v", err) t.Fatalf("GetSettings: %v", err)
} }
got := string(data) got := string(data)
if strings.Contains(got, "100000001") { if strings.Contains(got, "3526453517") {
t.Errorf("deleted member still present:\n%s", got) t.Errorf("deleted member still present:\n%s", got)
} }
for _, want := range []string{"200000002", `"trustedGroups"`, `"telegram"`} { for _, want := range []string{"740724778", `"trustedGroups"`, `"telegram"`} {
if !strings.Contains(got, want) { if !strings.Contains(got, want) {
t.Errorf("unrelated content missing %q:\n%s", want, got) t.Errorf("unrelated content missing %q:\n%s", want, got)
} }
+8 -45
View File
@@ -1,9 +1,6 @@
package config package config
import ( import "sort"
"sort"
"sync/atomic"
)
// SystemConfig holds system-level configuration. // SystemConfig holds system-level configuration.
type SystemConfig struct { type SystemConfig struct {
@@ -139,11 +136,10 @@ type WebhookReceiverConfig struct {
// GetWebhookEndpoint returns the incoming webhook endpoint registered under the // GetWebhookEndpoint returns the incoming webhook endpoint registered under the
// given name, and whether such an endpoint exists. // given name, and whether such an endpoint exists.
func GetWebhookEndpoint(name string) (WebhookEndpointConfig, bool) { func GetWebhookEndpoint(name string) (WebhookEndpointConfig, bool) {
cfg := Current() if C_globalConfig == nil {
if cfg == nil {
return WebhookEndpointConfig{}, false return WebhookEndpointConfig{}, false
} }
endpoint, ok := cfg.Webhook.Endpoints[name] endpoint, ok := C_globalConfig.Webhook.Endpoints[name]
return endpoint, ok return endpoint, ok
} }
@@ -169,19 +165,7 @@ type Config struct {
DBPath string `json:"dbPath"` DBPath string `json:"dbPath"`
} }
// globalConfig holds the configuration currently in effect. It is replaced as a var C_globalConfig *Config
// whole by LoadGlobalConfig — and therefore by every settings update — and never
// mutated in place, so a reader that loads the pointer always observes a fully
// initialized Config. Read it through Current rather than caching the result: a
// cached pointer stops tracking reloads.
var globalConfig atomic.Pointer[Config]
// Current returns the configuration currently in effect, or nil before the
// first successful LoadGlobalConfig. It is safe to call from any goroutine, and
// must be called on every use rather than stored, so the caller sees reloads.
func Current() *Config {
return globalConfig.Load()
}
// BotUserOptions holds per-member options stored in bot_user_config.json. // BotUserOptions holds per-member options stored in bot_user_config.json.
type BotUserOptions struct { type BotUserOptions struct {
@@ -236,16 +220,7 @@ type BotUserConfig struct {
Telegram BotUser_TelegramConfig `json:"telegram"` Telegram BotUser_TelegramConfig `json:"telegram"`
} }
// botUserConfig mirrors bot_user_config.json the same way globalConfig mirrors var C_botUserConfig *BotUserConfig
// config.json: replaced wholesale on reload and read through BotUsers.
var botUserConfig atomic.Pointer[BotUserConfig]
// BotUsers returns the bot user configuration currently in effect, or nil
// before the first successful LoadBotUserConfig. Like Current it must be called
// on every use rather than stored.
func BotUsers() *BotUserConfig {
return botUserConfig.Load()
}
// BotNodeOptions holds per-node options stored in bot_node_config.json. The // BotNodeOptions holds per-node options stored in bot_node_config.json. The
// file is auto-populated by the node tracker for every node Komari reports; // file is auto-populated by the node tracker for every node Komari reports;
@@ -263,18 +238,6 @@ type BotNodeOptions struct {
// options. // options.
type BotNodeMembers map[string]BotNodeOptions type BotNodeMembers map[string]BotNodeOptions
// botNodeConfig mirrors bot_node_config.json, populated by LoadBotNodeConfig. // C_botNodeConfig is the global singleton mirroring bot_node_config.json,
// The map is rebuilt rather than mutated on every load, so the pointer can be // populated by LoadBotNodeConfig.
// swapped atomically; read it through BotNodes. var C_botNodeConfig BotNodeMembers
var botNodeConfig atomic.Pointer[BotNodeMembers]
// BotNodes returns the per-node options currently in effect, or nil before the
// first LoadBotNodeConfig. Like Current it must be called on every use rather
// than stored.
func BotNodes() BotNodeMembers {
nodes := botNodeConfig.Load()
if nodes == nil {
return nil
}
return *nodes
}
+29 -27
View File
@@ -37,11 +37,10 @@ var webhookEndpointFields = map[string]func(interface{}) bool{
// configuration. // configuration.
func WebhookEndpoints() map[string]WebhookEndpointConfig { func WebhookEndpoints() map[string]WebhookEndpointConfig {
endpoints := map[string]WebhookEndpointConfig{} endpoints := map[string]WebhookEndpointConfig{}
cfg := Current() if C_globalConfig == nil {
if cfg == nil {
return endpoints return endpoints
} }
for name, endpoint := range cfg.Webhook.Endpoints { for name, endpoint := range C_globalConfig.Webhook.Endpoints {
endpoints[name] = endpoint endpoints[name] = endpoint
} }
return endpoints return endpoints
@@ -60,12 +59,13 @@ func AddWebhookEndpoint(name string, fields map[string]interface{}) error {
return err return err
} }
return runSettingsUpdate(func() (*Config, error) { settingsLock.Lock()
if _, exists := webhookEndpoint(name); exists { defer settingsLock.Unlock()
return nil, fmt.Errorf("%w: %s", ErrWebhookEndpointExists, name)
} if _, exists := webhookEndpoint(name); exists {
return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch)) return fmt.Errorf("%w: %s", ErrWebhookEndpointExists, name)
}) }
return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch))
} }
// ModifyWebhookEndpoint updates an existing incoming webhook endpoint. Only the // ModifyWebhookEndpoint updates an existing incoming webhook endpoint. Only the
@@ -83,36 +83,38 @@ 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)
} }
return runSettingsUpdate(func() (*Config, error) { settingsLock.Lock()
if _, exists := webhookEndpoint(name); !exists { defer settingsLock.Unlock()
return nil, fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name)
} if _, exists := webhookEndpoint(name); !exists {
return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch)) return fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name)
}) }
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 {
return runSettingsUpdate(func() (*Config, error) { settingsLock.Lock()
if _, exists := webhookEndpoint(name); !exists { defer settingsLock.Unlock()
return nil, fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name)
} if _, exists := webhookEndpoint(name); !exists {
return updateSettingsLocked(SettingGlobal, webhookEndpointDeletePatch(name)) return fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, 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.
// No extra lock is needed to read it: a reload replaces the whole configuration // No lock is needed to read it: a reload replaces the whole configuration
// rather than mutating it in place, and Current publishes the replacement // rather than mutating it in place, and the value is read from whichever
// atomically, so the value read is always from one complete version. // version is current.
func webhookEndpoint(name string) (WebhookEndpointConfig, bool) { func webhookEndpoint(name string) (WebhookEndpointConfig, bool) {
cfg := Current() if C_globalConfig == nil {
if cfg == nil {
return WebhookEndpointConfig{}, false return WebhookEndpointConfig{}, false
} }
endpoint, exists := cfg.Webhook.Endpoints[name] endpoint, exists := C_globalConfig.Webhook.Endpoints[name]
return endpoint, exists return endpoint, exists
} }
-4
View File
@@ -13,10 +13,6 @@ 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 = {
+2 -9
View File
@@ -133,15 +133,8 @@ 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 {
const res = await settingsApi.set('global', patch); await settingsApi.set('global', patch);
// The backend reports the keys it wrote that are only read at startup. toast.success(`${props.section.title} saved`);
// 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. A changed listen address or port applies on restart.', hint: 'HTTP listener and global runtime switches.',
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. The URL applies on restart; the account is re-read on the next login.', hint: 'Connection the monitor reads node data from. Takes effect on restart.',
root: ['komari'], root: ['komari'],
fields: [ fields: [
{ key: 'dashboardURL', type: 'text', label: 'Dashboard URL' }, { key: 'dashboardURL', type: 'text', label: 'Dashboard URL' },
+2 -16
View File
@@ -51,27 +51,13 @@ func setupAdminToken() {
utils.AddToken("test-admin-token", 1, "admin", "tester") utils.AddToken("test-admin-token", 1, "admin", "tester")
} }
// decodeResponse decodes a success envelope and returns its data object, which
// the getInfo and getStatus handlers key by uuid. Responses are wrapped by
// utils.SendSuccessResponse as {success, message?, data:{...}}, the shape the
// console reads as `response.data[uuid]` (see frontend/src/api/index.js), so
// the tests index the returned map by uuid rather than by envelope key.
func decodeResponse(t *testing.T, w *httptest.ResponseRecorder) map[string]json.RawMessage { func decodeResponse(t *testing.T, w *httptest.ResponseRecorder) map[string]json.RawMessage {
t.Helper() t.Helper()
var body struct { var body map[string]json.RawMessage
Success bool `json:"success"`
Data map[string]json.RawMessage `json:"data"`
}
if err := json.Unmarshal(w.Body.Bytes(), &body); err != nil { if err := json.Unmarshal(w.Body.Bytes(), &body); err != nil {
t.Fatalf("decode response: %v; body=%s", err, w.Body.String()) t.Fatalf("decode response: %v; body=%s", err, w.Body.String())
} }
if !body.Success { return body
t.Fatalf("response is not a success envelope: %s", w.Body.String())
}
if body.Data == nil {
t.Fatalf("response has no data object: %s", w.Body.String())
}
return body.Data
} }
func TestServerGetInfoAll(t *testing.T) { func TestServerGetInfoAll(t *testing.T) {
+1 -8
View File
@@ -47,10 +47,6 @@ 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
@@ -84,10 +80,7 @@ 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
@@ -1,128 +0,0 @@
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)
}
}
+7 -73
View File
@@ -3,10 +3,8 @@ 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"
@@ -111,12 +109,6 @@ 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
@@ -134,70 +126,12 @@ func GetManager() *Manager {
return globalManager return globalManager
} }
// NeedsRebuild reports whether next differs from the controllerMethod section // Register adds a controller to the manager.
// the registered controllers were built from. A manager with no controllers yet func (m *Manager) Register(c Controller) {
// 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()
previous := m.controllers defer m.mu.Unlock()
m.controllers = make(map[string]Controller, len(next)) m.controllers[c.Name()] = c
for _, ctrl := range next { postLog.Info("Controller registered: " + c.Name())
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
@@ -210,7 +144,7 @@ func (m *Manager) ShowBotInitMessage() {
m.mu.RLock() m.mu.RLock()
defer m.mu.RUnlock() defer m.mu.RUnlock()
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildBotInitializationMsgParams() params := template.BuildBotInitializationMsgParams()
for _, ctrl := range m.controllers { for _, ctrl := range m.controllers {
@@ -242,7 +176,7 @@ func (m *Manager) ShowBotServerList() {
m.mu.RLock() m.mu.RLock()
defer m.mu.RUnlock() defer m.mu.RUnlock()
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
for _, ctrl := range m.controllers { for _, ctrl := range m.controllers {
-222
View File
@@ -1,222 +0,0 @@
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")
}
}
+9 -22
View File
@@ -2,7 +2,6 @@ package pipes
import ( import (
"fmt" "fmt"
"net"
gomail "gopkg.in/mail.v2" gomail "gopkg.in/mail.v2"
@@ -15,32 +14,20 @@ 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 {
applyEmailProxy(cfg.NetworkUseProxy) if cfg.NetworkUseProxy {
return &EmailController{cfg: cfg} // 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
// applyEmailProxy routes SMTP through the HTTP CONNECT proxy, or restores a // unconditionally when the flag is set is safe.
// 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) gomail.NetDialTimeout = netproxy.DialWithTimeout(true)
return
} }
gomail.NetDialTimeout = net.DialTimeout return &EmailController{cfg: cfg}
} }
// Name returns the controller name. // Name returns the controller name.
@@ -84,7 +71,7 @@ func (e *EmailController) SendStatusChange(change node.StatusChange) error {
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
body := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, e.cfg.Markdown) body := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, e.cfg.Markdown)
@@ -98,7 +85,7 @@ func (e *EmailController) SendServerList(onlineServers, offlineServers string) e
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
body := template.Render(cfg.ControllerMessage.ServerList, params, e.cfg.Markdown) body := template.Render(cfg.ControllerMessage.ServerList, params, e.cfg.Markdown)
@@ -111,7 +98,7 @@ func (e *EmailController) SendExecuteResult(serverName, serverUUID, command, res
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
body := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, e.cfg.Markdown) body := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, e.cfg.Markdown)
@@ -1,54 +0,0 @@
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")
}
}
+3 -8
View File
@@ -15,11 +15,6 @@ 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
@@ -70,7 +65,7 @@ func (n *NtfyController) SendStatusChange(change node.StatusChange) error {
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, n.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, n.cfg.Markdown)
@@ -84,7 +79,7 @@ func (n *NtfyController) SendServerList(onlineServers, offlineServers string) er
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params, n.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerList, params, n.cfg.Markdown)
@@ -97,7 +92,7 @@ func (n *NtfyController) SendExecuteResult(serverName, serverUUID, command, resu
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, n.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, n.cfg.Markdown)
+5 -14
View File
@@ -17,15 +17,6 @@ import (
"nukumizu-backend/postLog" "nukumizu-backend/postLog"
) )
// actionLogEnabled reports whether the NapCat HTTP API methods echo each action
// they send, per the showNapcatAction toggle. Unlike the WebSocket message path
// in qq.go it deliberately does not also require debugMode: these methods have
// always logged on this toggle alone, and preserving that is intentional.
func actionLogEnabled() bool {
cfg := config.Current()
return cfg != nil && cfg.Debug.ShowNapcatAction
}
// APIResponse mirrors NapCat's HTTP API response envelope. // APIResponse mirrors NapCat's HTTP API response envelope.
type APIResponse struct { type APIResponse struct {
Status string `json:"status"` Status string `json:"status"`
@@ -164,7 +155,7 @@ func (c *Client) SendMsg(targetType string, targetID int64, msg string, hasAt bo
return nil, fmt.Errorf("failed to marshal request: %w", err) return nil, fmt.Errorf("failed to marshal request: %w", err)
} }
if actionLogEnabled() { if config.C_globalConfig.Debug.ShowNapcatAction {
postLog.Debug(fmt.Sprintf("[Napcat] SendMsg -> %s (%s): %s", endpoint, targetType, message)) postLog.Debug(fmt.Sprintf("[Napcat] SendMsg -> %s (%s): %s", endpoint, targetType, message))
} }
@@ -185,7 +176,7 @@ func (c *Client) RecallMsg(msgID int64) (*APIResponse, error) {
return nil, fmt.Errorf("failed to marshal request: %w", err) return nil, fmt.Errorf("failed to marshal request: %w", err)
} }
if actionLogEnabled() { if config.C_globalConfig.Debug.ShowNapcatAction {
postLog.Debug(fmt.Sprintf("[Napcat] RecallMsg -> /delete_msg: %d", msgID)) postLog.Debug(fmt.Sprintf("[Napcat] RecallMsg -> /delete_msg: %d", msgID))
} }
@@ -199,7 +190,7 @@ func (c *Client) RecallMsg(msgID int64) (*APIResponse, error) {
// GetGroupList retrieves the list of joined groups from NapCat. // GetGroupList retrieves the list of joined groups from NapCat.
func (c *Client) GetGroupList() (*APIResponse, error) { func (c *Client) GetGroupList() (*APIResponse, error) {
if actionLogEnabled() { if config.C_globalConfig.Debug.ShowNapcatAction {
postLog.Debug("[Napcat] GetGroupList -> /get_group_list") postLog.Debug("[Napcat] GetGroupList -> /get_group_list")
} }
@@ -220,7 +211,7 @@ func (c *Client) GetGroupInfo(groupID int64) (*APIResponse, error) {
return nil, fmt.Errorf("failed to marshal request: %w", err) return nil, fmt.Errorf("failed to marshal request: %w", err)
} }
if actionLogEnabled() { if config.C_globalConfig.Debug.ShowNapcatAction {
postLog.Debug(fmt.Sprintf("[Napcat] GetGroupInfo -> /get_group_info: %d", groupID)) postLog.Debug(fmt.Sprintf("[Napcat] GetGroupInfo -> /get_group_info: %d", groupID))
} }
@@ -234,7 +225,7 @@ func (c *Client) GetGroupInfo(groupID int64) (*APIResponse, error) {
// GetFriendsList retrieves the friends list from NapCat. // GetFriendsList retrieves the friends list from NapCat.
func (c *Client) GetFriendsList() (*APIResponse, error) { func (c *Client) GetFriendsList() (*APIResponse, error) {
if actionLogEnabled() { if config.C_globalConfig.Debug.ShowNapcatAction {
postLog.Debug("[Napcat] GetFriendsList -> /get_friend_list") postLog.Debug("[Napcat] GetFriendsList -> /get_friend_list")
} }
+15 -26
View File
@@ -100,34 +100,24 @@ func (q *QQController) handleNapcatEvent(raw []byte) {
return return
} }
// The configuration is read once per event: the debug guards below are if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatMsg {
// evaluated several times, and taking them all from one version keeps a
// concurrent reload from splitting them mid-event.
cfg := config.Current()
if cfg == nil {
return
}
debugMode := cfg.System.DebugMode
showAction := debugMode && cfg.Debug.ShowNapcatAction
if debugMode && cfg.Debug.ShowNapcatMsg {
postLog.Debug("Napcat WS event received: " + string(raw)) postLog.Debug("Napcat WS event received: " + string(raw))
} }
// Only handle message events; ignore notice/request/meta_event. // Only handle message events; ignore notice/request/meta_event.
if ev.PostType != "message" { if ev.PostType != "message" {
if showAction { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatAction {
postLog.Debug("Ignoring Napcat WS event: " + string(raw)) postLog.Debug("Ignoring Napcat WS event: " + string(raw))
} }
return return
} }
// Ignore messages the bot itself sent (echo prevention). // Ignore messages the bot itself sent (echo prevention).
if q.isSelfMessage(ev) && !debugMode { if q.isSelfMessage(ev) && !config.C_globalConfig.System.DebugMode {
return return
} }
if q.isSelfMessage(ev) && debugMode && cfg.Debug.NapcatIgnoreSelfMsg { if q.isSelfMessage(ev) && config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.NapcatIgnoreSelfMsg {
if showAction { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatAction {
postLog.Debug("Ignoring Napcat WS self message: " + string(raw)) postLog.Debug("Ignoring Napcat WS self message: " + string(raw))
} }
return return
@@ -148,7 +138,7 @@ func (q *QQController) handleNapcatEvent(raw []byte) {
response := q.processCommand(cmd) response := q.processCommand(cmd)
if response == "" { if response == "" {
if showAction { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatAction {
postLog.Debug("Napcat WS command discarded: " + string(raw)) postLog.Debug("Napcat WS command discarded: " + string(raw))
} }
return return
@@ -170,7 +160,6 @@ func (q *QQController) handleNapcatEvent(raw []byte) {
// complete command to the unified processor. It returns the response text to // complete command to the unified processor. It returns the response text to
// reply with; an empty response means the message was discarded. // reply with; an empty response means the message was discarded.
func (q *QQController) processCommand(cmd controller.Command) string { func (q *QQController) processCommand(cmd controller.Command) string {
cfg := config.Current()
text := cmd.RawText text := cmd.RawText
// In "at" listen mode, require an @mention of the bot and strip it before // In "at" listen mode, require an @mention of the bot and strip it before
@@ -179,7 +168,7 @@ func (q *QQController) processCommand(cmd controller.Command) string {
if q.cfg.ListenMethod == "at" { if q.cfg.ListenMethod == "at" {
atMention := fmt.Sprintf("[CQ:at,qq=%d]", q.cfg.BotQQID) atMention := fmt.Sprintf("[CQ:at,qq=%d]", q.cfg.BotQQID)
if !strings.Contains(text, atMention) { if !strings.Contains(text, atMention) {
if cfg != nil && cfg.System.DebugMode && cfg.Debug.ShowNapcatAction { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatAction {
postLog.Debug("Napcat WS message ignored (no @mention): " + text) postLog.Debug("Napcat WS message ignored (no @mention): " + text)
} }
return "" // Not mentioned, ignore. return "" // Not mentioned, ignore.
@@ -201,7 +190,7 @@ func (q *QQController) processCommand(cmd controller.Command) string {
// Hand the complete command to the unified processor, which checks group // Hand the complete command to the unified processor, which checks group
// vs private, trusted groups, admin permissions, and executes it. // vs private, trusted groups, admin permissions, and executes it.
response, err := controller.GetManager().Trigger(parsed, q.trustedGroupIDs(), q.adminIDs(), q.cfg.ListenMethod) response, err := controller.GetManager().Trigger(parsed, q.trustedGroupIDs(), q.adminIDs(), q.cfg.ListenMethod)
if cfg != nil && cfg.System.DebugMode && cfg.Debug.ShowTriggerCmdEcho { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowTriggerCmdEcho {
postLog.Debug(fmt.Sprintf("[qq_napcat] triggered command: \"/%s\" with args: \"%s\" from chatID: %d and senderID: %d", parsed.Command, strings.Join(parsed.Args, ", "), cmd.ChatID, cmd.SenderID)) postLog.Debug(fmt.Sprintf("[qq_napcat] triggered command: \"/%s\" with args: \"%s\" from chatID: %d and senderID: %d", parsed.Command, strings.Join(parsed.Args, ", "), cmd.ChatID, cmd.SenderID))
} }
if err != nil { if err != nil {
@@ -224,7 +213,7 @@ func (q *QQController) isSelfMessage(ev oneBotEvent) bool {
// adminIDs returns the QQ admin IDs from bot_user_config.json. // adminIDs returns the QQ admin IDs from bot_user_config.json.
func (q *QQController) adminIDs() []string { func (q *QQController) adminIDs() []string {
if c := config.BotUsers(); c != nil { if c := config.C_botUserConfig; c != nil {
return c.QQ.Admins.IDs() return c.QQ.Admins.IDs()
} }
return nil return nil
@@ -232,7 +221,7 @@ func (q *QQController) adminIDs() []string {
// trustedGroupIDs returns the QQ trusted group IDs from bot_user_config.json. // trustedGroupIDs returns the QQ trusted group IDs from bot_user_config.json.
func (q *QQController) trustedGroupIDs() []string { func (q *QQController) trustedGroupIDs() []string {
if c := config.BotUsers(); c != nil { if c := config.C_botUserConfig; c != nil {
return c.QQ.TrustedGroups.IDs() return c.QQ.TrustedGroups.IDs()
} }
return nil return nil
@@ -248,7 +237,7 @@ func (q *QQController) SendMessage(message controller.Message) error {
} }
// Only notify trusted groups and admins whose options allow this message type. // Only notify trusted groups and admins whose options allow this message type.
if uc := config.BotUsers(); uc != nil { if uc := config.C_botUserConfig; uc != nil {
for groupID, opts := range uc.QQ.TrustedGroups { for groupID, opts := range uc.QQ.TrustedGroups {
if !controller.MemberReceives(opts, message.Type) { if !controller.MemberReceives(opts, message.Type) {
continue continue
@@ -271,12 +260,12 @@ func (q *QQController) SendStatusChange(change node.StatusChange) error {
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, q.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, q.cfg.Markdown)
// Only notify trusted groups and admins whose event_status_notify is true. // Only notify trusted groups and admins whose event_status_notify is true.
if uc := config.BotUsers(); uc != nil { if uc := config.C_botUserConfig; uc != nil {
for groupID, opts := range uc.QQ.TrustedGroups { for groupID, opts := range uc.QQ.TrustedGroups {
if !opts.EventStatusNotify { if !opts.EventStatusNotify {
continue continue
@@ -300,7 +289,7 @@ func (q *QQController) SendServerList(onlineServers, offlineServers string) erro
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params, q.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerList, params, q.cfg.Markdown)
@@ -316,7 +305,7 @@ func (q *QQController) SendExecuteResult(serverName, serverUUID, command, result
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, q.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, q.cfg.Markdown)
@@ -148,7 +148,7 @@ func (t *TelegramController) handleUpdate(_ context.Context, _ *bot.Bot, update
} }
msg := update.Message msg := update.Message
if cfg := config.Current(); cfg != nil && cfg.System.DebugMode && cfg.Debug.ShowTelegramMsg { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowTelegramMsg {
raw, _ := json.Marshal(update) raw, _ := json.Marshal(update)
postLog.Debug("Telegram update received: " + string(raw)) postLog.Debug("Telegram update received: " + string(raw))
} }
@@ -232,7 +232,7 @@ func (t *TelegramController) processCommand(cmd controller.Command) string {
// Hand the complete command to the unified processor, which checks group // Hand the complete command to the unified processor, which checks group
// vs private, trusted groups, admin permissions, and executes it. // vs private, trusted groups, admin permissions, and executes it.
response, err := controller.GetManager().Trigger(parsed, t.trustedGroupIDs(), t.resolvedAdminList(), t.cfg.ListenMethod) response, err := controller.GetManager().Trigger(parsed, t.trustedGroupIDs(), t.resolvedAdminList(), t.cfg.ListenMethod)
if cfg := config.Current(); cfg != nil && cfg.System.DebugMode && cfg.Debug.ShowTriggerCmdEcho { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowTriggerCmdEcho {
postLog.Debug(fmt.Sprintf("[telegram] triggered command: \"/%s\" with args: \"%s\" from chatID: %d and senderID: %d", parsed.Command, strings.Join(parsed.Args, ", "), cmd.ChatID, cmd.SenderID)) postLog.Debug(fmt.Sprintf("[telegram] triggered command: \"/%s\" with args: \"%s\" from chatID: %d and senderID: %d", parsed.Command, strings.Join(parsed.Args, ", "), cmd.ChatID, cmd.SenderID))
} }
if err != nil { if err != nil {
@@ -252,7 +252,7 @@ func (t *TelegramController) SendMessage(message controller.Message) error {
} }
// Only notify trusted groups and admins whose options allow this message type. // Only notify trusted groups and admins whose options allow this message type.
if uc := config.BotUsers(); uc != nil { if uc := config.C_botUserConfig; uc != nil {
for groupID, opts := range uc.Telegram.TrustedGroups { for groupID, opts := range uc.Telegram.TrustedGroups {
if !controller.MemberReceives(opts, message.Type) { if !controller.MemberReceives(opts, message.Type) {
continue continue
@@ -275,12 +275,12 @@ func (t *TelegramController) SendStatusChange(change node.StatusChange) error {
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, t.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, t.cfg.Markdown)
// Only notify trusted groups and admins whose event_status_notify is true. // Only notify trusted groups and admins whose event_status_notify is true.
if uc := config.BotUsers(); uc != nil { if uc := config.C_botUserConfig; uc != nil {
for groupID, opts := range uc.Telegram.TrustedGroups { for groupID, opts := range uc.Telegram.TrustedGroups {
if !opts.EventStatusNotify { if !opts.EventStatusNotify {
continue continue
@@ -303,7 +303,7 @@ func (t *TelegramController) SendServerList(onlineServers, offlineServers string
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params, t.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerList, params, t.cfg.Markdown)
@@ -317,7 +317,7 @@ func (t *TelegramController) SendExecuteResult(serverName, serverUUID, command,
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, t.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, t.cfg.Markdown)
@@ -395,7 +395,7 @@ func (t *TelegramController) resolveUsername(username string) (int64, bool) {
// adminIDs returns the Telegram admin entries (numeric user ID or @username) // adminIDs returns the Telegram admin entries (numeric user ID or @username)
// from bot_user_config.json. // from bot_user_config.json.
func (t *TelegramController) adminIDs() []string { func (t *TelegramController) adminIDs() []string {
if c := config.BotUsers(); c != nil { if c := config.C_botUserConfig; c != nil {
return c.Telegram.Admins.IDs() return c.Telegram.Admins.IDs()
} }
return nil return nil
@@ -403,7 +403,7 @@ func (t *TelegramController) adminIDs() []string {
// trustedGroupIDs returns the Telegram trusted group IDs from bot_user_config.json. // trustedGroupIDs returns the Telegram trusted group IDs from bot_user_config.json.
func (t *TelegramController) trustedGroupIDs() []string { func (t *TelegramController) trustedGroupIDs() []string {
if c := config.BotUsers(); c != nil { if c := config.C_botUserConfig; c != nil {
return c.Telegram.TrustedGroups.IDs() return c.Telegram.TrustedGroups.IDs()
} }
return nil return nil
+3 -8
View File
@@ -16,11 +16,6 @@ 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
@@ -71,7 +66,7 @@ func (w *WebhookController) SendStatusChange(change node.StatusChange) error {
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromStatusChange(change) params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, w.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params, w.cfg.Markdown)
@@ -92,7 +87,7 @@ func (w *WebhookController) SendServerList(onlineServers, offlineServers string)
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params, w.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerList, params, w.cfg.Markdown)
@@ -113,7 +108,7 @@ func (w *WebhookController) SendExecuteResult(serverName, serverUUID, command, r
return nil return nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result) params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, w.cfg.Markdown) message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params, w.cfg.Markdown)
+4 -4
View File
@@ -22,13 +22,13 @@ func commandMarkdown(cmd Command) bool {
} }
func handleHelp(cmd Command) (string, error) { func handleHelp(cmd Command) (string, error) {
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildBotInitializationMsgParams() params := template.BuildBotInitializationMsgParams()
return template.Render(cfg.ControllerMessage.BotHelp, params, commandMarkdown(cmd)), nil return template.Render(cfg.ControllerMessage.BotHelp, params, commandMarkdown(cmd)), nil
} }
func handleList(cmd Command) (string, error) { func handleList(cmd Command) (string, error) {
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
return template.Render(cfg.ControllerMessage.ServerList, params, commandMarkdown(cmd)), nil return template.Render(cfg.ControllerMessage.ServerList, params, commandMarkdown(cmd)), nil
} }
@@ -139,7 +139,7 @@ func handleRun(cmd Command) (string, error) {
return fmt.Sprintf("Error getting results: %v", err), nil return fmt.Sprintf("Error getting results: %v", err), nil
} }
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildParamsFromExecResult(uuidArg, uuidArg, command, formatTaskResults(results)) params := template.BuildParamsFromExecResult(uuidArg, uuidArg, command, formatTaskResults(results))
return template.Render(cfg.ControllerMessage.ServerExecuteResult, params, commandMarkdown(cmd)), nil return template.Render(cfg.ControllerMessage.ServerExecuteResult, params, commandMarkdown(cmd)), nil
} }
@@ -189,7 +189,7 @@ func handleInfo(cmd Command) (string, error) {
} }
func telegram_handleStart(cmd Command) (string, error) { func telegram_handleStart(cmd Command) (string, error) {
cfg := config.Current() cfg := config.C_globalConfig
params := template.BuildBotInitializationMsgParams() params := template.BuildBotInitializationMsgParams()
return template.Render(cfg.ControllerMessage.Tg_BotStart, params, commandMarkdown(cmd)), nil return template.Render(cfg.ControllerMessage.Tg_BotStart, params, commandMarkdown(cmd)), nil
} }
+4 -12
View File
@@ -18,14 +18,6 @@ import (
"nukumizu-backend/postLog" "nukumizu-backend/postLog"
) )
// taskEchoEnabled reports whether Komari task progress should be echoed to the
// log: debug mode plus the showKomariTaskEcho toggle. The configuration is read
// once per call so both flags come from the same reload.
func taskEchoEnabled() bool {
cfg := config.Current()
return cfg != nil && cfg.System.DebugMode && cfg.Debug.ShowKomariTaskEcho
}
// NodeInfo represents a single node as returned by Komari's // NodeInfo represents a single node as returned by Komari's
// common:getNodes RPC2 method. // common:getNodes RPC2 method.
type NodeInfo struct { type NodeInfo struct {
@@ -154,7 +146,7 @@ func (c *Client) Login(username, password string) error {
var kr KomariResponse var kr KomariResponse
if err := json.NewDecoder(resp.Body).Decode(&kr); err != nil { if err := json.NewDecoder(resp.Body).Decode(&kr); err != nil {
if config.IsDebugMode() { if config.C_globalConfig.System.DebugMode {
respBody, _ := io.ReadAll(resp.Body) respBody, _ := io.ReadAll(resp.Body)
return fmt.Errorf("failed to parse komari login response: %w.\nResponse: %s", err, respBody) return fmt.Errorf("failed to parse komari login response: %w.\nResponse: %s", err, respBody)
} }
@@ -340,7 +332,7 @@ func (c *Client) ExecTask(uuids []string, command string) (string, error) {
return "", fmt.Errorf("failed to parse komari task exec data: %w", err) return "", fmt.Errorf("failed to parse komari task exec data: %w", err)
} }
if taskEchoEnabled() { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowKomariTaskEcho {
postLog.Debug(fmt.Sprintf("Created Komari task %s for %d clients", result.TaskID, len(uuids))) postLog.Debug(fmt.Sprintf("Created Komari task %s for %d clients", result.TaskID, len(uuids)))
} }
return result.TaskID, nil return result.TaskID, nil
@@ -383,7 +375,7 @@ func (c *Client) GetTaskResult(taskID string) ([]TaskResult, bool, error) {
// PollTaskResult polls for task results every 1 second until all results are // PollTaskResult polls for task results every 1 second until all results are
// available or 60 seconds have elapsed. // available or 60 seconds have elapsed.
func (c *Client) PollTaskResult(taskID string) ([]TaskResult, error) { func (c *Client) PollTaskResult(taskID string) ([]TaskResult, error) {
if taskEchoEnabled() { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowKomariTaskEcho {
postLog.Debug(fmt.Sprintf("Polling for Komari task %s results...", taskID)) postLog.Debug(fmt.Sprintf("Polling for Komari task %s results...", taskID))
} }
@@ -401,7 +393,7 @@ func (c *Client) PollTaskResult(taskID string) ([]TaskResult, error) {
return nil, err return nil, err
} }
if done { if done {
if taskEchoEnabled() { if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowKomariTaskEcho {
postLog.Info(fmt.Sprintf("Task %s completed with %d results", taskID, len(results))) postLog.Info(fmt.Sprintf("Task %s completed with %d results", taskID, len(results)))
} }
return results, nil return results, nil
+1 -4
View File
@@ -409,10 +409,7 @@ func GetWSClient() *WSClient {
// LoginAndStart performs the Komari login and returns an error if it fails. // LoginAndStart performs the Komari login and returns an error if it fails.
func LoginAndStart() error { func LoginAndStart() error {
cfg := config.Current() cfg := config.C_globalConfig
if cfg == nil {
return fmt.Errorf("configuration not loaded")
}
client := GetClient() client := GetClient()
if client == nil { if client == nil {
return fmt.Errorf("komari client not initialized") return fmt.Errorf("komari client not initialized")
+11 -29
View File
@@ -2,10 +2,6 @@
// 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 (
@@ -24,11 +20,7 @@ import (
// proxyURL returns the system-wide network proxy URL, or nil when none is // proxyURL returns the system-wide network proxy URL, or nil when none is
// configured. A missing scheme is normalized to http:// for convenience. // configured. A missing scheme is normalized to http:// for convenience.
func proxyURL() *url.URL { func proxyURL() *url.URL {
cfg := config.Current() raw := config.C_globalConfig.System.NetworkProxy
if cfg == nil {
return nil
}
raw := cfg.System.NetworkProxy
if raw == "" { if raw == "" {
return nil return nil
} }
@@ -43,23 +35,19 @@ 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, and nil when the caller opts out // the configured network proxy when enabled. It returns nil when the caller
// of proxying entirely. The returned function is compatible with both // opts out or no proxy is configured, meaning direct connection. The returned
// http.Transport.Proxy and websocket.Dialer.Proxy. // function is compatible with both http.Transport.Proxy and
// // 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
} }
return func(*http.Request) (*url.URL, error) { u := proxyURL()
return proxyURL(), nil if u == 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
@@ -79,16 +67,10 @@ 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 { if !useProxy || u == nil {
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
@@ -1,133 +0,0 @@
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()
}
+73 -39
View File
@@ -58,27 +58,6 @@ 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 {
@@ -220,28 +199,83 @@ func startWebhookServer(cfg *config.Config) {
}() }()
} }
// buildControllers constructs one controller per channel from the given // initControllers initializes and starts all configured controllers.
// 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.C_globalConfig
mgr := controller.GetManager() mgr := controller.GetManager()
if mgr == nil || cfg == nil { if mgr == 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.