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

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

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

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

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

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

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

DialWithTimeout reads the address per dial for the same reason, which is what
lets the gomail NetDialTimeout hook follow a settings change.
2026-09-28 23:25:13 +08:00
NanamiAdmin 7dc6089871 docs(readme): correct how system.networkProxy takes effect
Build / ubuntu-latest (push) Canceled after 3m11s
Build / windows-latest (push) Canceled after 3m23s
The rebuild only fires when the controllerMethod section changed, so editing
networkProxy on its own does not rebuild anything and the channels keep the
proxy they connected with. The previous wording read as though a rebuild would
follow on its own.
2026-09-28 23:22:07 +08:00
NanamiAdmin 2a5a46a984 refactor(controller): apply controller settings by rebuilding the channel set
Build / ubuntu-latest (push) Canceled after 20s
Build / windows-latest (push) Canceled after 24s
Giving each controller a Reload method worked for the notification pipes, which
only hold values, but not for the two bot channels: the NapCat client is stopped
through a sync.Once and the Telegram polling context is created with the
controller, so neither can be re-pointed once running. It also meant five
bespoke implementations of the same idea.

Replace the whole set instead. Manager.ReplaceAll swaps in a freshly built set
under the registry lock, stops the outgoing controllers outside it, and starts
the incoming ones off the calling goroutine so a slow handshake does not hold
the settings request open. ReplaceAll is now the only way controllers are
installed: Register is gone, because a lone registration would slip a controller
in without recording the settings it was built from, which is what NeedsRebuild
compares against.

The rebuild runs from the reload hook and only fires when the controllerMethod
section actually changed, so saving a message template or a webhook endpoint
leaves the channels alone. NeedsRebuild compares the whole section by value, and
a test pins that reloading a file reproduces it exactly — a default applied
inconsistently would make every save look like a controller change and tear down
every channel.

With controllers immutable after construction and replaced wholesale, the atomic
pointers and Reload methods have no writers left, so they are gone and the pipes
return to plain fields. notifyReload now serializes hook execution for the same
reason: a hook mutates process-wide state, and two overlapping settings updates
would otherwise each decide to rebuild from their own view of what was applied.

Kept from that work: applyEmailProxy restores gomail's default dialer when
networkUseProxy is turned off. The old constructor only ever installed the proxy
dialer, so switching the flag back left SMTP tunnelled through a proxy with no
way to undo it short of a restart.

Documents the resulting behaviour in the README under "What applies without a
restart".
2026-09-28 23:21:35 +08:00
NanamiAdmin 0691e85ccf feat(controller): hot-reload the notification pipes
Build / ubuntu-latest (push) Failing after 3m13s
Build / windows-latest (push) Canceled after 6m23s
Email, ntfy and webhook kept their settings in a plain struct field copied at
construction, so editing controllerMethod in config.json wrote the file and
replaced the in-memory configuration while the running channels stayed on their
boot values.

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

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

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

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

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

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

A panicking hook is logged and skipped: by then the new configuration is on disk
and published, so reporting the write as failed would be a lie, and the hooks
registered after the broken one still need to run.
2026-09-28 23:10:33 +08:00
NanamiAdmin ac5587d0c8 fix(handler): read the data envelope in the server endpoint tests
Build / windows-latest (push) Failing after 1m49s
Build / ubuntu-latest (push) Failing after 5m45s
TestServerGetInfoAll, TestServerGetInfoSingleAndMissing and
TestServerGetStatusAll indexed the decoded response by uuid at the top level,
but utils.SendSuccessResponse wraps the payload as {success, message?, data:{...}}
— the shape the console reads as response.data[uuid] (frontend/src/api/index.js).
The uuid lookups therefore always missed, and unmarshalling the nil raw message
surfaced as "unexpected end of JSON input".

decodeResponse now decodes the envelope, asserts it is a success carrying a data
object, and returns that object. The handlers are unchanged: they already return
the documented shape.
2026-09-28 23:03:45 +08:00
NanamiAdminandClaude Code 2524e40566 test(config): restore settings tests with fixture IDs redacted
config/settings_test.go was removed from the whole history because its fixtures
used a real QQ number and group ID. This adds the file back as a new commit
whose fixtures are placeholder IDs, so the deep-merge, null-delete, array-replace
and settings-type validation coverage is kept without the identifiers.

The file has no history before this commit by design: the old blobs are gone
from every ref, and the working copy has been rebuilt from the redacted
version.

Co-Authored-By: Claude Code <noreply@anthropic.com>
2026-09-28 22:49:54 +08:00
NanamiAdminandClaude Code 5167d799bc chore(gitignore): narrow the test-file ignore rule to the repo root
`*test.go` matched every test file at any depth, which kept the real suite out
of version control along with the scratch files it was meant to catch. Anchor
the pattern to the repository root so ad-hoc test files dropped there stay
untracked while the tests that live next to the code they cover are tracked.

Co-Authored-By: Claude Code <noreply@anthropic.com>
2026-09-28 22:47:42 +08:00
NanamiAdminandClaude Code cb2df5076d refactor(config): publish config singletons through atomic pointers
The three configuration singletons (C_globalConfig, C_botUserConfig,
C_botNodeConfig) were plain variables: LoadGlobalConfig and friends assigned
them from the goroutine handling a settings update, while bot pipes, the node
tracker and the HTTP handlers read them from their own goroutines. That is an
unsynchronized read of a concurrently written variable — a data race the race
detector reports, and one that already existed before any hot-reload work
because /api/webhook/add reloads the configuration while the bots run.

Replace them with atomic.Pointer values behind Current(), BotUsers() and
BotNodes(). Each reload builds a fresh value and publishes it atomically, so a
reader either sees the previous configuration or the new one, never a partial
one. Callers read through the accessor on every use instead of caching it.

Two spots that read several fields of one guard now snapshot once per call, so
a reload cannot split a combined check mid-flight:
- qq_napcat.handleNapcatEvent, which evaluates the debug guards per event
- the komari task-echo guard, now behind taskEchoEnabled()

The NapCat HTTP methods keep logging on showNapcatAction alone (without
requiring debugMode), matching their existing behaviour; that inconsistency
with the WebSocket path is preserved, not introduced, and is called out in
actionLogEnabled.

Also adds TestConcurrentReloadAndRead, which drives every accessor from four
reader goroutines while two writers reload the configuration, and sanitises the
member IDs used as test fixtures in config/settings_test.go.

Co-Authored-By: Claude Code <noreply@anthropic.com>
2026-09-28 22:47:37 +08:00
NanamiAdmin 2d5cf39151 fix(changelog): update version format to include commit hash for pre-release 2026-09-24 15:32:36 +08:00
NanamiAdmin 05dcc3769d feat(changelog): add webhook API settings and configuration feature 2026-09-24 15:03:57 +08:00
NanamiAdmin f36f17bd2a feat(frontend): enhance webhook functionality and UI
- 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 3877c6d128 feat(webhook): implement management API for incoming webhook endpoints 2026-09-24 12:32:22 +08:00
NanamiAdmin 3ce42b3cca docs(changelog): add incoming webhook API feature to version 0.2.0.5 2026-09-24 11:24:13 +08:00
NanamiAdmin 057d7d584d feat: add incoming webhook API support with configurable endpoints
- 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 90ca08066b feat(release): update version to 0.2.0 and enhance changelog with new features 2026-09-23 11:22:46 +08:00
NanamiAdmin d481f13d1d feat(auth): implement WebSocket authentication for admin access to logs 2026-09-23 11:20:52 +08:00
NanamiAdmin 7ecf51eba5 docs(changelog): add Ver.0.1.2.4-59c0f5d.pre-release change log 2026-09-23 11:01:47 +08:00
NanamiAdmin 59c0f5d5d8 fix(variables): update version and build number in SoftwareInfo 2026-09-22 22:58:43 +08:00
NanamiAdmin 16455027ec fix(changelog): add missing URL for frontend build integration entry 2026-09-22 22:57:30 +08:00
NanamiAdmin 4be6d2add0 feat(build): integrate frontend build into backend binary and update scripts 2026-09-22 22:57:08 +08:00
NanamiAdmin 15144fb9d4 docs(changelog): add initial changelog entries for version 0.1.2.3 2026-09-22 22:11:04 +08:00
NanamiAdmin a89ad8143a fix(build): correct build script paths for Windows and Linux 2026-09-22 21:00:15 +08:00
NanamiAdmin fb5b46a548 chore(go.mod): update Go version to 1.27.1 2026-09-22 20:59:46 +08:00
NanamiAdmin f6092cd178 docs(readme): add frontend development instructions and update requirements 2026-09-10 21:02:43 +08:00
NanamiAdmin 3cba5df5c7 docs(readme): update build instructions 2026-09-10 21:01:13 +08:00
NanamiAdmin b2566d45f8 feat(build): add scripts for building frontend and backend for Linux and Windows 2026-09-10 20:59:21 +08:00
NanamiAdmin 8b88377d4f chore(web): directly read frontend/dist folder but not read it from web folder 2026-09-10 20:51:11 +08:00
NanamiAdmin f6af57cf58 feat(web): add static file serving for the frontend 2026-09-10 20:43:46 +08:00
NanamiAdmin 160918a93e fix(frontend): adjust the first node card height to avoid higher a little then others 2026-09-10 20:15:26 +08:00
NanamiAdmin f622baadbd fix(frontend/settings): fix json marshal fault, now the settings page could be loaded correctly 2026-09-10 20:12:21 +08:00
NanamiAdmin 78284ed816 feat(frontend): add debugMode global variable to store debug mode state 2026-09-10 19:36:29 +08:00
NanamiAdmin d96b90b5bf feat: update API response structure to nest payloads under a single "data" key 2026-09-09 16:15:33 +08:00
NanamiAdmin c0eada9bcc remove claude skills file 2026-09-08 23:18:31 +08:00
NanamiAdmin 378727ac57 feat(frontend): implement basic frontend interface 2026-09-08 23:15:50 +08:00
NanamiAdmin 8b43e8b2ea feat(config): make resolver support null value to delete a
key.
2026-09-08 23:10:26 +08:00
NanamiAdmin 786f364743 feat(user): implement registration for the first user with concurrency handling 2026-09-08 22:53:02 +08:00
NanamiAdmin a0e2df615c 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 1641 additions and 231 deletions
+4
View File
@@ -6,6 +6,10 @@ 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/
+13 -1
View File
@@ -294,6 +294,18 @@ Admins and trusted groups are defined **per bot channel** and map a member ID to
| `{{ list.offlineServers }}` | Formatted list of offline servers | | `{{ list.offlineServers }}` | Formatted list of offline servers |
| `{{ softwareVersion }}`, `{{ softwareBuildVer }}`, `{{ softwareCommitHash }}`, `{{ softwareBuildType }}`, `{{ softwareBuildTime }}`, `{{ softwareDeveloper }}`, `{{ softwareDescription }}` | Build metadata (commit hash and build time are injected at compile time) | | `{{ softwareVersion }}`, `{{ softwareBuildVer }}`, `{{ softwareCommitHash }}`, `{{ softwareBuildType }}`, `{{ softwareBuildTime }}`, `{{ softwareDeveloper }}`, `{{ softwareDescription }}` | Build metadata (commit hash and build time are injected at compile time) |
### What applies without a restart
Saving settings applies most of them immediately. How each group takes effect:
| Settings | How it applies |
|---|---|
| `controllerMethod` (all five channels) | Every channel is **rebuilt**: the running controllers are stopped and a fresh set is built from the new settings. A channel that owns a connection reconnects — Telegram re-runs its `getMe` handshake, NapCat opens a new WebSocket — so notifications sent during the swap are lost. The rebuild only happens when this section actually changed; saving a message template does not disturb the channels. |
| `controllerMessage`, `debug`, `bot_user_config.json`, `bot_node_config.json`, `webhook.endpoints` | Picked up as they are used; nothing is restarted. |
| `system.debugMode` | Applies to both behavior and log filtering. |
| `system.networkProxy` | Read on every request and every dial, so a new address reaches channels that are already running. Only the per-channel `networkUseProxy` opt-in is fixed when a channel is built, so toggling that still needs the rebuild above. |
| `system.listenAddr` / `listenPort`, `webhook.enabled` / `listenAddr` / `listenPort`, `dataPath`, `dbPath`, `komari.dashboardURL` | **Applied at startup only.** Saving them changes the file and the in-memory configuration but not the running listener, database or Komari client — restart to apply. |
## API ## API
Success responses follow the envelope `{"success": true, "message": "...", "data": {...}}`, with the payload nested under a single `data` key. Error responses use `{"success": false, "message": "..."}`. `message` may be omitted on success when there is nothing to report. Success responses follow the envelope `{"success": true, "message": "...", "data": {...}}`, with the payload nested under a single `data` key. Error responses use `{"success": false, "message": "..."}`. `message` may be omitted on success when there is nothing to report.
@@ -322,7 +334,7 @@ Browser WebSocket handshakes cannot carry custom headers, so `/api/system/getLog
| `/api/server/getStatus` | GET | admin | Live server status (mirrors the Bot's `/status`). Query `?uuid=<uuid>` (or `all`). Returns `data: {<uuid>: {uuid, name, online, report}}`; `report` is `null` when the node has not reported yet. `404` for an unknown single uuid. | | `/api/server/getStatus` | GET | admin | Live server status (mirrors the Bot's `/status`). Query `?uuid=<uuid>` (or `all`). Returns `data: {<uuid>: {uuid, name, online, report}}`; `report` is `null` when the node has not reported yet. `404` for an unknown single uuid. |
| `/api/server/exec` | POST | bot / admin | Execute a command. Body `{uuid: [<uuid>...], command}`. Dispatches a Komari task and polls until completion (or timeout). Returns `data: {taskID, results}`. | | `/api/server/exec` | POST | bot / admin | Execute a command. Body `{uuid: [<uuid>...], command}`. Dispatches a Komari task and polls until completion (or timeout). Returns `data: {taskID, results}`. |
| `/api/settings/get` | GET | admin | `?type=global\|bot_user_config\|bot_node_config` | Returns `data: {config}`, where `config` is the selected config file's content (same layout as the JSON file). | | `/api/settings/get` | GET | admin | `?type=global\|bot_user_config\|bot_node_config` | Returns `data: {config}`, where `config` is the selected config file's content (same layout as the JSON file). |
| `/api/settings/set` | POST | admin | `?type=<same types>` + JSON body of partial updates, e.g. `{"system":{"debugMode":true}}` | Deep-merges the body into the selected config file, persists it, and reloads it in memory. Only the given keys change; arrays replace. | | `/api/settings/set` | POST | admin | `?type=<same types>` + JSON body of partial updates, e.g. `{"system":{"debugMode":true}}` | Deep-merges the body into the selected config file, persists it, and reloads it in memory. Only the given keys change; arrays replace. Returns `data: {type, restartRequired}`, where `restartRequired` lists the keys the update changed that are only read at startup (see [What applies without a restart](#what-applies-without-a-restart)) — the write succeeds regardless, this only says which edits are not live yet. Always an array, empty when everything took effect. |
| `/api/webhook/add` | POST | admin | Add an incoming webhook endpoint. Body `{name, enabled?, token?, notifyPipes?}` — only the fields given are stored, the rest start at their defaults. `409` when the name is already configured. | | `/api/webhook/add` | POST | admin | Add an incoming webhook endpoint. Body `{name, enabled?, token?, notifyPipes?}` — only the fields given are stored, the rest start at their defaults. `409` when the name is already configured. |
| `/api/webhook/modify` | POST | admin | Change an existing endpoint. Body `{name, ...}` — the fields given are the fields that change (same partial-update rule as `/api/settings/set`, but scoped to one endpoint). `404` for an unknown name, `400` when no other field is given. | | `/api/webhook/modify` | POST | admin | Change an existing endpoint. Body `{name, ...}` — the fields given are the fields that change (same partial-update rule as `/api/settings/set`, but scoped to one endpoint). `404` for an unknown name, `400` when no other field is given. |
| `/api/webhook/delete` | POST | admin | Remove an endpoint. Body `{name}`. `404` for an unknown name. | | `/api/webhook/delete` | POST | admin | Remove an endpoint. Body `{name}`. `404` for an unknown name. |
+164
View File
@@ -0,0 +1,164 @@
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)
}
}
+10 -11
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) {
C_botNodeConfig = cfg botNodeConfig.Store(&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,19 +27,17 @@ 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)
} }
} }
C_botNodeConfig = cfg botNodeConfig.Store(&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. // explicitly sets the flag to false. A missing configuration (nil map) yields
// the same default.
func NodeStatusNotifyEnabled(uuid string) bool { func NodeStatusNotifyEnabled(uuid string) bool {
if C_botNodeConfig == nil { opts, ok := BotNodes()[uuid]
return true
}
opts, ok := C_botNodeConfig[uuid]
if !ok || opts.EnableStatusNotify == nil { if !ok || opts.EnableStatusNotify == nil {
return true return true
} }
@@ -146,7 +144,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 }}"
} }
C_globalConfig = &cfg globalConfig.Store(&cfg)
return &cfg, nil return &cfg, nil
} }
@@ -162,7 +160,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)
} }
C_botUserConfig = &cfg botUserConfig.Store(&cfg)
return &cfg, nil return &cfg, nil
} }
@@ -213,8 +211,9 @@ 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 {
if C_globalConfig == nil { cfg := Current()
if cfg == nil {
return false return false
} }
return C_globalConfig.System.DebugMode return cfg.System.DebugMode
} }
+81
View File
@@ -0,0 +1,81 @@
package config
import (
"fmt"
"sync"
"nukumizu-backend/postLog"
)
// reloadHooks are the callbacks run after the global configuration has been
// reloaded, i.e. after every settings update that touches config.json.
//
// They exist so that packages which already depend on config — the controller
// manager, the logger — can react to an update without config importing them,
// which would be an import cycle. Everything a hook needs is passed in.
var (
reloadHooksMu sync.Mutex
reloadHooks []func(*Config)
// reloadRunMu serializes hook execution. A hook mutates process-wide state
// — the logger's debug flag, the controller registry — so two overlapping
// settings updates must not run them at the same time, or both would decide
// to rebuild the controllers from their own view of what was applied.
reloadRunMu sync.Mutex
)
// OnReload registers a hook to run after every reload of the global
// configuration, receiving the configuration now in effect. Hooks run in
// registration order.
//
// Register once, at startup, before the first settings update can arrive: a
// hook registered later has already missed the updates that came before it, and
// the configuration it would have seen is not replayed.
//
// A hook runs on the goroutine serving /api/settings/set, so it must not block
// for long. It may be called concurrently by two overlapping updates.
func OnReload(hook func(*Config)) {
reloadHooksMu.Lock()
defer reloadHooksMu.Unlock()
reloadHooks = append(reloadHooks, hook)
}
// notifyReload runs every registered hook with cfg. A nil cfg is the signal
// that the update touched one of the other settings files, which have no
// hook-visible reload, and it is ignored.
//
// A panicking hook is logged and skipped rather than allowed to unwind through
// UpdateSettings: by the time hooks run the new configuration is already on
// disk and published in memory, so reporting the write as failed would be a
// lie, and the hooks registered after the broken one must still run.
func notifyReload(cfg *Config) {
if cfg == nil {
return
}
// Copy under the registry lock, then release it before running anything: a
// hook is free to register another hook without deadlocking.
reloadHooksMu.Lock()
hooks := make([]func(*Config), len(reloadHooks))
copy(hooks, reloadHooks)
reloadHooksMu.Unlock()
reloadRunMu.Lock()
defer reloadRunMu.Unlock()
for _, hook := range hooks {
runReloadHook(hook, cfg)
}
}
// runReloadHook runs one hook, isolating a panic to that hook. Recovering in a
// separate function rather than inline is deliberate: a deferred recover placed
// in the loop body would not run until notifyReload itself returned, which
// would abandon the remaining hooks.
func runReloadHook(hook func(*Config), cfg *Config) {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Configuration reload hook panicked: %v", r))
}
}()
hook(cfg)
}
+173
View File
@@ -0,0 +1,173 @@
package config
import (
"reflect"
"testing"
"nukumizu-backend/global"
)
// swapReloadHooks replaces the registered reload hooks for the duration of a
// test and restores the previous set afterwards, so the package-level registry
// cannot leak into another test.
func swapReloadHooks(t *testing.T, hooks ...func(*Config)) {
t.Helper()
reloadHooksMu.Lock()
original := reloadHooks
reloadHooks = hooks
reloadHooksMu.Unlock()
t.Cleanup(func() {
reloadHooksMu.Lock()
reloadHooks = original
reloadHooksMu.Unlock()
})
}
func TestReloadHookRunsOnlyForGlobalSettings(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{"system":{"debugMode":false}}`)
writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{}`)
writeTempConfig(t, &global.ConfigPath.BotNodeConfig, `{}`)
var seen []bool
swapReloadHooks(t, func(cfg *Config) {
seen = append(seen, cfg.System.DebugMode)
})
// A config.json update runs the hook, with the values just written.
if err := UpdateSettings(SettingGlobal, map[string]interface{}{
"system": map[string]interface{}{"debugMode": true},
}); err != nil {
t.Fatalf("UpdateSettings(global): %v", err)
}
if len(seen) != 1 {
t.Fatalf("hook ran %d times for a config.json update, want 1", len(seen))
}
if !seen[0] {
t.Error("hook received a configuration without the updated debugMode")
}
// The other two files reload a singleton that callers read at the point of
// use, so they have nothing to notify.
for _, settingsType := range []string{SettingBotUserConfig, SettingBotNodeConfig} {
if err := UpdateSettings(settingsType, map[string]interface{}{
"unused": map[string]interface{}{"event_reply": true},
}); err != nil {
t.Fatalf("UpdateSettings(%s): %v", settingsType, err)
}
}
if len(seen) != 1 {
t.Errorf("hook ran %d times after updates to the other settings files, want 1", len(seen))
}
}
func TestReloadHookIsolatesPanic(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{"system":{"debugMode":false}}`)
reached := false
swapReloadHooks(t,
func(*Config) { panic("hook under test") },
func(*Config) { reached = true },
)
// The configuration is already on disk and published by the time hooks run,
// so a broken hook must not turn a successful write into a failed request.
if err := UpdateSettings(SettingGlobal, map[string]interface{}{
"system": map[string]interface{}{"debugMode": true},
}); err != nil {
t.Fatalf("a panicking hook must not fail the settings write: %v", err)
}
if !reached {
t.Error("a panicking hook stopped the hooks registered after it")
}
// The write itself must still have landed.
cfg := Current()
if cfg == nil || !cfg.System.DebugMode {
t.Error("the settings write did not take effect")
}
}
// TestUnrelatedUpdateLeavesControllerMethodAlone guards the trigger for a
// controller rebuild. Whether to rebuild is decided by comparing the whole
// controllerMethod section with the one the running controllers were built
// from, so reloading the file has to reproduce that section byte for byte. A
// default applied inconsistently — a nil recipient slice turned into an empty
// one on the second load, say — would make every settings edit look like a
// controller change and tear down every channel on each save.
func TestUnrelatedUpdateLeavesControllerMethodAlone(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{
"controllerMethod": {
"qq(napcat)": { "enabled": false },
"email": { "enabled": false, "to": [] },
"webhook": { "enabled": false, "headers": {} }
}
}`)
first, err := LoadGlobalConfig(global.ConfigPath.Global)
if err != nil {
t.Fatalf("LoadGlobalConfig: %v", err)
}
before := first.ControllerMethod
if err := UpdateSettings(SettingGlobal, map[string]interface{}{
"controllerMessage": map[string]interface{}{"BOT_STARTED": "hello"},
}); err != nil {
t.Fatalf("UpdateSettings: %v", err)
}
after := Current().ControllerMethod
if !reflect.DeepEqual(before, after) {
t.Errorf("an unrelated update changed the controllerMethod section:\nbefore: %+v\nafter: %+v", before, after)
}
}
// TestReloadHookSeesEveryWriterPath pins the invariant that each writer of
// config.json notifies, not just /api/settings/set: the webhook endpoint
// helpers go through the same channel.
func TestReloadHookSeesEveryWriterPath(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{
"webhook": { "endpoints": {} }
}`)
var seen int
swapReloadHooks(t, func(*Config) { seen++ })
if err := AddWebhookEndpoint("example", map[string]interface{}{
"enabled": true,
"token": "secret",
"notifyPipes": []interface{}{"ntfy"},
}); err != nil {
t.Fatalf("AddWebhookEndpoint: %v", err)
}
if seen != 1 {
t.Errorf("hook ran %d times after adding an endpoint, want 1", seen)
}
if err := ModifyWebhookEndpoint("example", map[string]interface{}{
"enabled": false,
}); err != nil {
t.Fatalf("ModifyWebhookEndpoint: %v", err)
}
if seen != 2 {
t.Errorf("hook ran %d times after modifying an endpoint, want 2", seen)
}
if err := DeleteWebhookEndpoint("example"); err != nil {
t.Fatalf("DeleteWebhookEndpoint: %v", err)
}
if seen != 3 {
t.Errorf("hook ran %d times after deleting an endpoint, want 3", seen)
}
// A rejected write changes nothing, so it must not notify either.
if err := DeleteWebhookEndpoint("ghost"); err == nil {
t.Error("deleting an unknown endpoint should fail")
}
if seen != 3 {
t.Errorf("hook ran for a rejected write (%d notifications, want 3)", seen)
}
}
+125 -14
View File
@@ -6,6 +6,7 @@ import (
"errors" "errors"
"fmt" "fmt"
"os" "os"
"strings"
"sync" "sync"
"nukumizu-backend/global" "nukumizu-backend/global"
@@ -23,6 +24,89 @@ const (
// accepted constants above. // accepted constants above.
var ErrUnsupportedSettingsType = errors.New("unsupported settings type") var ErrUnsupportedSettingsType = errors.New("unsupported settings type")
// startupOnlySettings are the config.json keys that are read once before the
// program starts serving and never again: the listener addresses, the paths the
// databases are opened from, and the Komari dashboard its client is built
// against. Editing one writes the file and replaces the in-memory
// configuration, but the running program keeps the old value, so an update that
// touches one is reported back to the caller instead of being silently
// accepted.
//
// Keep this in step with main: these are exactly the settings main reads before
// the HTTP server comes up. Everything else — controllerMethod, networkProxy,
// the message templates, the debug switches — is picked up at runtime.
var startupOnlySettings = []string{
"system.listenAddr",
"system.listenPort",
"webhook.enabled",
"webhook.listenAddr",
"webhook.listenPort",
"komari.dashboardURL",
"dataPath",
"dbPath",
}
// RestartRequiredKeys lists the settings in patch that only take effect at
// startup, as dot-separated paths, in the order startupOnlySettings declares
// them. Only config.json carries such settings; an update to one of the other
// files always reports nothing.
//
// The write itself succeeds either way — this is advice for the user, not a
// rejection. The result is never nil, so a caller can put it straight into a
// JSON response and get [] rather than null.
func RestartRequiredKeys(settingsType string, patch map[string]interface{}) []string {
keys := []string{}
if settingsType != SettingGlobal {
return keys
}
patched := patchPaths(patch)
for _, watched := range startupOnlySettings {
for _, path := range patched {
if pathsOverlap(path, watched) {
keys = append(keys, watched)
break
}
}
}
return keys
}
// patchPaths expands a nested settings patch into the dot-separated paths of its
// leaves. An object is descended into rather than reported, so a patch that only
// names sections still resolves to the keys it changes, and a JSON null is a
// leaf because it deletes the key it names.
func patchPaths(patch map[string]interface{}) []string {
paths := []string{}
var walk func(prefix string, node map[string]interface{})
walk = func(prefix string, node map[string]interface{}) {
for key, value := range node {
path := key
if prefix != "" {
path = prefix + "." + key
}
if nested, ok := value.(map[string]interface{}); ok && nested != nil {
walk(path, nested)
continue
}
paths = append(paths, path)
}
}
walk("", patch)
return paths
}
// pathsOverlap reports whether a patched path and a watched setting can affect
// each other: they are the same key, the patch names something inside the
// watched setting, or the patch names a section the watched setting lives in.
// The last case matters because a patch may replace a whole section, which
// changes every key under it.
func pathsOverlap(patched, watched string) bool {
return patched == watched ||
strings.HasPrefix(patched, watched+".") ||
strings.HasPrefix(watched, patched+".")
}
// settingsLock serializes read-modify-write access to the on-disk configuration // settingsLock serializes read-modify-write access to the on-disk configuration
// files so concurrent admin edits (UpdateSettings) and the node tracker's // files so concurrent admin edits (UpdateSettings) and the node tracker's
// background save (SaveBotNodeConfig) cannot lose each other's updates. // background save (SaveBotNodeConfig) cannot lose each other's updates.
@@ -82,19 +166,45 @@ func GetSettings(settingsType string) ([]byte, error) {
// written the matching in-memory singleton is reloaded so runtime code observes // written the matching in-memory singleton is reloaded so runtime code observes
// the new values. // the new values.
func UpdateSettings(settingsType string, patch map[string]interface{}) error { func UpdateSettings(settingsType string, patch map[string]interface{}) error {
return runSettingsUpdate(func() (*Config, error) {
return updateSettingsLocked(settingsType, patch)
})
}
// runSettingsUpdate runs fn under settingsLock and then, once the lock is
// released, runs the reload hooks with whatever configuration fn reports (nil
// when the update did not touch config.json).
//
// The hooks deliberately run outside settingsLock. A hook rebuilds controllers,
// which can wait on a network call, while settingsLock is also held by the node
// tracker's background save (SaveBotNodeConfig); holding it across a hook would
// stall node registration behind an unrelated settings edit.
func runSettingsUpdate(fn func() (*Config, error)) error {
cfg, err := func() (*Config, error) {
settingsLock.Lock() settingsLock.Lock()
defer settingsLock.Unlock() defer settingsLock.Unlock()
return fn()
}()
if err != nil {
return err
}
return updateSettingsLocked(settingsType, patch) notifyReload(cfg)
return nil
} }
// updateSettingsLocked is UpdateSettings without the locking, for callers that // updateSettingsLocked is UpdateSettings without the locking, for callers that
// need to inspect the loaded configuration and write in one critical section // need to inspect the loaded configuration and write in one critical section
// (see the incoming webhook endpoint helpers). Callers must hold settingsLock. // (see the incoming webhook endpoint helpers). Callers must hold settingsLock.
func updateSettingsLocked(settingsType string, patch map[string]interface{}) error { //
// It returns the freshly loaded global configuration, or nil when the settings
// type is one of the other files. The caller is responsible for handing that
// value to notifyReload once settingsLock is released — which runSettingsUpdate
// does for every writer.
func updateSettingsLocked(settingsType string, patch map[string]interface{}) (*Config, error) {
path, err := settingsPath(settingsType) path, err := settingsPath(settingsType)
if err != nil { if err != nil {
return err return nil, err
} }
// Start from whatever is already on disk so nothing is dropped. A missing or // Start from whatever is already on disk so nothing is dropped. A missing or
@@ -104,23 +214,23 @@ func updateSettingsLocked(settingsType string, patch map[string]interface{}) err
if err == nil { if err == nil {
if len(bytes.TrimSpace(data)) > 0 { if len(bytes.TrimSpace(data)) > 0 {
if err := json.Unmarshal(data, &current); err != nil { if err := json.Unmarshal(data, &current); err != nil {
return fmt.Errorf("failed to parse existing %s settings file %s: %w", settingsType, path, err) return nil, fmt.Errorf("failed to parse existing %s settings file %s: %w", settingsType, path, err)
} }
} }
} else if !os.IsNotExist(err) { } else if !os.IsNotExist(err) {
return fmt.Errorf("failed to read existing %s settings file %s: %w", settingsType, path, err) return nil, fmt.Errorf("failed to read existing %s settings file %s: %w", settingsType, path, err)
} }
deepMergeSettings(current, patch) deepMergeSettings(current, patch)
data, err = json.MarshalIndent(current, "", " ") data, err = json.MarshalIndent(current, "", " ")
if err != nil { if err != nil {
return fmt.Errorf("failed to marshal %s settings: %w", settingsType, err) return nil, fmt.Errorf("failed to marshal %s settings: %w", settingsType, err)
} }
data = append(data, '\n') data = append(data, '\n')
if err := os.WriteFile(path, data, 0o644); err != nil { if err := os.WriteFile(path, data, 0o644); err != nil {
return fmt.Errorf("failed to write %s settings file %s: %w", settingsType, path, err) return nil, fmt.Errorf("failed to write %s settings file %s: %w", settingsType, path, err)
} }
return reloadSettings(settingsType, path) return reloadSettings(settingsType, path)
@@ -151,18 +261,19 @@ func deepMergeSettings(dst, src map[string]interface{}) {
} }
// reloadSettings refreshes the in-memory singleton for the given settings type // reloadSettings refreshes the in-memory singleton for the given settings type
// so the running program observes the values just persisted to disk. // so the running program observes the values just persisted to disk. Only
func reloadSettings(settingsType, path string) error { // config.json has a hook-visible reload, so for SettingGlobal it returns the
// configuration now in effect and for the other types it returns nil.
func reloadSettings(settingsType, path string) (*Config, error) {
switch settingsType { switch settingsType {
case SettingGlobal: case SettingGlobal:
_, err := LoadGlobalConfig(path) return LoadGlobalConfig(path)
return err
case SettingBotUserConfig: case SettingBotUserConfig:
_, err := LoadBotUserConfig(path) _, err := LoadBotUserConfig(path)
return err return nil, err
case SettingBotNodeConfig: case SettingBotNodeConfig:
return LoadBotNodeConfig(path) return nil, LoadBotNodeConfig(path)
default: default:
return ErrUnsupportedSettingsType return nil, ErrUnsupportedSettingsType
} }
} }
+172 -14
View File
@@ -4,6 +4,7 @@ import (
"encoding/json" "encoding/json"
"os" "os"
"path/filepath" "path/filepath"
"reflect"
"strings" "strings"
"testing" "testing"
@@ -27,13 +28,13 @@ func TestUpdateSettingsDeepMerge(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{ writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{
"qq(napcat)": { "qq(napcat)": {
"admins": { "admins": {
"3526453517": { "100000001": {
"event_status_notify": true, "event_status_notify": true,
"event_bot_started": true "event_bot_started": true
} }
}, },
"trustedGroups": { "trustedGroups": {
"740724778": { "200000002": {
"event_status_notify": true, "event_status_notify": true,
"event_bot_started": false "event_bot_started": false
} }
@@ -44,7 +45,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{}{
"3526453517": map[string]interface{}{ "100000001": 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
}, },
@@ -70,7 +71,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 740724778 kept `"event_status_notify": true`, // sibling under trustedGroups 200000002 kept
`"12345"`, `"12345"`,
} { } {
if !strings.Contains(got, want) { if !strings.Contains(got, want) {
@@ -79,16 +80,17 @@ 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.
if C_botUserConfig == nil { users := BotUsers()
t.Fatal("C_botUserConfig not reloaded") if users == nil {
t.Fatal("bot user config not reloaded")
} }
if C_botUserConfig.QQ.Admins["3526453517"].EventStatusNotify { if users.QQ.Admins["100000001"].EventStatusNotify {
t.Error("expected reloaded admin event_status_notify = false") t.Error("expected reloaded admin event_status_notify = false")
} }
if !C_botUserConfig.QQ.Admins["3526453517"].EventReply { if !users.QQ.Admins["100000001"].EventReply {
t.Error("expected reloaded admin event_reply = true") t.Error("expected reloaded admin event_reply = true")
} }
if !C_botUserConfig.QQ.TrustedGroups["12345"].EventBotStarted { if !users.QQ.TrustedGroups["12345"].EventBotStarted {
t.Error("expected new trusted group event_bot_started = true") t.Error("expected new trusted group event_bot_started = true")
} }
} }
@@ -147,6 +149,162 @@ func TestUpdateSettingsReplacesArraysAndKeepsNumbers(t *testing.T) {
} }
} }
func TestRestartRequiredKeys(t *testing.T) {
cases := []struct {
name string
settingsType string
patch map[string]interface{}
want []string
}{
{
name: "a runtime switch needs no restart",
settingsType: SettingGlobal,
patch: map[string]interface{}{"system": map[string]interface{}{"debugMode": true}},
want: []string{},
},
{
name: "the listen port does",
settingsType: SettingGlobal,
patch: map[string]interface{}{"system": map[string]interface{}{"listenPort": "9090"}},
want: []string{"system.listenPort"},
},
{
name: "only the startup key of a mixed patch is reported",
settingsType: SettingGlobal,
patch: map[string]interface{}{
"system": map[string]interface{}{"listenPort": "9090", "debugMode": true},
},
want: []string{"system.listenPort"},
},
{
name: "a top-level path is reported",
settingsType: SettingGlobal,
patch: map[string]interface{}{"dataPath": "/srv/data"},
want: []string{"dataPath"},
},
{
name: "results follow the declared order, not the patch order",
settingsType: SettingGlobal,
patch: map[string]interface{}{"dbPath": "/srv/db", "dataPath": "/srv/data"},
want: []string{"dataPath", "dbPath"},
},
{
name: "deleting a startup key with null is reported",
settingsType: SettingGlobal,
patch: map[string]interface{}{"system": map[string]interface{}{"listenPort": nil}},
want: []string{"system.listenPort"},
},
{
name: "replacing a whole section reports the startup keys inside it",
settingsType: SettingGlobal,
patch: map[string]interface{}{"webhook": map[string]interface{}{"listenAddr": "127.0.0.1"}},
want: []string{"webhook.listenAddr"},
},
{
// Deleting the section resets the URL to its built-in default.
name: "deleting a section the startup key lives in reports it",
settingsType: SettingGlobal,
patch: map[string]interface{}{"komari": nil},
want: []string{"komari.dashboardURL"},
},
{
name: "deleting a section reports every startup key inside it",
settingsType: SettingGlobal,
patch: map[string]interface{}{"webhook": nil},
want: []string{"webhook.enabled", "webhook.listenAddr", "webhook.listenPort"},
},
{
// An empty object merges nothing, so it changes no key and needs no
// restart — surprising enough to pin.
name: "an empty object changes nothing",
settingsType: SettingGlobal,
patch: map[string]interface{}{"komari": map[string]interface{}{}},
want: []string{},
},
{
// webhook.endpoints must not be mistaken for webhook.enabled.
name: "a sibling subtree is not mistaken for the startup key",
settingsType: SettingGlobal,
patch: map[string]interface{}{
"webhook": map[string]interface{}{
"endpoints": map[string]interface{}{"example": map[string]interface{}{"enabled": true}},
},
},
want: []string{},
},
{
// The Komari credentials are re-read on the next login, so only the
// dashboard URL is startup-only.
name: "komari credentials are not startup-only",
settingsType: SettingGlobal,
patch: map[string]interface{}{
"komari": map[string]interface{}{
"account": map[string]interface{}{"username": "admin", "password": "x"},
},
},
want: []string{},
},
{
name: "controller settings are not startup-only",
settingsType: SettingGlobal,
patch: map[string]interface{}{
"controllerMethod": map[string]interface{}{
"telegram": map[string]interface{}{"enabled": true, "botToken": "t"},
},
},
want: []string{},
},
{
name: "an empty patch reports nothing",
settingsType: SettingGlobal,
patch: map[string]interface{}{},
want: []string{},
},
{
// Only config.json has settings that are read once at startup.
name: "the other settings files never need a restart",
settingsType: SettingBotUserConfig,
patch: map[string]interface{}{"dataPath": "/srv/data"},
want: []string{},
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
got := RestartRequiredKeys(tc.settingsType, tc.patch)
if !reflect.DeepEqual(got, tc.want) {
t.Errorf("RestartRequiredKeys() = %v, want %v", got, tc.want)
}
if got == nil {
t.Error("the result must never be nil, so it serializes as [] rather than null")
}
})
}
}
// TestRestartRequiredKeysIsAdvisory pins that reporting a startup-only key does
// not stop the write: the caller is told, the file is still updated.
func TestRestartRequiredKeysIsAdvisory(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.Global, `{"system":{"listenPort":"8080"}}`)
keys := RestartRequiredKeys(SettingGlobal, map[string]interface{}{
"system": map[string]interface{}{"listenPort": "9090"},
})
if len(keys) != 1 {
t.Fatalf("expected the listen port to be reported, got %v", keys)
}
if err := UpdateSettings(SettingGlobal, map[string]interface{}{
"system": map[string]interface{}{"listenPort": "9090"},
}); err != nil {
t.Fatalf("UpdateSettings: %v", err)
}
if got := Current().System.ListenPort; got != "9090" {
t.Errorf("the update was not applied: listenPort = %q", got)
}
}
func TestSettingsTypeValidation(t *testing.T) { func TestSettingsTypeValidation(t *testing.T) {
for _, valid := range []string{SettingGlobal, SettingBotUserConfig, SettingBotNodeConfig} { for _, valid := range []string{SettingGlobal, SettingBotUserConfig, SettingBotNodeConfig} {
if !IsValidSettingsType(valid) { if !IsValidSettingsType(valid) {
@@ -181,8 +339,8 @@ func TestUpdateSettingsRemovesKeysWithNull(t *testing.T) {
writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{ writeTempConfig(t, &global.ConfigPath.BotUserConfig, `{
"qq(napcat)": { "qq(napcat)": {
"admins": { "admins": {
"3526453517": { "event_status_notify": true }, "100000001": { "event_status_notify": true },
"740724778": { "event_status_notify": false } "200000002": { "event_status_notify": false }
}, },
"trustedGroups": { "trustedGroups": {
"999": { "event_bot_started": true } "999": { "event_bot_started": true }
@@ -197,7 +355,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{}{
"3526453517": nil, "100000001": nil,
}, },
}, },
} }
@@ -210,10 +368,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, "3526453517") { if strings.Contains(got, "100000001") {
t.Errorf("deleted member still present:\n%s", got) t.Errorf("deleted member still present:\n%s", got)
} }
for _, want := range []string{"740724778", `"trustedGroups"`, `"telegram"`} { for _, want := range []string{"200000002", `"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)
} }
+45 -8
View File
@@ -1,6 +1,9 @@
package config package config
import "sort" import (
"sort"
"sync/atomic"
)
// SystemConfig holds system-level configuration. // SystemConfig holds system-level configuration.
type SystemConfig struct { type SystemConfig struct {
@@ -136,10 +139,11 @@ 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) {
if C_globalConfig == nil { cfg := Current()
if cfg == nil {
return WebhookEndpointConfig{}, false return WebhookEndpointConfig{}, false
} }
endpoint, ok := C_globalConfig.Webhook.Endpoints[name] endpoint, ok := cfg.Webhook.Endpoints[name]
return endpoint, ok return endpoint, ok
} }
@@ -165,7 +169,19 @@ type Config struct {
DBPath string `json:"dbPath"` DBPath string `json:"dbPath"`
} }
var C_globalConfig *Config // globalConfig holds the configuration currently in effect. It is replaced as a
// 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 {
@@ -220,7 +236,16 @@ type BotUserConfig struct {
Telegram BotUser_TelegramConfig `json:"telegram"` Telegram BotUser_TelegramConfig `json:"telegram"`
} }
var C_botUserConfig *BotUserConfig // botUserConfig mirrors bot_user_config.json the same way globalConfig mirrors
// 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;
@@ -238,6 +263,18 @@ type BotNodeOptions struct {
// options. // options.
type BotNodeMembers map[string]BotNodeOptions type BotNodeMembers map[string]BotNodeOptions
// C_botNodeConfig is the global singleton mirroring bot_node_config.json, // botNodeConfig mirrors bot_node_config.json, populated by LoadBotNodeConfig.
// populated by LoadBotNodeConfig. // The map is rebuilt rather than mutated on every load, so the pointer can be
var C_botNodeConfig BotNodeMembers // swapped atomically; read it through BotNodes.
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
}
+18 -20
View File
@@ -37,10 +37,11 @@ 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{}
if C_globalConfig == nil { cfg := Current()
if cfg == nil {
return endpoints return endpoints
} }
for name, endpoint := range C_globalConfig.Webhook.Endpoints { for name, endpoint := range cfg.Webhook.Endpoints {
endpoints[name] = endpoint endpoints[name] = endpoint
} }
return endpoints return endpoints
@@ -59,13 +60,12 @@ func AddWebhookEndpoint(name string, fields map[string]interface{}) error {
return err return err
} }
settingsLock.Lock() return runSettingsUpdate(func() (*Config, error) {
defer settingsLock.Unlock()
if _, exists := webhookEndpoint(name); exists { if _, exists := webhookEndpoint(name); exists {
return fmt.Errorf("%w: %s", ErrWebhookEndpointExists, name) return nil, fmt.Errorf("%w: %s", ErrWebhookEndpointExists, name)
} }
return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch)) return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch))
})
} }
// ModifyWebhookEndpoint updates an existing incoming webhook endpoint. Only the // ModifyWebhookEndpoint updates an existing incoming webhook endpoint. Only the
@@ -83,38 +83,36 @@ func ModifyWebhookEndpoint(name string, fields map[string]interface{}) error {
return fmt.Errorf("%w: no fields to update", ErrWebhookEndpointInvalid) return fmt.Errorf("%w: no fields to update", ErrWebhookEndpointInvalid)
} }
settingsLock.Lock() return runSettingsUpdate(func() (*Config, error) {
defer settingsLock.Unlock()
if _, exists := webhookEndpoint(name); !exists { if _, exists := webhookEndpoint(name); !exists {
return fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name) return nil, fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name)
} }
return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch)) return updateSettingsLocked(SettingGlobal, webhookEndpointsPatch(name, patch))
})
} }
// DeleteWebhookEndpoint removes the incoming webhook endpoint registered under // DeleteWebhookEndpoint removes the incoming webhook endpoint registered under
// name. The endpoint stops accepting requests as soon as the configuration is // name. The endpoint stops accepting requests as soon as the configuration is
// reloaded. // reloaded.
func DeleteWebhookEndpoint(name string) error { func DeleteWebhookEndpoint(name string) error {
settingsLock.Lock() return runSettingsUpdate(func() (*Config, error) {
defer settingsLock.Unlock()
if _, exists := webhookEndpoint(name); !exists { if _, exists := webhookEndpoint(name); !exists {
return fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name) return nil, fmt.Errorf("%w: %s", ErrWebhookEndpointNotFound, name)
} }
return updateSettingsLocked(SettingGlobal, webhookEndpointDeletePatch(name)) return updateSettingsLocked(SettingGlobal, webhookEndpointDeletePatch(name))
})
} }
// webhookEndpoint returns the named endpoint held by the loaded configuration. // webhookEndpoint returns the named endpoint held by the loaded configuration.
// No lock is needed to read it: a reload replaces the whole configuration // No extra lock is needed to read it: a reload replaces the whole configuration
// rather than mutating it in place, and the value is read from whichever // rather than mutating it in place, and Current publishes the replacement
// version is current. // atomically, so the value read is always from one complete version.
func webhookEndpoint(name string) (WebhookEndpointConfig, bool) { func webhookEndpoint(name string) (WebhookEndpointConfig, bool) {
if C_globalConfig == nil { cfg := Current()
if cfg == nil {
return WebhookEndpointConfig{}, false return WebhookEndpointConfig{}, false
} }
endpoint, exists := C_globalConfig.Webhook.Endpoints[name] endpoint, exists := cfg.Webhook.Endpoints[name]
return endpoint, exists return endpoint, exists
} }
+4
View File
@@ -13,6 +13,10 @@ export const serverApi = {
// /api/settings/get?type=… / /api/settings/set?type=… // /api/settings/get?type=… / /api/settings/set?type=…
// get → { success, message, data: { config } }. // get → { success, message, data: { config } }.
// set → { success, message, data: { type, restartRequired } }, where
// restartRequired lists the keys the update changed that are only read at
// startup, so the caller can say which edits are not live yet. It is
// always an array, empty when the whole update took effect.
// `type` is one of global | bot_user_config | bot_node_config. // `type` is one of global | bot_user_config | bot_node_config.
// For set, pass a partial object; a JSON null value removes that key. // For set, pass a partial object; a JSON null value removes that key.
export const settingsApi = { export const settingsApi = {
+8 -1
View File
@@ -133,8 +133,15 @@ async function save() {
const patch = wrapRoot(props.section, nest(obj)); const patch = wrapRoot(props.section, nest(obj));
saving.value = true; saving.value = true;
try { try {
await settingsApi.set('global', patch); const res = await settingsApi.set('global', patch);
// The backend reports the keys it wrote that are only read at startup.
// A plain "saved" would suggest those are live too.
const pending = (res && res.data && res.data.restartRequired) || [];
if (pending.length) {
toast.warn(`${props.section.title} saved — restart to apply: ${pending.join(', ')}`, 7000);
} else {
toast.success(`${props.section.title} saved`); toast.success(`${props.section.title} saved`);
}
emit('saved'); emit('saved');
} catch (e) { } catch (e) {
toast.error('Failed to save: ' + e.message); toast.error('Failed to save: ' + e.message);
+2 -2
View File
@@ -11,7 +11,7 @@ const sections = [
{ {
id: 'system', id: 'system',
title: 'System', title: 'System',
hint: 'HTTP listener and global runtime switches.', hint: 'HTTP listener and global runtime switches. A changed listen address or port applies on restart.',
root: ['system'], root: ['system'],
fields: [ fields: [
{ key: 'debugMode', type: 'bool', label: 'Debug mode', help: 'Skipped X-Timestamp checks and verbose debug logging.' }, { key: 'debugMode', type: 'bool', label: 'Debug mode', help: 'Skipped X-Timestamp checks and verbose debug logging.' },
@@ -37,7 +37,7 @@ const sections = [
{ {
id: 'komari', id: 'komari',
title: 'Komari dashboard', title: 'Komari dashboard',
hint: 'Connection the monitor reads node data from. Takes effect on restart.', hint: 'Connection the monitor reads node data from. The URL applies on restart; the account is re-read on the next login.',
root: ['komari'], root: ['komari'],
fields: [ fields: [
{ key: 'dashboardURL', type: 'text', label: 'Dashboard URL' }, { key: 'dashboardURL', type: 'text', label: 'Dashboard URL' },
+16 -2
View File
@@ -51,13 +51,27 @@ 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 map[string]json.RawMessage var body struct {
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())
} }
return body if !body.Success {
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) {
+7
View File
@@ -47,6 +47,10 @@ func SettingsGetHandler(w http.ResponseWriter, r *http.Request) {
// {"system": {"debugMode": true}} // {"system": {"debugMode": true}}
// //
// Multiple entries may be given at once; only the provided keys are changed. // Multiple entries may be given at once; only the provided keys are changed.
//
// The response carries data.restartRequired: the keys the update changed that
// are only read at startup, so the caller can say which edits are not live yet.
// It is empty for an update that took effect in full.
func SettingsSetHandler(w http.ResponseWriter, r *http.Request) { func SettingsSetHandler(w http.ResponseWriter, r *http.Request) {
if !utils.Auth(w, r, "POST", "admin") { if !utils.Auth(w, r, "POST", "admin") {
return return
@@ -80,7 +84,10 @@ func SettingsSetHandler(w http.ResponseWriter, r *http.Request) {
return return
} }
// The write landed either way. This only tells the caller which of the keys
// it changed will not be live until the program is restarted.
utils.SendSuccessResponse(w, "settings updated successfully", map[string]interface{}{ utils.SendSuccessResponse(w, "settings updated successfully", map[string]interface{}{
"type": settingsType, "type": settingsType,
"restartRequired": config.RestartRequiredKeys(settingsType, patch),
}) })
} }
+128
View File
@@ -0,0 +1,128 @@
package handler
import (
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strconv"
"strings"
"testing"
"time"
"nukumizu-backend/config"
"nukumizu-backend/global"
)
// tempConfigFile points one of the settings paths at a throwaway file, so a
// test can drive the settings API without touching the config files in the
// working directory.
func tempConfigFile(t *testing.T, field *string, content string) {
t.Helper()
path := filepath.Join(t.TempDir(), "config.json")
if err := os.WriteFile(path, []byte(content), 0o644); err != nil {
t.Fatalf("write temp config: %v", err)
}
original := *field
*field = path
t.Cleanup(func() { *field = original })
}
// settingsSetRequest builds an authenticated POST for the settings endpoint.
func settingsSetRequest(t *testing.T, settingsType, body string) *http.Request {
t.Helper()
req := httptest.NewRequest(
http.MethodPost,
"/api/settings/set?type="+settingsType,
strings.NewReader(body),
)
req.Header.Set("X-Token", "test-admin-token")
req.Header.Set("X-Timestamp", strconv.FormatInt(time.Now().Unix(), 10))
return req
}
// settingsSetResponse is the envelope /api/settings/set answers with.
type settingsSetResponse struct {
Success bool `json:"success"`
Data struct {
Type string `json:"type"`
RestartRequired []string `json:"restartRequired"`
} `json:"data"`
}
func decodeSettingsSetResponse(t *testing.T, w *httptest.ResponseRecorder) settingsSetResponse {
t.Helper()
var body settingsSetResponse
if err := json.Unmarshal(w.Body.Bytes(), &body); err != nil {
t.Fatalf("decode response: %v; body=%s", err, w.Body.String())
}
return body
}
// TestSettingsSetReportsStartupOnlyKeys covers the field the console reads to
// tell the user which of their edits are not live yet.
func TestSettingsSetReportsStartupOnlyKeys(t *testing.T) {
setupAdminToken()
tempConfigFile(t, &global.ConfigPath.Global, `{"system":{"debugMode":false,"listenPort":"8080"}}`)
w := httptest.NewRecorder()
SettingsSetHandler(w, settingsSetRequest(t, "global",
`{"system":{"listenPort":"9090","debugMode":true}}`))
if w.Code != http.StatusOK {
t.Fatalf("status = %d, body = %s", w.Code, w.Body.String())
}
body := decodeSettingsSetResponse(t, w)
if !body.Success {
t.Fatalf("not a success envelope: %s", w.Body.String())
}
if body.Data.Type != "global" {
t.Errorf("data.type = %q, want global", body.Data.Type)
}
if len(body.Data.RestartRequired) != 1 || body.Data.RestartRequired[0] != "system.listenPort" {
t.Errorf("restartRequired = %v, want [system.listenPort]", body.Data.RestartRequired)
}
// Being reported as startup-only must not stop the write.
data, err := config.GetSettings(config.SettingGlobal)
if err != nil {
t.Fatalf("GetSettings: %v", err)
}
for _, want := range []string{`"listenPort": "9090"`, `"debugMode": true`} {
if !strings.Contains(string(data), want) {
t.Errorf("the patch was not applied, missing %s:\n%s", want, data)
}
}
}
// TestSettingsSetRestartRequiredIsAlwaysAnArray pins the shape a client
// iterates over: an update with nothing to report must answer [] and not null.
func TestSettingsSetRestartRequiredIsAlwaysAnArray(t *testing.T) {
setupAdminToken()
tempConfigFile(t, &global.ConfigPath.Global, `{"system":{"debugMode":false}}`)
w := httptest.NewRecorder()
SettingsSetHandler(w, settingsSetRequest(t, "global", `{"system":{"debugMode":true}}`))
if w.Code != http.StatusOK {
t.Fatalf("status = %d, body = %s", w.Code, w.Body.String())
}
if !strings.Contains(w.Body.String(), `"restartRequired":[]`) {
t.Errorf("restartRequired should serialize as an empty array: %s", w.Body.String())
}
}
// TestSettingsSetRejectsUnknownType keeps the failure path intact now that the
// success path computes an extra field.
func TestSettingsSetRejectsUnknownType(t *testing.T) {
setupAdminToken()
tempConfigFile(t, &global.ConfigPath.Global, `{}`)
w := httptest.NewRecorder()
SettingsSetHandler(w, settingsSetRequest(t, "nonsense", `{}`))
if w.Code != http.StatusBadRequest {
t.Errorf("status = %d, want 400", w.Code)
}
}
+73 -7
View File
@@ -3,8 +3,10 @@ package controller
import ( import (
"errors" "errors"
"fmt" "fmt"
"reflect"
"strings" "strings"
"sync" "sync"
"sync/atomic"
"nukumizu-backend/config" "nukumizu-backend/config"
"nukumizu-backend/internal/node" "nukumizu-backend/internal/node"
@@ -109,6 +111,12 @@ type BotController interface {
type Manager struct { type Manager struct {
mu sync.RWMutex mu sync.RWMutex
controllers map[string]Controller controllers map[string]Controller
// builtFrom records the controllerMethod section the registered controllers
// were built from, so a settings update that concerns them can be told apart
// from one that does not. It is read on the settings-update goroutine and
// written when the set is replaced.
builtFrom atomic.Pointer[config.ControllerMethodConfig]
} }
var globalManager *Manager var globalManager *Manager
@@ -126,12 +134,70 @@ func GetManager() *Manager {
return globalManager return globalManager
} }
// Register adds a controller to the manager. // NeedsRebuild reports whether next differs from the controllerMethod section
func (m *Manager) Register(c Controller) { // the registered controllers were built from. A manager with no controllers yet
// always reports true, so the first call installs the initial set.
//
// Controllers are rebuilt wholesale rather than reconfigured in place: each one
// reads its settings into fields at construction, and two of them own
// connections that cannot be re-pointed (the NapCat WebSocket listener is
// stopped through a sync.Once, the Telegram polling context is created with the
// controller). Replacing the set keeps every channel on the same footing.
func (m *Manager) NeedsRebuild(next config.ControllerMethodConfig) bool {
built := m.builtFrom.Load()
if built == nil {
return true
}
// The section carries a header map and a recipient slice, so it is not
// comparable with ==.
return !reflect.DeepEqual(*built, next)
}
// ReplaceAll stops every registered controller and swaps in next, which the
// caller built from method. It is the only way controllers are installed, at
// startup and after a settings change alike.
//
// The swap happens under the registry lock so routing flips to the new set
// atomically; stopping and starting happen outside it. Both can block — Stop
// closes sockets, Start performs a handshake — and holding m.mu across them
// would stall every notification for the duration.
func (m *Manager) ReplaceAll(next []Controller, method config.ControllerMethodConfig) {
m.mu.Lock() m.mu.Lock()
defer m.mu.Unlock() previous := m.controllers
m.controllers[c.Name()] = c m.controllers = make(map[string]Controller, len(next))
postLog.Info("Controller registered: " + c.Name()) for _, ctrl := range next {
m.controllers[ctrl.Name()] = ctrl
}
m.mu.Unlock()
m.builtFrom.Store(&method)
names := make([]string, 0, len(next))
for _, ctrl := range next {
names = append(names, ctrl.Name())
}
postLog.Info("Controller set installed: " + strings.Join(names, ", "))
for _, ctrl := range previous {
ctrl.Stop()
}
// Start off the calling goroutine, the way startup does: Telegram's getMe
// and the NapCat WebSocket handshake would otherwise hold the settings
// request open for as long as they take. The new controllers are already
// routable, and each one can send before Start returns.
for _, ctrl := range next {
go func(ctrl Controller) {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Controller %s panicked on start: %v", ctrl.Name(), r))
}
}()
if err := ctrl.Start(); err != nil {
postLog.Error(fmt.Sprintf("Failed to start controller %s: %v", ctrl.Name(), err))
}
}(ctrl)
}
} }
// ShowBotInitMessage sends the bot initialization message to all enabled // ShowBotInitMessage sends the bot initialization message to all enabled
@@ -144,7 +210,7 @@ func (m *Manager) ShowBotInitMessage() {
m.mu.RLock() m.mu.RLock()
defer m.mu.RUnlock() defer m.mu.RUnlock()
cfg := config.C_globalConfig cfg := config.Current()
params := template.BuildBotInitializationMsgParams() params := template.BuildBotInitializationMsgParams()
for _, ctrl := range m.controllers { for _, ctrl := range m.controllers {
@@ -176,7 +242,7 @@ func (m *Manager) ShowBotServerList() {
m.mu.RLock() m.mu.RLock()
defer m.mu.RUnlock() defer m.mu.RUnlock()
cfg := config.C_globalConfig cfg := config.Current()
params := template.BuildParamsFromServerList() params := template.BuildParamsFromServerList()
for _, ctrl := range m.controllers { for _, ctrl := range m.controllers {
+222
View File
@@ -0,0 +1,222 @@
package controller
import (
"sync"
"testing"
"time"
"nukumizu-backend/config"
"nukumizu-backend/internal/node"
)
// fakeController records the lifecycle calls it receives so a test can assert
// what ReplaceAll did to it. Start signals on started, because ReplaceAll starts
// controllers off the calling goroutine.
type fakeController struct {
name string
started chan struct{}
mu sync.Mutex
starts int
stops int
}
func newFake(name string) *fakeController {
return &fakeController{name: name, started: make(chan struct{}, 4)}
}
func (f *fakeController) Name() string { return f.name }
func (f *fakeController) IsEnabled() bool { return true }
func (f *fakeController) IsMarkdown() bool { return false }
func (f *fakeController) Start() error {
f.mu.Lock()
f.starts++
f.mu.Unlock()
select {
case f.started <- struct{}{}:
default:
}
return nil
}
func (f *fakeController) Stop() {
f.mu.Lock()
defer f.mu.Unlock()
f.stops++
}
func (f *fakeController) lifecycle() (starts, stops int) {
f.mu.Lock()
defer f.mu.Unlock()
return f.starts, f.stops
}
func (f *fakeController) SendStatusChange(node.StatusChange) error { return nil }
func (f *fakeController) SendServerList(string, string) error { return nil }
func (f *fakeController) SendExecuteResult(string, string, string, string) error { return nil }
func (f *fakeController) SendAlert(Alert) error { return nil }
// waitStarted blocks until the controller's Start has run.
func (f *fakeController) waitStarted(t *testing.T) {
t.Helper()
select {
case <-f.started:
case <-time.After(5 * time.Second):
t.Fatalf("controller %s was never started", f.name)
}
}
func newTestManager() *Manager {
return &Manager{controllers: make(map[string]Controller)}
}
// routedTo returns the controller the manager currently routes the given name
// to.
func (m *Manager) routedTo(name string) Controller {
m.mu.RLock()
defer m.mu.RUnlock()
return m.controllers[name]
}
func TestReplaceAllInstallsAndStarts(t *testing.T) {
m := newTestManager()
alpha, beta := newFake("alpha"), newFake("beta")
m.ReplaceAll([]Controller{alpha, beta}, config.ControllerMethodConfig{})
alpha.waitStarted(t)
beta.waitStarted(t)
if got := m.routedTo("alpha"); got != Controller(alpha) {
t.Errorf("alpha is not routable after ReplaceAll: %v", got)
}
if got := m.routedTo("beta"); got != Controller(beta) {
t.Errorf("beta is not routable after ReplaceAll: %v", got)
}
}
// TestReplaceAllStopsTheOutgoingSet is the invariant the rebuild relies on: the
// old controllers must be shut down, or a rebuilt NapCat or Telegram controller
// would leave its previous connection running.
func TestReplaceAllStopsTheOutgoingSet(t *testing.T) {
m := newTestManager()
outgoing := newFake("alpha")
m.ReplaceAll([]Controller{outgoing}, config.ControllerMethodConfig{})
outgoing.waitStarted(t)
incoming := newFake("alpha")
m.ReplaceAll([]Controller{incoming}, config.ControllerMethodConfig{})
incoming.waitStarted(t)
if starts, stops := outgoing.lifecycle(); starts != 1 || stops != 1 {
t.Errorf("outgoing controller lifecycle = %d starts / %d stops, want 1/1", starts, stops)
}
if _, stops := incoming.lifecycle(); stops != 0 {
t.Errorf("incoming controller was stopped %d times", stops)
}
if got := m.routedTo("alpha"); got != Controller(incoming) {
t.Error("routing still points at the outgoing controller")
}
}
func TestReplaceAllDropsChannelsLeftOut(t *testing.T) {
m := newTestManager()
m.ReplaceAll([]Controller{newFake("alpha"), newFake("beta")}, config.ControllerMethodConfig{})
m.ReplaceAll([]Controller{newFake("beta")}, config.ControllerMethodConfig{})
if got := m.routedTo("alpha"); got != nil {
t.Errorf("a channel missing from the new set is still routable: %v", got)
}
if got := m.routedTo("beta"); got == nil {
t.Error("the surviving channel is not routable")
}
}
func TestNeedsRebuild(t *testing.T) {
m := newTestManager()
base := config.ControllerMethodConfig{
Email: config.EmailConfig{Enabled: true, SMTPPort: 587, To: []string{"a@example.com"}},
}
if !m.NeedsRebuild(base) {
t.Error("a manager with nothing installed must report that a rebuild is needed")
}
m.ReplaceAll(nil, base)
if m.NeedsRebuild(base) {
t.Error("the settings the set was built from must not ask for another rebuild")
}
changed := base
changed.Email.Enabled = false
if !m.NeedsRebuild(changed) {
t.Error("a changed email setting must ask for a rebuild")
}
// The section carries maps and slices, so it is compared by value rather
// than by identity: equal content must not trigger a rebuild.
withHeaders := config.ControllerMethodConfig{
Webhook: config.WebhookConfig{Headers: map[string]string{"X-Token": "t"}},
}
m.ReplaceAll(nil, withHeaders)
equalHeaders := config.ControllerMethodConfig{
Webhook: config.WebhookConfig{Headers: map[string]string{"X-Token": "t"}},
}
if m.NeedsRebuild(equalHeaders) {
t.Error("equal header maps must not ask for a rebuild")
}
differentHeaders := config.ControllerMethodConfig{
Webhook: config.WebhookConfig{Headers: map[string]string{"X-Token": "other"}},
}
if !m.NeedsRebuild(differentHeaders) {
t.Error("a changed header must ask for a rebuild")
}
}
// TestReplaceAllDuringNotification drives ReplaceAll while notifications are
// being routed. Run with -race: the swap replaces the map the routing path
// reads, which is what the registry lock exists to make safe.
func TestReplaceAllDuringNotification(t *testing.T) {
m := newTestManager()
m.ReplaceAll([]Controller{newFake("alpha")}, config.ControllerMethodConfig{})
var readers, writers sync.WaitGroup
stop := make(chan struct{})
for i := 0; i < 3; i++ {
readers.Add(1)
go func() {
defer readers.Done()
for {
select {
case <-stop:
return
default:
}
m.NotifyStatusChange(node.StatusChange{UUID: "u1", Name: "alpha", Event: "Online"})
_ = m.IsMarkdown("alpha")
}
}()
}
writers.Add(1)
go func() {
defer writers.Done()
for i := 0; i < 20; i++ {
m.ReplaceAll([]Controller{newFake("alpha"), newFake("beta")}, config.ControllerMethodConfig{})
}
}()
// Let the writer finish, then release the readers: they only return once
// stop is closed.
writers.Wait()
close(stop)
readers.Wait()
if got := m.routedTo("beta"); got == nil {
t.Error("the last installed set is not routable")
}
}
+23 -10
View File
@@ -2,6 +2,7 @@ package pipes
import ( import (
"fmt" "fmt"
"net"
gomail "gopkg.in/mail.v2" gomail "gopkg.in/mail.v2"
@@ -14,22 +15,34 @@ import (
) )
// EmailController handles email notifications via SMTP. // EmailController handles email notifications via SMTP.
//
// cfg is written once, by the constructor, and never again: a controller that
// needs different settings is replaced wholesale by the manager rather than
// reconfigured in place, so the send methods can read it without locking.
type EmailController struct { type EmailController struct {
cfg config.EmailConfig cfg config.EmailConfig
} }
// NewEmailController creates a new Email controller. // NewEmailController creates a new Email controller.
func NewEmailController(cfg config.EmailConfig) *EmailController { func NewEmailController(cfg config.EmailConfig) *EmailController {
if cfg.NetworkUseProxy { applyEmailProxy(cfg.NetworkUseProxy)
// Route SMTP through the HTTP CONNECT proxy. NetDialTimeout is
// gomail's documented hook for overriding how the SMTP connection is
// dialed. There is a single global email channel, so overriding it
// unconditionally when the flag is set is safe.
gomail.NetDialTimeout = netproxy.DialWithTimeout(true)
}
return &EmailController{cfg: cfg} return &EmailController{cfg: cfg}
} }
// applyEmailProxy routes SMTP through the HTTP CONNECT proxy, or restores a
// direct dial. gomail exposes the dial path as the package-level
// NetDialTimeout, whose own default is net.DialTimeout, so turning the proxy
// off has to put that back rather than leave the hook in place. There is a
// single global email channel, so setting a package-level hook here is
// unambiguous.
func applyEmailProxy(useProxy bool) {
if useProxy {
gomail.NetDialTimeout = netproxy.DialWithTimeout(true)
return
}
gomail.NetDialTimeout = net.DialTimeout
}
// Name returns the controller name. // Name returns the controller name.
func (e *EmailController) Name() string { func (e *EmailController) Name() string {
return "email" return "email"
@@ -71,7 +84,7 @@ func (e *EmailController) SendStatusChange(change node.StatusChange) error {
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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)
@@ -85,7 +98,7 @@ func (e *EmailController) SendServerList(onlineServers, offlineServers string) e
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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)
@@ -98,7 +111,7 @@ func (e *EmailController) SendExecuteResult(serverName, serverUUID, command, res
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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)
@@ -0,0 +1,54 @@
package pipes
import (
"net"
"reflect"
"testing"
gomail "gopkg.in/mail.v2"
"nukumizu-backend/config"
)
// dialerIsDefault reports whether gomail still dials directly. Function values
// are not comparable in Go, so the code pointers are compared instead.
func dialerIsDefault() bool {
return reflect.ValueOf(gomail.NetDialTimeout).Pointer() ==
reflect.ValueOf(net.DialTimeout).Pointer()
}
// TestApplyEmailProxyRestoresDefaultDialer covers the case the constructor used
// to get wrong. gomail exposes its dial hook as a package-level variable with no
// unset, and the code only installed the proxy dialer when the flag was set, so
// turning networkUseProxy back off left SMTP tunnelled through a proxy nobody
// had asked for — with no way to undo it short of a restart.
func TestApplyEmailProxyRestoresDefaultDialer(t *testing.T) {
t.Cleanup(func() { gomail.NetDialTimeout = net.DialTimeout })
applyEmailProxy(true)
if dialerIsDefault() {
t.Fatal("applyEmailProxy(true) did not install the proxy dialer")
}
applyEmailProxy(false)
if !dialerIsDefault() {
t.Error("applyEmailProxy(false) left the proxy dialer in place")
}
}
// TestNewEmailControllerAppliesProxySetting pins that the dialer follows the
// settings a controller is built with, which is what makes a rebuilt controller
// pick up a changed proxy flag.
func TestNewEmailControllerAppliesProxySetting(t *testing.T) {
t.Cleanup(func() { gomail.NetDialTimeout = net.DialTimeout })
NewEmailController(config.EmailConfig{NetworkUseProxy: true})
if dialerIsDefault() {
t.Error("a controller built with networkUseProxy: true did not install the proxy dialer")
}
NewEmailController(config.EmailConfig{NetworkUseProxy: false})
if !dialerIsDefault() {
t.Error("a controller rebuilt with networkUseProxy: false left the proxy dialer in place")
}
}
+8 -3
View File
@@ -15,6 +15,11 @@ import (
) )
// NtfyController handles notifications via ntfy.sh or a self-hosted ntfy server. // NtfyController handles notifications via ntfy.sh or a self-hosted ntfy server.
//
// cfg and httpClient are written once, by the constructor, and never again: a
// controller that needs different settings is replaced wholesale by the manager
// rather than reconfigured in place, so the send methods can read them without
// locking.
type NtfyController struct { type NtfyController struct {
cfg config.NtfyConfig cfg config.NtfyConfig
httpClient *http.Client httpClient *http.Client
@@ -65,7 +70,7 @@ func (n *NtfyController) SendStatusChange(change node.StatusChange) error {
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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)
@@ -79,7 +84,7 @@ func (n *NtfyController) SendServerList(onlineServers, offlineServers string) er
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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)
@@ -92,7 +97,7 @@ func (n *NtfyController) SendExecuteResult(serverName, serverUUID, command, resu
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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)
+14 -5
View File
@@ -17,6 +17,15 @@ 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"`
@@ -155,7 +164,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 config.C_globalConfig.Debug.ShowNapcatAction { if actionLogEnabled() {
postLog.Debug(fmt.Sprintf("[Napcat] SendMsg -> %s (%s): %s", endpoint, targetType, message)) postLog.Debug(fmt.Sprintf("[Napcat] SendMsg -> %s (%s): %s", endpoint, targetType, message))
} }
@@ -176,7 +185,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 config.C_globalConfig.Debug.ShowNapcatAction { if actionLogEnabled() {
postLog.Debug(fmt.Sprintf("[Napcat] RecallMsg -> /delete_msg: %d", msgID)) postLog.Debug(fmt.Sprintf("[Napcat] RecallMsg -> /delete_msg: %d", msgID))
} }
@@ -190,7 +199,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 config.C_globalConfig.Debug.ShowNapcatAction { if actionLogEnabled() {
postLog.Debug("[Napcat] GetGroupList -> /get_group_list") postLog.Debug("[Napcat] GetGroupList -> /get_group_list")
} }
@@ -211,7 +220,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 config.C_globalConfig.Debug.ShowNapcatAction { if actionLogEnabled() {
postLog.Debug(fmt.Sprintf("[Napcat] GetGroupInfo -> /get_group_info: %d", groupID)) postLog.Debug(fmt.Sprintf("[Napcat] GetGroupInfo -> /get_group_info: %d", groupID))
} }
@@ -225,7 +234,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 config.C_globalConfig.Debug.ShowNapcatAction { if actionLogEnabled() {
postLog.Debug("[Napcat] GetFriendsList -> /get_friend_list") postLog.Debug("[Napcat] GetFriendsList -> /get_friend_list")
} }
+26 -15
View File
@@ -100,24 +100,34 @@ func (q *QQController) handleNapcatEvent(raw []byte) {
return return
} }
if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatMsg { // The configuration is read once per event: the debug guards below are
// 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 config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatAction { if showAction {
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) && !config.C_globalConfig.System.DebugMode { if q.isSelfMessage(ev) && !debugMode {
return return
} }
if q.isSelfMessage(ev) && config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.NapcatIgnoreSelfMsg { if q.isSelfMessage(ev) && debugMode && cfg.Debug.NapcatIgnoreSelfMsg {
if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatAction { if showAction {
postLog.Debug("Ignoring Napcat WS self message: " + string(raw)) postLog.Debug("Ignoring Napcat WS self message: " + string(raw))
} }
return return
@@ -138,7 +148,7 @@ func (q *QQController) handleNapcatEvent(raw []byte) {
response := q.processCommand(cmd) response := q.processCommand(cmd)
if response == "" { if response == "" {
if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatAction { if showAction {
postLog.Debug("Napcat WS command discarded: " + string(raw)) postLog.Debug("Napcat WS command discarded: " + string(raw))
} }
return return
@@ -160,6 +170,7 @@ 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
@@ -168,7 +179,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 config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowNapcatAction { if cfg != nil && cfg.System.DebugMode && cfg.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.
@@ -190,7 +201,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 config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowTriggerCmdEcho { if cfg != nil && cfg.System.DebugMode && cfg.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 {
@@ -213,7 +224,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.C_botUserConfig; c != nil { if c := config.BotUsers(); c != nil {
return c.QQ.Admins.IDs() return c.QQ.Admins.IDs()
} }
return nil return nil
@@ -221,7 +232,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.C_botUserConfig; c != nil { if c := config.BotUsers(); c != nil {
return c.QQ.TrustedGroups.IDs() return c.QQ.TrustedGroups.IDs()
} }
return nil return nil
@@ -237,7 +248,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.C_botUserConfig; uc != nil { if uc := config.BotUsers(); 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
@@ -260,12 +271,12 @@ func (q *QQController) SendStatusChange(change node.StatusChange) error {
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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.C_botUserConfig; uc != nil { if uc := config.BotUsers(); uc != nil {
for groupID, opts := range uc.QQ.TrustedGroups { for groupID, opts := range uc.QQ.TrustedGroups {
if !opts.EventStatusNotify { if !opts.EventStatusNotify {
continue continue
@@ -289,7 +300,7 @@ func (q *QQController) SendServerList(onlineServers, offlineServers string) erro
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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)
@@ -305,7 +316,7 @@ func (q *QQController) SendExecuteResult(serverName, serverUUID, command, result
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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 config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowTelegramMsg { if cfg := config.Current(); cfg != nil && cfg.System.DebugMode && cfg.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 config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowTriggerCmdEcho { if cfg := config.Current(); cfg != nil && cfg.System.DebugMode && cfg.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.C_botUserConfig; uc != nil { if uc := config.BotUsers(); 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.C_globalConfig cfg := config.Current()
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.C_botUserConfig; uc != nil { if uc := config.BotUsers(); 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.C_globalConfig cfg := config.Current()
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.C_globalConfig cfg := config.Current()
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.C_botUserConfig; c != nil { if c := config.BotUsers(); 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.C_botUserConfig; c != nil { if c := config.BotUsers(); c != nil {
return c.Telegram.TrustedGroups.IDs() return c.Telegram.TrustedGroups.IDs()
} }
return nil return nil
+8 -3
View File
@@ -16,6 +16,11 @@ import (
) )
// WebhookController handles notifications via generic HTTP webhooks. // WebhookController handles notifications via generic HTTP webhooks.
//
// cfg and httpClient are written once, by the constructor, and never again: a
// controller that needs different settings is replaced wholesale by the manager
// rather than reconfigured in place, so the send methods can read them without
// locking.
type WebhookController struct { type WebhookController struct {
cfg config.WebhookConfig cfg config.WebhookConfig
httpClient *http.Client httpClient *http.Client
@@ -66,7 +71,7 @@ func (w *WebhookController) SendStatusChange(change node.StatusChange) error {
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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)
@@ -87,7 +92,7 @@ func (w *WebhookController) SendServerList(onlineServers, offlineServers string)
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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)
@@ -108,7 +113,7 @@ func (w *WebhookController) SendExecuteResult(serverName, serverUUID, command, r
return nil return nil
} }
cfg := config.C_globalConfig cfg := config.Current()
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.C_globalConfig cfg := config.Current()
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.C_globalConfig cfg := config.Current()
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.C_globalConfig cfg := config.Current()
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.C_globalConfig cfg := config.Current()
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
} }
+12 -4
View File
@@ -18,6 +18,14 @@ 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 {
@@ -146,7 +154,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.C_globalConfig.System.DebugMode { if config.IsDebugMode() {
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)
} }
@@ -332,7 +340,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 config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowKomariTaskEcho { if taskEchoEnabled() {
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
@@ -375,7 +383,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 config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowKomariTaskEcho { if taskEchoEnabled() {
postLog.Debug(fmt.Sprintf("Polling for Komari task %s results...", taskID)) postLog.Debug(fmt.Sprintf("Polling for Komari task %s results...", taskID))
} }
@@ -393,7 +401,7 @@ func (c *Client) PollTaskResult(taskID string) ([]TaskResult, error) {
return nil, err return nil, err
} }
if done { if done {
if config.C_globalConfig.System.DebugMode && config.C_globalConfig.Debug.ShowKomariTaskEcho { if taskEchoEnabled() {
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
+4 -1
View File
@@ -409,7 +409,10 @@ 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.C_globalConfig cfg := config.Current()
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")
+29 -11
View File
@@ -2,6 +2,10 @@
// network proxy configured in the system config. Each caller decides whether // network proxy configured in the system config. Each caller decides whether
// to use the proxy by passing its own useProxy flag (the per-channel // to use the proxy by passing its own useProxy flag (the per-channel
// networkUseProxy setting), so proxying is opt-in per channel. // networkUseProxy setting), so proxying is opt-in per channel.
//
// The opt-in is captured when a client is built, but the proxy address is not:
// it is read again on every request and every dial, so editing
// system.networkProxy takes effect on clients that already exist.
package netproxy package netproxy
import ( import (
@@ -20,7 +24,11 @@ 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 {
raw := config.C_globalConfig.System.NetworkProxy cfg := config.Current()
if cfg == nil {
return nil
}
raw := cfg.System.NetworkProxy
if raw == "" { if raw == "" {
return nil return nil
} }
@@ -35,19 +43,23 @@ func proxyURL() *url.URL {
} }
// ProxyFunc returns a transport proxy function that routes requests through // ProxyFunc returns a transport proxy function that routes requests through
// the configured network proxy when enabled. It returns nil when the caller // the configured network proxy when enabled, and nil when the caller opts out
// opts out or no proxy is configured, meaning direct connection. The returned // of proxying entirely. The returned function is compatible with both
// function is compatible with both http.Transport.Proxy and // http.Transport.Proxy and websocket.Dialer.Proxy.
// websocket.Dialer.Proxy. //
// The proxy address is resolved on every call rather than once here, so a
// settings update that changes system.networkProxy reaches a client that was
// already built. That is also why opting out is the only case that returns nil:
// a function resolved to nothing at construction time would pin its client to
// whatever was configured then. A nil URL from the returned function means no
// proxy is configured and the request goes direct.
func ProxyFunc(useProxy bool) func(*http.Request) (*url.URL, error) { func ProxyFunc(useProxy bool) func(*http.Request) (*url.URL, error) {
if !useProxy { if !useProxy {
return nil return nil
} }
u := proxyURL() return func(*http.Request) (*url.URL, error) {
if u == nil { return proxyURL(), nil
return nil
} }
return http.ProxyURL(u)
} }
// HTTPClient builds an http.Client that sends traffic through the configured // HTTPClient builds an http.Client that sends traffic through the configured
@@ -67,10 +79,16 @@ func HTTPClient(useProxy bool, timeout time.Duration) *http.Client {
// through the configured HTTP CONNECT proxy when enabled. Its signature // through the configured HTTP CONNECT proxy when enabled. Its signature
// matches net.DialTimeout so it can be plugged into gomail's NetDialTimeout // matches net.DialTimeout so it can be plugged into gomail's NetDialTimeout
// to send SMTP over the proxy. // to send SMTP over the proxy.
//
// Like ProxyFunc it reads the proxy address per dial, so clearing or changing
// system.networkProxy reaches a dialer that already exists.
func DialWithTimeout(useProxy bool) func(network, addr string, timeout time.Duration) (net.Conn, error) { func DialWithTimeout(useProxy bool) func(network, addr string, timeout time.Duration) (net.Conn, error) {
u := proxyURL()
return func(network, addr string, timeout time.Duration) (net.Conn, error) { return func(network, addr string, timeout time.Duration) (net.Conn, error) {
if !useProxy || u == nil { if !useProxy {
return net.DialTimeout(network, addr, timeout)
}
u := proxyURL()
if u == nil {
return net.DialTimeout(network, addr, timeout) return net.DialTimeout(network, addr, timeout)
} }
return dialViaProxy(u, addr, timeout) return dialViaProxy(u, addr, timeout)
+133
View File
@@ -0,0 +1,133 @@
package netproxy
import (
"net"
"net/http"
"net/url"
"os"
"path/filepath"
"testing"
"time"
"nukumizu-backend/config"
)
// publishConfig writes a config.json and makes it the configuration in effect,
// which is what a settings update does.
func publishConfig(t *testing.T, body string) {
t.Helper()
path := filepath.Join(t.TempDir(), "config.json")
if err := os.WriteFile(path, []byte(body), 0o644); err != nil {
t.Fatalf("write temp config: %v", err)
}
if _, err := config.LoadGlobalConfig(path); err != nil {
t.Fatalf("LoadGlobalConfig: %v", err)
}
}
// resolve runs a proxy function and returns the URL it chose, or "" when it
// chose a direct connection.
func resolve(t *testing.T, proxy func(*http.Request) (*url.URL, error)) string {
t.Helper()
req, err := http.NewRequest(http.MethodGet, "https://example.com/", nil)
if err != nil {
t.Fatalf("NewRequest: %v", err)
}
u, err := proxy(req)
if err != nil {
t.Fatalf("proxy function: %v", err)
}
if u == nil {
return ""
}
return u.String()
}
// TestProxyFuncResolvesPerCall is the property that makes system.networkProxy
// hot-reloadable: the function handed to a transport keeps reading the live
// configuration instead of the address that was configured when it was built.
func TestProxyFuncResolvesPerCall(t *testing.T) {
publishConfig(t, `{"system":{"networkProxy":"http://127.0.0.1:7890"}}`)
proxy := ProxyFunc(true)
if proxy == nil {
t.Fatal("ProxyFunc(true) returned nil, so the channel would never proxy")
}
if got := resolve(t, proxy); got != "http://127.0.0.1:7890" {
t.Errorf("first resolution = %q", got)
}
// The same function must follow a settings update.
publishConfig(t, `{"system":{"networkProxy":"http://127.0.0.1:8888"}}`)
if got := resolve(t, proxy); got != "http://127.0.0.1:8888" {
t.Errorf("after a settings update the same function resolved %q", got)
}
// Clearing the proxy falls back to a direct connection.
publishConfig(t, `{"system":{"networkProxy":""}}`)
if got := resolve(t, proxy); got != "" {
t.Errorf("a cleared proxy still resolved %q", got)
}
}
func TestProxyFuncOptOutReturnsNil(t *testing.T) {
publishConfig(t, `{"system":{"networkProxy":"http://127.0.0.1:7890"}}`)
// A channel with networkUseProxy off must not be handed a function at all,
// so its transport keeps the default direct dialing.
if ProxyFunc(false) != nil {
t.Error("ProxyFunc(false) must return nil")
}
}
func TestProxyFuncNormalizesMissingScheme(t *testing.T) {
publishConfig(t, `{"system":{"networkProxy":"127.0.0.1:7890"}}`)
if got := resolve(t, ProxyFunc(true)); got != "http://127.0.0.1:7890" {
t.Errorf("resolved %q, want the http:// prefix added", got)
}
}
func TestProxyFuncIgnoresUnusableProxy(t *testing.T) {
// A value that cannot be parsed must leave the client dialing directly
// rather than failing every request.
publishConfig(t, `{"system":{"networkProxy":"://missing-scheme"}}`)
if got := resolve(t, ProxyFunc(true)); got != "" {
t.Errorf("an unparseable proxy resolved %q, want a direct connection", got)
}
}
// TestDialWithTimeoutDialsDirectlyWithoutProxy covers the path a cleared
// system.networkProxy takes: the dialer was built while a proxy was configured,
// and must fall back to a direct dial once there is none.
func TestDialWithTimeoutDialsDirectlyWithoutProxy(t *testing.T) {
publishConfig(t, `{"system":{"networkProxy":"http://127.0.0.1:7890"}}`)
dial := DialWithTimeout(true)
if dial == nil {
t.Fatal("DialWithTimeout(true) returned nil")
}
// No proxy is listening on that address, so a dial attempted now would
// fail; clearing the setting is what makes the direct path reachable.
publishConfig(t, `{"system":{"networkProxy":""}}`)
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen: %v", err)
}
defer ln.Close()
go func() {
conn, err := ln.Accept()
if err == nil {
conn.Close()
}
}()
conn, err := dial("tcp", ln.Addr().String(), 5*time.Second)
if err != nil {
t.Fatalf("dial through a cleared proxy: %v", err)
}
conn.Close()
}
+39 -73
View File
@@ -58,6 +58,27 @@ func main() {
postLog.SetDebugMode(cfg.System.DebugMode) postLog.SetDebugMode(cfg.System.DebugMode)
postLog.InitLogBroadcaster() postLog.InitLogBroadcaster()
// A settings update replaces the configuration in memory; these hooks push
// the new values into the state that was derived from the old one. The
// logger's debug flag is process-wide rather than read at every log call,
// and each controller holds its own copy of its channel settings plus the
// clients built from them.
config.OnReload(func(updated *config.Config) {
postLog.SetDebugMode(updated.System.DebugMode)
mgr := controller.GetManager()
if mgr == nil || !mgr.NeedsRebuild(updated.ControllerMethod) {
return
}
// Controller settings changed. Each controller reads its settings into
// fields when it is built, and the NapCat and Telegram ones own
// connections that cannot be re-pointed, so the change is applied by
// replacing the whole set rather than reconfiguring it in place.
postLog.Info("Controller settings changed; rebuilding every channel")
mgr.ReplaceAll(buildControllers(updated), updated.ControllerMethod)
})
dbPath := cfg.DBPath dbPath := cfg.DBPath
if err := postLog.InitLogsDatabase(fmt.Sprintf("%s/log.db", dbPath)); err != nil { if err := postLog.InitLogsDatabase(fmt.Sprintf("%s/log.db", dbPath)); err != nil {
@@ -199,83 +220,28 @@ func startWebhookServer(cfg *config.Config) {
}() }()
} }
// initControllers initializes and starts all configured controllers. // buildControllers constructs one controller per channel from the given
// configuration. The set is built fresh whenever controller settings change,
// because a controller reads its settings once at construction and two of them
// own connections that cannot be re-pointed.
func buildControllers(cfg *config.Config) []controller.Controller {
return []controller.Controller{
qq_napcat.NewQQController(cfg.ControllerMethod.QQ),
telegram.NewTelegramController(cfg.ControllerMethod.Telegram),
pipes.NewEmailController(cfg.ControllerMethod.Email),
pipes.NewNtfyController(cfg.ControllerMethod.Ntfy),
pipes.NewWebhookController(cfg.ControllerMethod.Webhook),
}
}
// initControllers installs the initial controller set.
func initControllers() { func initControllers() {
cfg := config.C_globalConfig cfg := config.Current()
mgr := controller.GetManager() mgr := controller.GetManager()
if mgr == nil { if mgr == nil || cfg == nil {
return return
} }
mgr.ReplaceAll(buildControllers(cfg), cfg.ControllerMethod)
// QQ (Napcat) controller.
qqCtrl := qq_napcat.NewQQController(cfg.ControllerMethod.QQ)
mgr.Register(qqCtrl)
go func() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("QQ controller panic: %v", r))
}
}()
if err := qqCtrl.Start(); err != nil {
postLog.Error("Failed to start QQ controller: " + err.Error())
}
}()
// Telegram controller.
tgCtrl := telegram.NewTelegramController(cfg.ControllerMethod.Telegram)
mgr.Register(tgCtrl)
go func() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Telegram controller panic: %v", r))
}
}()
if err := tgCtrl.Start(); err != nil {
postLog.Error("Failed to start Telegram controller: " + err.Error())
}
}()
// Email controller (status-only).
emailCtrl := pipes.NewEmailController(cfg.ControllerMethod.Email)
mgr.Register(emailCtrl)
go func() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Email controller panic: %v", r))
}
}()
if err := emailCtrl.Start(); err != nil {
postLog.Error("Failed to start Email controller: " + err.Error())
}
}()
// Ntfy controller (status-only).
ntfyCtrl := pipes.NewNtfyController(cfg.ControllerMethod.Ntfy)
mgr.Register(ntfyCtrl)
go func() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Ntfy controller panic: %v", r))
}
}()
if err := ntfyCtrl.Start(); err != nil {
postLog.Error("Failed to start Ntfy controller: " + err.Error())
}
}()
// Webhook controller (status-only).
webhookCtrl := pipes.NewWebhookController(cfg.ControllerMethod.Webhook)
mgr.Register(webhookCtrl)
go func() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Webhook controller panic: %v", r))
}
}()
if err := webhookCtrl.Start(); err != nil {
postLog.Error("Failed to start Webhook controller: " + err.Error())
}
}()
} }
// startBackgroundTasks starts periodic background goroutines. // startBackgroundTasks starts periodic background goroutines.