diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 500d7e5fd..13944ea1e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -141,6 +141,22 @@ jobs: # internal/web/service runs ~10x slower under -race and overruns the 10m default. go test -race -shuffle=on -count=1 -timeout 25m $(cat /tmp/go-packages.txt) + # A real master and node panel, each its own process, driven through node sync. + # A SKIP here means the binary was never handed over, so it fails the job. + node-e2e: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v7 + - uses: actions/setup-go@v7 + with: + go-version-file: go.mod + cache: true + - name: Master + node end to end + run: | + set -o pipefail + make node-e2e 2>&1 | tee /tmp/node-e2e.log + if grep -q -- '--- SKIP' /tmp/node-e2e.log; then echo "node-e2e skipped"; exit 1; fi + # Brief native-fuzz smoke on the security-/parser-critical decoders. Each runs the # generated corpus plus 30s of exploration; a crash here is a real input-handling bug. fuzz-smoke: diff --git a/CLAUDE.md b/CLAUDE.md index cd3298d56..fbdb4684e 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -160,9 +160,10 @@ file locations when it can answer in one hop. `-race`); `httptest` for HTTP. Keep `database.InitDB` for reopening a file or migrating a hand-built legacy DB. `internal/sub`'s `initSubDB(t)` is the template. - Code must pass `golangci-lint run` (gofumpt + goimports formatting): `make lint`. -- Postgres, xray-gRPC-e2e and scale tests `t.Skip` unless `XUI_TEST_PG_DSN`, - `XUI_DB_TYPE`+`XUI_DB_DSN`, `XRAY_E2E_BINARY` or `XUI_SCALE_TEST` is set — a - green `go test ./...` does not mean those paths ran. +- Postgres, xray-gRPC-e2e, master+node and scale tests `t.Skip` unless + `XUI_TEST_PG_DSN`, `XUI_DB_TYPE`+`XUI_DB_DSN`, `XRAY_E2E_BINARY`, + `XUI_NODE_E2E_BINARY` or `XUI_SCALE_TEST` is set — a green `go test ./...` + does not mean those paths ran. ## Frontend conventions (summary; full version in frontend/CLAUDE.md) - Ant Design 6 only — no Tailwind/shadcn. Targeted tweaks, not rewrites. @@ -191,9 +192,13 @@ reads as a broken repo, not a missing step. Run `make dist-stub` once; every make verify # gen-check + lint + typecheck + test + build + build-storybook That is the *fast* gate, not all of CI. `ci.yml` also runs `make race`, -`make vulncheck`, a live-Postgres job (where a SKIP counts as a failure) and a +`make vulncheck`, a live-Postgres job (where a SKIP counts as a failure), +`make node-e2e` (a real master and node panel, `internal/nodee2e/`) and a 30s fuzz smoke on `FuzzParseLink`/`FuzzDecodeCertPin` — run those locally when -you touch DB/dialect or parser code. +you touch DB/dialect, node sync or parser code. Node sync has two layers: every +`runtime.Remote` call gets a cell in `internal/web/node_contract_test.go` (fast, +in `make test-go`; a method without one fails it), and a multi-tick flow (cron, +adopt, node down) gets one in `internal/nodee2e/node_sync_test.go`. Common targets: `make gen` (regenerate Zod/OpenAPI), `make lint` (Go + frontend), `make test` (Go `-shuffle=on` + frontend), `make race`, `make build`. See `Makefile`. diff --git a/Makefile b/Makefile index 8b36b610d..26ec51307 100644 --- a/Makefile +++ b/Makefile @@ -58,6 +58,13 @@ test-go: dist-stub ## Go tests (shuffle, no cache) race: dist-stub ## Go tests with the race detector (needs a C compiler) go test -race -shuffle=on -count=1 -timeout 25m $(GO_PKGS) +.PHONY: node-e2e +# Two real panel processes (master + node); test-go only runs nodee2e as a skip. +NODE_E2E_BIN = $(CURDIR)/.cache/node-e2e/x-ui$(shell go env GOEXE) +node-e2e: dist-stub ## Master+node sync end to end with two real panel processes + go build -o $(NODE_E2E_BIN) . + XUI_NODE_E2E_BINARY=$(NODE_E2E_BIN) go test -count=1 -timeout 20m -v ./internal/nodee2e/ + .PHONY: test-fe test-fe: ## Frontend tests (vitest) cd $(FRONTEND) && npm test diff --git a/frontend/src/components/command-palette/CommandPalette.tsx b/frontend/src/components/command-palette/CommandPalette.tsx index 321d639dd..5861356ff 100644 --- a/frontend/src/components/command-palette/CommandPalette.tsx +++ b/frontend/src/components/command-palette/CommandPalette.tsx @@ -64,6 +64,7 @@ interface PaletteItem { export default function CommandPalette() { const { t } = useTranslation(); + const [messageApi, messageContextHolder] = message.useMessage(); const navigate = useNavigate(); const { isDark, isUltra, toggleTheme, toggleUltra, antdThemeConfig } = useTheme(); const { isOpen, close } = useCommandPalette(); @@ -194,14 +195,14 @@ export default function CommandPalette() { const copySubscription = useCallback( async (client: ClientRecord) => { if (!client.subId || !allSetting.subURI) { - message.warning(t('pages.clients.noSubId')); + messageApi.warning(t('pages.clients.noSubId')); return; } const link = `${allSetting.subURI}${client.subId}`; const ok = await ClipboardManager.copyText(link); - if (ok) message.success(t('copied')); + if (ok) messageApi.success(t('copied')); }, - [allSetting.subURI, t], + [allSetting.subURI, messageApi, t], ); const restartXray = useCallback(async () => { @@ -210,9 +211,9 @@ export default function CommandPalette() { silentSuccess: true, }); if (msg?.success) { - message.success(t('commandPalette.restartXraySuccess')); + messageApi.success(t('commandPalette.restartXraySuccess')); } - }, [close, t]); + }, [close, messageApi, t]); const cycleTheme = useCallback(() => { if (!isDark) { @@ -679,13 +680,15 @@ export default function CommandPalette() { } }; - if (!isOpen) return null; + // Kept mounted while closed: restartXray closes the palette before its toast. + if (!isOpen) return messageContextHolder; let lastCategory = ''; const themeModeClass = isUltra ? 'ultra' : isDark ? 'dark' : 'light'; return ( + {messageContextHolder}
{ const cleanSni = (sni || '').trim(); if (!cleanSni) { - message.warning(t('pages.xray.tuic.sniRequired')); + messageApi.warning(t('pages.xray.tuic.sniRequired')); return; } setValue('settings.server.certificate', `/root/cert/${cleanSni}/fullchain.pem`); @@ -51,12 +52,12 @@ export default function TuicFields() { ? await HttpUtil.get(`/panel/api/nodes/webCert/${nodeId}`, undefined, { silent: true }) : await HttpUtil.post('/panel/api/setting/all', undefined, { silent: true }); if (!msg?.success) { - message.warning(msg?.msg || t('pages.inbounds.setDefaultCertEmpty')); + messageApi.warning(msg?.msg || t('pages.inbounds.setDefaultCertEmpty')); return; } const obj = msg.obj as { webCertFile?: string; webKeyFile?: string }; if (!obj?.webCertFile && !obj?.webKeyFile) { - message.warning(t('pages.inbounds.setDefaultCertEmpty')); + messageApi.warning(t('pages.inbounds.setDefaultCertEmpty')); return; } if (obj.webCertFile) { @@ -65,9 +66,9 @@ export default function TuicFields() { if (obj.webKeyFile) { setValue('settings.server.private_key', obj.webKeyFile); } - message.success(t('pages.inbounds.setSuccess')); + messageApi.success(t('pages.inbounds.setSuccess')); } catch { - message.error(t('somethingWentWrong')); + messageApi.error(t('somethingWentWrong')); } finally { setLoadingPanelCert(false); } @@ -145,6 +146,7 @@ export default function TuicFields() { return ( <> + {messageContextHolder} ('iran-bypass'); const [includeAdblock, setIncludeAdblock] = useState(false); @@ -37,18 +38,19 @@ export default function HappSettingsContent({ const payload = buildHappPresetDeeplink(selectedPreset, includeAdblock); if (payload) { updateSetting({ subRoutingRules: payload }); - message.success(t('pages.settings.subHappPresetApplied')); + messageApi.success(t('pages.settings.subHappPresetApplied')); } }; const handleBuildDeeplink = (deeplink: string) => { updateSetting({ subRoutingRules: deeplink }); setIsModalOpen(false); - message.success(t('pages.settings.subHappDeeplinkGenerated')); + messageApi.success(t('pages.settings.subHappDeeplinkGenerated')); }; return ( <> + {messageContextHolder} (null); @@ -179,7 +180,7 @@ export default function RoutingTab({ try { parsed = JSON.parse(value); } catch { - message.error(t('pages.xray.importInvalidJson')); + messageApi.error(t('pages.xray.importInvalidJson')); return; } const obj = parsed as { rules?: unknown; routing?: { rules?: unknown } }; @@ -191,7 +192,7 @@ export default function RoutingTab({ ? obj.routing!.rules : null; if (!list) { - message.error(t('pages.xray.importInvalidJson')); + messageApi.error(t('pages.xray.importInvalidJson')); return; } mutate((tt) => { @@ -347,6 +348,7 @@ export default function RoutingTab({ return ( <> {modalContextHolder} + {messageContextHolder} { fireEvent.click(screen.getByRole('button', { name: 'Generate Deeplink' })); expect(generatedProfile()).toEqual(minimal); }); + + // The static message API outlived the test file and logged act() warnings + // after teardown, failing CI with "Closing rpc while onUserConsoleLog was pending". + it('takes its toast down with it when unmounted', async () => { + renderSettings(); + openEditor(); + fireEvent.click(screen.getByRole('button', { name: 'Generate Deeplink' })); + await screen.findByText('Deeplink generated and applied to routing rules'); + + cleanup(); + + expect(document.body.textContent).not.toContain( + 'Deeplink generated and applied to routing rules', + ); + }); }); diff --git a/frontend/src/test/happ-settings-presets.test.tsx b/frontend/src/test/happ-settings-presets.test.tsx index 3b7c34266..845776f54 100644 --- a/frontend/src/test/happ-settings-presets.test.tsx +++ b/frontend/src/test/happ-settings-presets.test.tsx @@ -1,15 +1,12 @@ import { useState } from 'react'; import { describe, expect, it, vi } from 'vitest'; import { fireEvent, screen } from '@testing-library/react'; -import { message } from 'antd'; import { AllSetting } from '@/models/setting'; import HappSettingsContent from '@/pages/settings/HappSettingsContent'; import { renderWithProviders } from './test-utils'; -vi.spyOn(message, 'success').mockImplementation(() => undefined as never); - const chinaProfile = { Name: 'Bypass-CN', GlobalProxy: 'true', diff --git a/frontend/src/test/no-static-message.test.ts b/frontend/src/test/no-static-message.test.ts new file mode 100644 index 000000000..343f8b705 --- /dev/null +++ b/frontend/src/test/no-static-message.test.ts @@ -0,0 +1,31 @@ +import { readFileSync, readdirSync, statSync } from 'node:fs'; +import { fileURLToPath } from 'node:url'; +import { join, relative, resolve } from 'node:path'; + +import { describe, expect, it } from 'vitest'; + +const srcRoot = resolve(fileURLToPath(import.meta.url), '../..'); +const staticCall = /(? { + const path = join(dir, name); + if (statSync(path).isDirectory()) return name === 'test' ? [] : sourceFiles(path); + return /\.tsx?$/.test(name) ? [path] : []; + }); +} + +// antd's static message renders outside React: it ignores the theme and its +// timers outlive the component, which broke CI after the Happ tests tore down. +describe('antd message', () => { + it('is only used through message.useMessage()', () => { + const offenders = sourceFiles(srcRoot).flatMap((file) => + readFileSync(file, 'utf8') + .split('\n') + .flatMap((line, i) => + staticCall.test(line) ? [`${relative(srcRoot, file)}:${i + 1}: ${line.trim()}`] : [], + ), + ); + expect(offenders).toEqual([]); + }); +}); diff --git a/internal/nodee2e/harness_test.go b/internal/nodee2e/harness_test.go new file mode 100644 index 000000000..2aa28ec36 --- /dev/null +++ b/internal/nodee2e/harness_test.go @@ -0,0 +1,365 @@ +// Package nodee2e drives a real master panel and a real node panel, each its own +// process, through the node-sync paths. Gated by XUI_NODE_E2E_BINARY. +package nodee2e + +import ( + "bytes" + "encoding/json" + "fmt" + "io" + "net" + "net/http" + "os" + "os/exec" + "path/filepath" + "regexp" + "strconv" + "strings" + "sync" + "testing" + "time" + + "github.com/mhsanaei/3x-ui/v3/internal/database" +) + +const settleTimeout = 30 * time.Second + +func panelBinary(t *testing.T) string { + t.Helper() + bin := os.Getenv("XUI_NODE_E2E_BINARY") + if bin == "" { + t.Skip("XUI_NODE_E2E_BINARY not set; run `make node-e2e`") + } + abs, err := filepath.Abs(bin) + if err != nil { + t.Fatalf("resolve %s: %v", bin, err) + } + return abs +} + +type panel struct { + t *testing.T + name string + bin string + dir string + port int + token string + cmd *exec.Cmd + logOut *os.File +} + +func freePort(t *testing.T) int { + t.Helper() + l, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("free port: %v", err) + } + defer l.Close() + return l.Addr().(*net.TCPAddr).Port +} + +func (p *panel) env() []string { + return append(os.Environ(), + "XUI_DB_FOLDER="+filepath.Join(p.dir, "db"), + "XUI_LOG_FOLDER="+filepath.Join(p.dir, "log"), + "XUI_BIN_FOLDER="+filepath.Join(p.dir, "bin"), + "XUI_ENABLE_FAIL2BAN=false", + "MSYS_NO_PATHCONV=1", + ) +} + +func (p *panel) cli(args ...string) string { + p.t.Helper() + cmd := exec.Command(p.bin, args...) + cmd.Env = p.env() + out, err := cmd.CombinedOutput() + if err != nil { + p.t.Fatalf("%s %v: %v\n%s", p.name, args, err, out) + } + return string(out) +} + +var apiTokenLine = regexp.MustCompile(`(?m)^apiToken:\s*(\S+)`) + +func (p *panel) mintToken(name, scope string) string { + p.t.Helper() + out := p.cli("setting", "-getApiToken", "-tokenName", name, "-tokenScope", scope) + m := apiTokenLine.FindStringSubmatch(out) + if m == nil { + p.t.Fatalf("%s: no apiToken in output:\n%s", p.name, out) + } + return m[1] +} + +// newPanel prepares a panel's database: credentials, a private port, its own +// sub-server port (two panels on one host would race for 2096) and an admin token. +func newPanel(t *testing.T, bin, name string) *panel { + t.Helper() + p := preparePanel(t, bin, name) + p.token = p.mintToken("e2e-driver", "admin") + return p +} + +// sharedDBMu guards the process-global database handle the harness borrows +// while the scopes run in parallel. +var sharedDBMu sync.Mutex + +func preparePanel(t *testing.T, bin, name string) *panel { + t.Helper() + p := &panel{t: t, name: name, bin: bin, dir: t.TempDir(), port: freePort(t)} + for _, d := range []string{"db", "log", "bin"} { + if err := os.MkdirAll(filepath.Join(p.dir, d), 0o755); err != nil { + t.Fatal(err) + } + } + p.cli("setting", "-username", "e2e", "-password", "e2e-pass", "-port", strconv.Itoa(p.port), "-webBasePath", "/") + sharedDBMu.Lock() + defer sharedDBMu.Unlock() + if err := database.InitDB(filepath.Join(p.dir, "db", "x-ui.db")); err != nil { + t.Fatalf("%s: open db: %v", name, err) + } + db := database.GetDB() + db.Exec("DELETE FROM settings WHERE key = ?", "subPort") + if err := db.Exec("INSERT INTO settings(key, value) VALUES (?, ?)", "subPort", strconv.Itoa(freePort(t))).Error; err != nil { + t.Fatalf("%s: set subPort: %v", name, err) + } + if err := database.CloseDB(); err != nil { + t.Fatalf("%s: close db: %v", name, err) + } + t.Cleanup(p.stop) + return p +} + +func (p *panel) start() { + p.t.Helper() + logOut, err := os.OpenFile(filepath.Join(p.dir, "stdout.log"), os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o644) + if err != nil { + p.t.Fatal(err) + } + p.logOut = logOut + p.cmd = exec.Command(p.bin, "run") + p.cmd.Env = p.env() + p.cmd.Stdout = logOut + p.cmd.Stderr = logOut + if err := p.cmd.Start(); err != nil { + p.t.Fatalf("%s: start: %v", p.name, err) + } + eventually(p.t, settleTimeout, p.name+" answers /server/status", func() (bool, string) { + env, err := p.try(http.MethodGet, "/panel/api/server/status", nil) + if err != nil { + return false, err.Error() + } + return env.Success, env.Msg + }) +} + +func (p *panel) stop() { + if p.cmd == nil || p.cmd.Process == nil { + return + } + _ = p.cmd.Process.Kill() + _, _ = p.cmd.Process.Wait() + p.cmd = nil + if p.logOut != nil { + _ = p.logOut.Close() + p.logOut = nil + } + if p.t.Failed() { + if b, err := os.ReadFile(filepath.Join(p.dir, "stdout.log")); err == nil { + tail := string(b) + if len(tail) > 6000 { + tail = tail[len(tail)-6000:] + } + p.t.Logf("---- %s stdout tail ----\n%s", p.name, tail) + } + } +} + +// deleteInboundRow simulates a node that lost an inbound (restore, reinstall) +// while stopped; it must not run against a live panel. +func (p *panel) deleteInboundRow(id int) { + p.t.Helper() + if p.cmd != nil { + p.t.Fatalf("%s: deleteInboundRow on a running panel", p.name) + } + sharedDBMu.Lock() + defer sharedDBMu.Unlock() + if err := database.InitDB(filepath.Join(p.dir, "db", "x-ui.db")); err != nil { + p.t.Fatalf("%s: open db: %v", p.name, err) + } + defer func() { _ = database.CloseDB() }() + db := database.GetDB() + for _, q := range []string{"DELETE FROM client_inbounds WHERE inbound_id = ?", "DELETE FROM client_traffics WHERE inbound_id = ?", "DELETE FROM inbounds WHERE id = ?"} { + if err := db.Exec(q, id).Error; err != nil { + p.t.Fatalf("%s: %s: %v", p.name, q, err) + } + } +} + +func (p *panel) url() string { return "http://127.0.0.1:" + strconv.Itoa(p.port) } + +type envelope struct { + Success bool `json:"success"` + Msg string `json:"msg"` + Obj json.RawMessage `json:"obj"` +} + +func (p *panel) try(method, path string, body any) (*envelope, error) { + var rd io.Reader + if body != nil { + b, err := json.Marshal(body) + if err != nil { + return nil, err + } + rd = bytes.NewReader(b) + } + req, err := http.NewRequest(method, p.url()+path, rd) + if err != nil { + return nil, err + } + req.Header.Set("Authorization", "Bearer "+p.token) + if body != nil { + req.Header.Set("Content-Type", "application/json") + } + resp, err := (&http.Client{Timeout: 20 * time.Second}).Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + raw, err := io.ReadAll(resp.Body) + if err != nil { + return nil, err + } + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("HTTP %d: %s", resp.StatusCode, raw) + } + var env envelope + if err := json.Unmarshal(raw, &env); err != nil { + return nil, fmt.Errorf("decode %s: %w (%s)", path, err, raw) + } + return &env, nil +} + +// call fails the test on transport errors or success:false. +func (p *panel) call(method, path string, body any) json.RawMessage { + p.t.Helper() + env, err := p.try(method, path, body) + if err != nil { + p.t.Fatalf("%s %s %s: %v", p.name, method, path, err) + } + if !env.Success { + p.t.Fatalf("%s %s %s: success=false msg=%q", p.name, method, path, env.Msg) + } + return env.Obj +} + +func eventually(t *testing.T, timeout time.Duration, what string, check func() (bool, string)) { + t.Helper() + deadline := time.Now().Add(timeout) + last := "" + for { + ok, detail := check() + if ok { + return + } + last = detail + if time.Now().After(deadline) { + t.Fatalf("timed out after %s waiting for %s; last: %s", timeout, what, last) + } + time.Sleep(500 * time.Millisecond) + } +} + +// inboundView is the subset of an inbound row the scenarios assert on. +type inboundView struct { + Id int `json:"id"` + Remark string `json:"remark"` + Enable bool `json:"enable"` + Port int `json:"port"` + Tag string `json:"tag"` + NodeID *int `json:"nodeId"` + Settings json.RawMessage `json:"settings"` +} + +type clientEntry map[string]any + +func (c clientEntry) email() string { s, _ := c["email"].(string); return s } + +func (ib inboundView) clients() []clientEntry { + raw := ib.Settings + var asString string + if json.Unmarshal(raw, &asString) == nil { + raw = json.RawMessage(asString) + } + var s struct { + Clients []clientEntry `json:"clients"` + } + _ = json.Unmarshal(raw, &s) + return s.Clients +} + +func (ib inboundView) emails() []string { + out := []string{} + for _, c := range ib.clients() { + out = append(out, c.email()) + } + return out +} + +func (ib inboundView) client(email string) clientEntry { + for _, c := range ib.clients() { + if strings.EqualFold(c.email(), email) { + return c + } + } + return nil +} + +func (p *panel) inbounds() []inboundView { + p.t.Helper() + var list []inboundView + if err := json.Unmarshal(p.call(http.MethodGet, "/panel/api/inbounds/list", nil), &list); err != nil { + p.t.Fatalf("%s: decode inbound list: %v", p.name, err) + } + return list +} + +func (p *panel) inboundOnPort(port int) (inboundView, bool) { + p.t.Helper() + for _, ib := range p.inbounds() { + if ib.Port == port { + return ib, true + } + } + return inboundView{}, false +} + +const tcpStream = `{"network":"tcp","security":"none","tcpSettings":{"header":{"type":"none"}}}` + +func vlessInbound(remark string, port int, nodeID *int, clients ...map[string]any) map[string]any { + if clients == nil { + clients = []map[string]any{} + } + settings, _ := json.Marshal(map[string]any{"clients": clients, "decryption": "none"}) + body := map[string]any{ + "remark": remark, "enable": true, "port": port, "protocol": "vless", + "settings": string(settings), "streamSettings": tcpStream, "sniffing": `{}`, + } + if nodeID != nil { + body["nodeId"] = *nodeID + } + return body +} + +func vlessClient(email string) map[string]any { + return map[string]any{"email": email, "enable": true, "id": newUUID(email)} +} + +// newUUID derives a stable, valid UUID from a label so failures are reproducible. +func newUUID(label string) string { + var b [16]byte + copy(b[:], []byte(label+"________________")) + b[6] = (b[6] & 0x0f) | 0x40 + b[8] = (b[8] & 0x3f) | 0x80 + return fmt.Sprintf("%x-%x-%x-%x-%x", b[0:4], b[4:6], b[6:8], b[8:10], b[10:16]) +} diff --git a/internal/nodee2e/node_sync_test.go b/internal/nodee2e/node_sync_test.go new file mode 100644 index 000000000..410b13ea1 --- /dev/null +++ b/internal/nodee2e/node_sync_test.go @@ -0,0 +1,344 @@ +package nodee2e + +import ( + "encoding/json" + "fmt" + "net/http" + "slices" + "strconv" + "testing" + "time" +) + +// TestNodeSync walks one master/node pair per enrollment scope through every +// operation that must converge onto the node. Each subtest names its invariant. +func TestNodeSync(t *testing.T) { + bin := panelBinary(t) + for _, scope := range []string{"admin", "node-sync"} { + t.Run("enrolled with "+scope+" token", func(t *testing.T) { + t.Parallel() + runNodeSyncScenarios(t, bin, scope) + }) + } +} + +type pair struct { + t *testing.T + master *panel + node *panel + nodeID int +} + +func (pr *pair) nodeInbound(port int) (inboundView, bool) { return pr.node.inboundOnPort(port) } + +func (pr *pair) masterInbound(port int) (inboundView, bool) { + for _, ib := range pr.master.inbounds() { + if ib.Port == port && ib.NodeID != nil && *ib.NodeID == pr.nodeID { + return ib, true + } + } + return inboundView{}, false +} + +func (pr *pair) waitNode(what string, port int, ok func(inboundView) bool) { + pr.t.Helper() + eventually(pr.t, settleTimeout, what, func() (bool, string) { + ib, found := pr.nodeInbound(port) + if !found { + return ok(inboundView{}) && false, fmt.Sprintf("node has no inbound on %d", port) + } + return ok(ib), fmt.Sprintf("node inbound %d: enable=%v remark=%q emails=%v", port, ib.Enable, ib.Remark, ib.emails()) + }) +} + +func (pr *pair) waitNodeAbsent(what string, port int) { + pr.t.Helper() + eventually(pr.t, settleTimeout, what, func() (bool, string) { + ib, found := pr.nodeInbound(port) + return !found, fmt.Sprintf("node still has inbound %d with %v", port, ib.emails()) + }) +} + +func (pr *pair) bulkAttach(emails []string, masterInboundID int) { + pr.t.Helper() + var res struct { + Attached []string `json:"attached"` + Errors []string `json:"errors"` + } + obj := pr.master.call(http.MethodPost, "/panel/api/clients/bulkAttach", map[string]any{"emails": emails, "inboundIds": []int{masterInboundID}}) + if err := json.Unmarshal(obj, &res); err != nil { + pr.t.Fatalf("decode bulkAttach: %v", err) + } + if len(res.Errors) != 0 || len(res.Attached) != len(emails) { + pr.t.Fatalf("bulkAttach attached=%d/%d errors=%v", len(res.Attached), len(emails), res.Errors) + } +} + +func emailRange(prefix string, from, to int) []string { + out := make([]string, 0, to-from+1) + for i := from; i <= to; i++ { + out = append(out, prefix+strconv.Itoa(i)) + } + return out +} + +func hasAll(have []string, want ...string) bool { + for _, w := range want { + if !slices.Contains(have, w) { + return false + } + } + return true +} + +func runNodeSyncScenarios(t *testing.T, bin, scope string) { + master := newPanel(t, bin, "master") + node := newPanel(t, bin, "node") + linkToken := node.mintToken("master-link", scope) + master.start() + node.start() + pr := &pair{t: t, master: master, node: node} + + var ( + adoptedPort = freePort(t) + madePort = freePort(t) + lostPort = freePort(t) + droppedPort = freePort(t) + offlinePort = freePort(t) + unmanagedPort = freePort(t) + localPort = freePort(t) + ) + node.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("pre-existing", adoptedPort, nil)) + + local := master.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("local-pool", localPort, nil)) + var localIb inboundView + _ = json.Unmarshal(local, &localIb) + for _, email := range emailRange("p", 1, 45) { + master.call(http.MethodPost, "/panel/api/clients/add", map[string]any{ + "client": map[string]any{"email": email, "enable": true}, "inboundIds": []int{localIb.Id}, + }) + } + + var nodeView struct { + Id int `json:"id"` + } + obj := master.call(http.MethodPost, "/panel/api/nodes/add", map[string]any{ + "name": "n1", "scheme": "http", "address": "127.0.0.1", "port": node.port, "basePath": "/", + "apiToken": linkToken, "enable": true, "allowPrivateAddress": true, + }) + if err := json.Unmarshal(obj, &nodeView); err != nil || nodeView.Id == 0 { + t.Fatalf("decode node add: %v (%s)", err, obj) + } + pr.nodeID = nodeView.Id + + var adoptedID, madeID int + t.Run("an inbound already on the node is adopted by the master", func(t *testing.T) { + pr.t = t + eventually(t, settleTimeout, "master adopts the node inbound", func() (bool, string) { + ib, ok := pr.masterInbound(adoptedPort) + adoptedID = ib.Id + return ok, "not adopted yet" + }) + }) + if adoptedID == 0 { + t.Fatal("no adopted inbound; later scenarios depend on it") + } + + t.Run("an inbound created on the master for the node lands there with its clients", func(t *testing.T) { + pr.t = t + obj := master.call(http.MethodPost, "/panel/api/inbounds/add", + vlessInbound("made-on-master", madePort, &pr.nodeID, vlessClient("m1"), vlessClient("m2"))) + var ib inboundView + _ = json.Unmarshal(obj, &ib) + madeID = ib.Id + pr.waitNode("node holds the master-made inbound", madePort, func(ib inboundView) bool { + return hasAll(ib.emails(), "m1", "m2") + }) + }) + + t.Run("editing the inbound on the master updates the node and keeps its clients", func(t *testing.T) { + pr.t = t + body := vlessInbound("renamed-on-master", madePort, &pr.nodeID) + body["settings"] = `{"decryption":"none"}` + master.call(http.MethodPost, "/panel/api/inbounds/update/"+strconv.Itoa(madeID), body) + pr.waitNode("node shows the new remark with both clients", madePort, func(ib inboundView) bool { + return ib.Remark == "renamed-on-master" && hasAll(ib.emails(), "m1", "m2") + }) + }) + + t.Run("attaching a few existing clients reaches the node", func(t *testing.T) { + pr.t = t + pr.bulkAttach([]string{"p1", "p2", "p3"}, madeID) + pr.waitNode("node holds p1..p3", madePort, func(ib inboundView) bool { + return hasAll(ib.emails(), "m1", "m2", "p1", "p2", "p3") + }) + }) + + t.Run("attaching more clients than the per-client push limit reaches the node", func(t *testing.T) { + pr.t = t + emails := emailRange("p", 4, 43) + pr.bulkAttach(emails, adoptedID) + pr.waitNode("node holds all 40", adoptedPort, func(ib inboundView) bool { + return hasAll(ib.emails(), emails...) + }) + }) + + t.Run("disabling a client on the master disables it on the node", func(t *testing.T) { + pr.t = t + mib, _ := pr.masterInbound(madePort) + entry := mib.client("p1") + if entry == nil { + t.Fatalf("master inbound has no p1: %v", mib.emails()) + } + entry["enable"] = false + master.call(http.MethodPost, "/panel/api/clients/update/p1", entry) + pr.waitNode("node p1 disabled", madePort, func(ib inboundView) bool { + c := ib.client("p1") + return c != nil && c["enable"] == false + }) + }) + + t.Run("detaching a client from the node inbound removes it there", func(t *testing.T) { + pr.t = t + master.call(http.MethodPost, "/panel/api/clients/p2/detach", map[string]any{"inboundIds": []int{madeID}}) + pr.waitNode("node drops p2", madePort, func(ib inboundView) bool { + return ib.client("p2") == nil && ib.client("p3") != nil + }) + }) + + t.Run("deleting a client on the master removes it from the node", func(t *testing.T) { + pr.t = t + master.call(http.MethodPost, "/panel/api/clients/del/p4", nil) + pr.waitNode("node drops p4", adoptedPort, func(ib inboundView) bool { + return ib.client("p4") == nil && ib.client("p5") != nil + }) + }) + + t.Run("switching the inbound off on the master switches it off on the node", func(t *testing.T) { + pr.t = t + master.call(http.MethodPost, "/panel/api/inbounds/setEnable/"+strconv.Itoa(madeID), map[string]any{"enable": false}) + pr.waitNode("node inbound disabled", madePort, func(ib inboundView) bool { return !ib.Enable }) + }) + + t.Run("node traffic reaches the master and a master reset clears the node", func(t *testing.T) { + pr.t = t + node.call(http.MethodPost, "/panel/api/clients/updateTraffic/p5", map[string]any{"upload": 1000, "download": 2000}) + usage := func(p *panel) int64 { + var tr struct{ Up, Down int64 } + _ = json.Unmarshal(p.call(http.MethodGet, "/panel/api/clients/traffic/p5", nil), &tr) + return tr.Up + tr.Down + } + eventually(t, settleTimeout, "master sees p5's node traffic", func() (bool, string) { + u := usage(master) + return u == 3000, fmt.Sprintf("master p5 usage %d", u) + }) + master.call(http.MethodPost, "/panel/api/clients/resetTraffic/p5", nil) + eventually(t, settleTimeout, "node p5 usage reset", func() (bool, string) { + u := usage(node) + return u == 0, fmt.Sprintf("node p5 usage %d", u) + }) + time.Sleep(12 * time.Second) + if u := usage(master); u != 0 { + t.Fatalf("master p5 usage %d after reset settled, want 0", u) + } + }) + + t.Run("an inbound deleted on the node is removed from the master too", func(t *testing.T) { + pr.t = t + master.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("deleted-on-node", lostPort, &pr.nodeID, vlessClient("l1"))) + pr.waitNode("node holds the inbound", lostPort, func(ib inboundView) bool { return ib.client("l1") != nil }) + nib, _ := pr.nodeInbound(lostPort) + node.call(http.MethodPost, "/panel/api/inbounds/del/"+strconv.Itoa(nib.Id), nil) + eventually(t, settleTimeout, "master mirrors the node-side delete (#6219)", func() (bool, string) { + _, still := pr.masterInbound(lostPort) + return !still, "master still has the inbound" + }) + }) + + t.Run("deleting the inbound on the master removes it from the node", func(t *testing.T) { + pr.t = t + obj := master.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("deleted-on-master", droppedPort, &pr.nodeID, vlessClient("d1"))) + var ib inboundView + _ = json.Unmarshal(obj, &ib) + pr.waitNode("node holds the inbound", droppedPort, func(ib inboundView) bool { return ib.client("d1") != nil }) + master.call(http.MethodPost, "/panel/api/inbounds/del/"+strconv.Itoa(ib.Id), nil) + pr.waitNodeAbsent("node drops the inbound", droppedPort) + }) + + t.Run("a change made while the node is down reaches it once it is back", func(t *testing.T) { + pr.t = t + node.stop() + pr.bulkAttach([]string{"p44"}, adoptedID) + node.start() + pr.waitNode("node holds p44 after restart", adoptedPort, func(ib inboundView) bool { + return ib.client("p44") != nil + }) + }) + + t.Run("an inbound the node lost while down is re-created with the master's pending change", func(t *testing.T) { + pr.t = t + nib, ok := pr.nodeInbound(adoptedPort) + if !ok { + t.Fatal("node has no adopted inbound to lose") + } + node.stop() + node.deleteInboundRow(nib.Id) + pr.bulkAttach([]string{"p45"}, adoptedID) + node.start() + pr.waitNode("node re-creates the inbound with p44 and p45", adoptedPort, func(ib inboundView) bool { + return ib.client("p44") != nil && ib.client("p45") != nil + }) + }) + + t.Run("an inbound created and edited while the node is down lands once it is back", func(t *testing.T) { + pr.t = t + node.stop() + obj := master.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("made-while-down", offlinePort, &pr.nodeID, vlessClient("o1"))) + var ib inboundView + _ = json.Unmarshal(obj, &ib) + body := vlessInbound("edited-while-down", offlinePort, &pr.nodeID) + body["settings"] = `{"decryption":"none"}` + master.call(http.MethodPost, "/panel/api/inbounds/update/"+strconv.Itoa(ib.Id), body) + node.start() + pr.waitNode("node holds the edited inbound with o1", offlinePort, func(ib inboundView) bool { + return ib.Remark == "edited-while-down" && ib.client("o1") != nil + }) + }) + + t.Run("a change made while the node is disabled on the master lands once it is re-enabled", func(t *testing.T) { + pr.t = t + nodePath := "/panel/api/nodes/setEnable/" + strconv.Itoa(pr.nodeID) + master.call(http.MethodPost, nodePath, map[string]any{"enable": false}) + pr.bulkAttach([]string{"p6"}, madeID) + time.Sleep(6 * time.Second) + if ib, _ := pr.nodeInbound(madePort); ib.client("p6") != nil { + t.Fatal("a disabled node received a push") + } + master.call(http.MethodPost, nodePath, map[string]any{"enable": true}) + pr.waitNode("node holds p6 after re-enable", madePort, func(ib inboundView) bool { return ib.client("p6") != nil }) + }) + + t.Run("selected sync mode leaves the node's unselected inbounds alone", func(t *testing.T) { + pr.t = t + var selected []string + for _, ib := range master.inbounds() { + if ib.NodeID != nil && *ib.NodeID == pr.nodeID { + selected = append(selected, ib.Tag) + } + } + master.call(http.MethodPost, "/panel/api/nodes/update/"+strconv.Itoa(pr.nodeID), map[string]any{ + "name": "n1", "scheme": "http", "address": "127.0.0.1", "port": node.port, "basePath": "/", + "enable": true, "allowPrivateAddress": true, "inboundSyncMode": "selected", "inboundTags": selected, + }) + node.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("node-only", unmanagedPort, nil, vlessClient("u1"))) + pr.bulkAttach([]string{"p7"}, madeID) + pr.waitNode("selected inbound still converges", madePort, func(ib inboundView) bool { return ib.client("p7") != nil }) + time.Sleep(12 * time.Second) + if _, adopted := pr.masterInbound(unmanagedPort); adopted { + t.Fatal("master adopted an unselected node inbound") + } + if ib, ok := pr.nodeInbound(unmanagedPort); !ok || ib.client("u1") == nil { + t.Fatal("reconcile swept or rewrote an unselected node inbound") + } + }) +} diff --git a/internal/util/wirecodec/wirecodec.go b/internal/util/wirecodec/wirecodec.go index 6c34a0ad8..63e9e09c3 100644 --- a/internal/util/wirecodec/wirecodec.go +++ b/internal/util/wirecodec/wirecodec.go @@ -16,6 +16,8 @@ const ( HashHeader = "X-Config-Sha256" // CapsHeader is set by a node on its API responses to advertise support. CapsHeader = "X-3x-Node-Caps" + // MasterPushHeader marks a request as a master's push, whatever its token scope. + MasterPushHeader = "X-3x-Master-Push" // EncodingZstd is the Content-Encoding value for a zstd-compressed body. EncodingZstd = "zstd" // CapZstd is the capability token advertised in CapsHeader. diff --git a/internal/web/controller/api.go b/internal/web/controller/api.go index c02eb2cd1..202b58490 100644 --- a/internal/web/controller/api.go +++ b/internal/web/controller/api.go @@ -98,6 +98,7 @@ var nodeSyncScopeAllow = map[string]map[string]struct{}{ "/inbounds/add": {http.MethodPost: {}}, "/inbounds/del/:id": {http.MethodPost: {}}, "/inbounds/update/:id": {http.MethodPost: {}}, + "/inbounds/:id/subSortIndex": {http.MethodPost: {}}, "/clients/add": {http.MethodPost: {}}, "/clients/del/:email": {http.MethodPost: {}}, "/clients/:email/detach": {http.MethodPost: {}}, @@ -106,11 +107,13 @@ var nodeSyncScopeAllow = map[string]map[string]struct{}{ "/server/getWebCertFiles": {http.MethodGet: {}}, "/server/descendants": {http.MethodGet: {}}, "/clients/resetTraffic/:email": {http.MethodPost: {}}, + "/clients/bulkResetTraffic": {http.MethodPost: {}}, "/inbounds/resetAllTraffics": {http.MethodPost: {}}, "/inbounds/:id/resetTraffic": {http.MethodPost: {}}, "/clients/onlinesByGuid": {http.MethodPost: {}}, "/clients/onlines": {http.MethodPost: {}}, "/clients/lastOnline": {http.MethodPost: {}}, + "/clients/activeInbounds": {http.MethodPost: {}}, "/inbounds/pushClientTraffics": {http.MethodPost: {}}, "/server/clientIps": {http.MethodGet: {}, http.MethodPost: {}}, "/clients/clientIpsByGuid": {http.MethodPost: {}}, diff --git a/internal/web/controller/api_auth_test.go b/internal/web/controller/api_auth_test.go index 44e7cc7a1..e454d6597 100644 --- a/internal/web/controller/api_auth_test.go +++ b/internal/web/controller/api_auth_test.go @@ -7,7 +7,6 @@ import ( "net/http/cookiejar" "net/http/httptest" "path/filepath" - "reflect" "testing" "github.com/gin-contrib/sessions" @@ -138,39 +137,6 @@ func TestCheckAPIAuth_AcceptsVerifiedClientCert(t *testing.T) { } } -func TestNodeSyncScopeAllowlistMatchesRemoteInventory(t *testing.T) { - expected := map[string]map[string]struct{}{ - "/server/status": {http.MethodGet: {}}, - "/inbounds/list": {http.MethodGet: {}}, - "/inbounds/add": {http.MethodPost: {}}, - "/inbounds/del/:id": {http.MethodPost: {}}, - "/inbounds/update/:id": {http.MethodPost: {}}, - "/clients/add": {http.MethodPost: {}}, - "/clients/del/:email": {http.MethodPost: {}}, - "/clients/:email/detach": {http.MethodPost: {}}, - "/clients/update/:email": {http.MethodPost: {}}, - "/server/restartXrayService": {http.MethodPost: {}}, - "/server/getWebCertFiles": {http.MethodGet: {}}, - "/server/descendants": {http.MethodGet: {}}, - "/clients/resetTraffic/:email": {http.MethodPost: {}}, - "/inbounds/resetAllTraffics": {http.MethodPost: {}}, - "/inbounds/:id/resetTraffic": {http.MethodPost: {}}, - "/clients/onlinesByGuid": {http.MethodPost: {}}, - "/clients/onlines": {http.MethodPost: {}}, - "/clients/lastOnline": {http.MethodPost: {}}, - "/inbounds/pushClientTraffics": {http.MethodPost: {}}, - "/server/clientIps": {http.MethodGet: {}, http.MethodPost: {}}, - "/clients/clientIpsByGuid": {http.MethodPost: {}}, - "/hosts/list": {http.MethodGet: {}}, - } - if !reflect.DeepEqual(nodeSyncScopeAllow, expected) { - t.Fatalf("node-sync allowlist drift:\n got: %#v\nwant: %#v", nodeSyncScopeAllow, expected) - } - if _, ok := nodeSyncScopeAllow["/server/updatePanel"]; ok { - t.Fatal("node-sync must not include /server/updatePanel") - } -} - func TestNodeSyncScopeUsesFullPathPatterns(t *testing.T) { engine, _ := newAPIAuthTestEngine(t) cases := []struct { diff --git a/internal/web/controller/inbound.go b/internal/web/controller/inbound.go index ca4879dc1..e990fe25f 100644 --- a/internal/web/controller/inbound.go +++ b/internal/web/controller/inbound.go @@ -7,6 +7,7 @@ import ( "strings" "github.com/mhsanaei/3x-ui/v3/internal/database/model" + "github.com/mhsanaei/3x-ui/v3/internal/util/wirecodec" "github.com/mhsanaei/3x-ui/v3/internal/web/middleware" "github.com/mhsanaei/3x-ui/v3/internal/web/service" "github.com/mhsanaei/3x-ui/v3/internal/web/session" @@ -64,7 +65,9 @@ func (a *InboundController) broadcastInboundsUpdate(userId int) { func (a *InboundController) inboundServiceFor(c *gin.Context) *service.InboundService { svc := a.inboundService scope, _ := c.Get("api_token_scope") - svc.FromNodeSync = scope == model.ApiScopeNodeSync + // A master enrolled with an admin token (the -getApiToken default) has no + // node-sync scope, so it marks every request it sends instead. + svc.FromNodeSync = scope == model.ApiScopeNodeSync || c.GetHeader(wirecodec.MasterPushHeader) != "" return &svc } diff --git a/internal/web/node_contract_test.go b/internal/web/node_contract_test.go new file mode 100644 index 000000000..d76d7a939 --- /dev/null +++ b/internal/web/node_contract_test.go @@ -0,0 +1,382 @@ +package web + +import ( + "context" + "net/http" + "net/http/httptest" + "net/url" + "path/filepath" + "reflect" + "strconv" + "strings" + "sync" + "testing" + "time" + + "github.com/robfig/cron/v3" + + "github.com/mhsanaei/3x-ui/v3/internal/database" + "github.com/mhsanaei/3x-ui/v3/internal/database/dbtest" + "github.com/mhsanaei/3x-ui/v3/internal/database/model" + "github.com/mhsanaei/3x-ui/v3/internal/util/crypto" + "github.com/mhsanaei/3x-ui/v3/internal/web/global" + "github.com/mhsanaei/3x-ui/v3/internal/web/runtime" + "github.com/mhsanaei/3x-ui/v3/internal/web/service" + "github.com/mhsanaei/3x-ui/v3/internal/xray" +) + +// nodeUnderContract serves the production router as a node and records every +// request the node refused for auth or scope. +type nodeUnderContract struct { + srv *httptest.Server + mu sync.Mutex + refused []string +} + +func startContractNode(t *testing.T) *nodeUnderContract { + t.Helper() + dbDir := t.TempDir() + t.Setenv("XUI_DB_FOLDER", dbDir) + dbtest.InitDB(t, filepath.Join(dbDir, "x-ui.db")) + prevMgr := runtime.GetManager() + runtime.SetManager(runtime.NewManager(runtime.LocalDeps{APIPort: func() int { return 0 }, SetNeedRestart: func() {}})) + t.Cleanup(func() { runtime.SetManager(prevMgr) }) + + previous := global.GetWebServer() + s := NewServer() + s.cron = cron.New(cron.WithLocation(time.Local), cron.WithSeconds()) + global.SetWebServer(s) + t.Cleanup(func() { + s.cancel() + global.SetWebServer(previous) + }) + engine, err := s.initRouter() + if err != nil { + t.Fatalf("initRouter: %v", err) + } + n := &nodeUnderContract{} + n.srv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK} + engine.ServeHTTP(rec, r) + if rec.status == http.StatusUnauthorized || rec.status == http.StatusForbidden { + n.mu.Lock() + n.refused = append(n.refused, r.Method+" "+r.URL.Path+" -> "+strconv.Itoa(rec.status)) + n.mu.Unlock() + } + })) + t.Cleanup(n.srv.Close) + return n +} + +type statusRecorder struct { + http.ResponseWriter + status int +} + +func (r *statusRecorder) WriteHeader(code int) { + r.status = code + r.ResponseWriter.WriteHeader(code) +} + +func (n *nodeUnderContract) takeRefused() []string { + n.mu.Lock() + defer n.mu.Unlock() + out := n.refused + n.refused = nil + return out +} + +func (n *nodeUnderContract) masterWithToken(t *testing.T, scope string) *runtime.Remote { + t.Helper() + token := "contract-" + scope + if err := database.GetDB().Create(&model.ApiToken{ + Name: "master-" + scope, Token: crypto.HashTokenSHA256(token), Enabled: true, Scope: scope, + }).Error; err != nil { + t.Fatalf("seed %s token: %v", scope, err) + } + u, _ := url.Parse(n.srv.URL) + port, _ := strconv.Atoi(u.Port()) + return runtime.NewRemote(&model.Node{ + Id: 1, Name: "contract-node", Scheme: "http", Address: u.Hostname(), Port: port, + BasePath: "/", ApiToken: token, Enable: true, AllowPrivateAddress: true, + }, nil) +} + +func nodeRow(t *testing.T, tag string) (*model.Inbound, bool) { + t.Helper() + var ib model.Inbound + err := database.GetDB().Where("tag = ?", tag).First(&ib).Error + return &ib, err == nil +} + +func nodeTraffic(t *testing.T, email string) int64 { + t.Helper() + var ct xray.ClientTraffic + if err := database.GetDB().Where("email = ?", email).First(&ct).Error; err != nil { + t.Fatalf("client_traffics %s: %v", email, err) + } + return ct.Up + ct.Down +} + +func seedNodeTraffic(t *testing.T, emails ...string) { + t.Helper() + for _, e := range emails { + if err := database.GetDB().Model(&xray.ClientTraffic{}).Where("email = ?", e). + Updates(map[string]any{"up": 100, "down": 200}).Error; err != nil { + t.Fatalf("seed traffic %s: %v", e, err) + } + } + if err := database.GetDB().Model(&model.Inbound{}).Where("tag = ?", contractTag). + Updates(map[string]any{"up": 100, "down": 200}).Error; err != nil { + t.Fatalf("seed inbound traffic: %v", err) + } +} + +const contractTag = "in-51001-tcp" + +func masterInbound(remark string, enable bool, clients ...string) *model.Inbound { + entries := make([]string, 0, len(clients)) + for i, email := range clients { + entries = append(entries, `{"email":"`+email+`","enable":true,"subId":"s-`+email+ + `","id":"0b6d5c2e-7c1a-4f4e-9d3b-00000000000`+strconv.Itoa(i)+`"}`) + } + return &model.Inbound{ + Tag: contractTag, Remark: remark, Enable: enable, Port: 51001, Protocol: model.VLESS, + Settings: `{"clients":[` + strings.Join(entries, ",") + `],"decryption":"none"}`, + StreamSettings: `{"network":"tcp","security":"none","tcpSettings":{"header":{"type":"none"}}}`, + Sniffing: `{}`, + } +} + +func nodeEmails(t *testing.T) []string { + t.Helper() + ib, ok := nodeRow(t, contractTag) + if !ok { + t.Fatal("node has no contract inbound") + } + clients, err := (&service.InboundService{}).GetClients(ib) + if err != nil { + t.Fatalf("parse node clients: %v", err) + } + emails := make([]string, 0, len(clients)) + for _, c := range clients { + emails = append(emails, c.Email) + } + return emails +} + +// TestMasterNodeContract sends every node call the master makes through the production +// router, once per enrollment scope; UpdatePanel is excluded from node-sync on purpose. +func TestMasterNodeContract(t *testing.T) { + for _, scope := range []string{model.ApiScopeAdmin, model.ApiScopeNodeSync} { + t.Run(scope, func(t *testing.T) { + node := startContractNode(t) + master := node.masterWithToken(t, scope) + ctx := context.Background() + + cells := []struct { + name string + covers []string + run func() error + check func(t *testing.T) + }{ + {"AddInbound creates the inbound with its clients", []string{"AddInbound"}, func() error { + return master.AddInbound(ctx, masterInbound("added", true, "c0", "c1")) + }, func(t *testing.T) { + if got := nodeEmails(t); strings.Join(got, ",") != "c0,c1" { + t.Fatalf("node clients = %v, want c0,c1", got) + } + }}, + {"UpdateInbound applies remark, clients and enable", []string{"UpdateInbound", "AddUser", "RemoveUser", "ReconcileInbound"}, func() error { + ib := masterInbound("updated", false, "c0", "c1", "c2") + if err := master.AddUser(ctx, ib, nil); err != nil { + return err + } + if err := master.RemoveUser(ctx, ib, ""); err != nil { + return err + } + if _, err := master.ReconcileInbound(ctx, ib, true); err != nil { + return err + } + return master.UpdateInbound(ctx, ib, ib) + }, func(t *testing.T) { + ib, _ := nodeRow(t, contractTag) + if ib.Remark != "updated" || ib.Enable { + t.Fatalf("node remark=%q enable=%v, want updated/false", ib.Remark, ib.Enable) + } + if got := nodeEmails(t); strings.Join(got, ",") != "c0,c1,c2" { + t.Fatalf("node clients = %v, want c0,c1,c2", got) + } + }}, + {"SetInboundSubSortIndex reaches the node", []string{"SetInboundSubSortIndex"}, func() error { + return master.SetInboundSubSortIndex(ctx, masterInbound("updated", false), 7) + }, func(t *testing.T) { + if ib, _ := nodeRow(t, contractTag); ib.SubSortIndex != 7 { + t.Fatalf("node subSortIndex = %d, want 7", ib.SubSortIndex) + } + }}, + {"AddClient attaches one client", []string{"AddClient"}, func() error { + return master.AddClient(ctx, masterInbound("updated", false), model.Client{ + Email: "c3", ID: "0b6d5c2e-7c1a-4f4e-9d3b-000000000003", SubID: "s-c3", Enable: true, + }) + }, func(t *testing.T) { + if got := nodeEmails(t); !strings.Contains(strings.Join(got, ","), "c3") { + t.Fatalf("node clients = %v, want c3 among them", got) + } + }}, + {"UpdateUser changes the client's limits", []string{"UpdateUser"}, func() error { + return master.UpdateUser(ctx, masterInbound("updated", false), "c3", model.Client{ + Email: "c3", ID: "0b6d5c2e-7c1a-4f4e-9d3b-000000000003", SubID: "s-c3", Enable: true, TotalGB: 5 << 30, + }) + }, func(t *testing.T) { + var ct xray.ClientTraffic + database.GetDB().Where("email = ?", "c3").First(&ct) + if ct.Total != 5<<30 { + t.Fatalf("node c3 total = %d, want %d", ct.Total, int64(5<<30)) + } + }}, + {"ResetClientTraffic zeroes one client", []string{"ResetClientTraffic"}, func() error { + seedNodeTraffic(t, "c0") + return master.ResetClientTraffic(ctx, nil, "c0") + }, func(t *testing.T) { + if u := nodeTraffic(t, "c0"); u != 0 { + t.Fatalf("node c0 usage = %d, want 0", u) + } + }}, + {"ResetClientTraffics zeroes several clients", []string{"ResetClientTraffics"}, func() error { + seedNodeTraffic(t, "c1", "c2") + return master.ResetClientTraffics(ctx, []string{"c1", "c2"}) + }, func(t *testing.T) { + if u := nodeTraffic(t, "c1") + nodeTraffic(t, "c2"); u != 0 { + t.Fatalf("node c1+c2 usage = %d, want 0", u) + } + }}, + {"ResetInboundTraffic zeroes the inbound", []string{"ResetInboundTraffic"}, func() error { + seedNodeTraffic(t) + return master.ResetInboundTraffic(ctx, masterInbound("updated", false)) + }, func(t *testing.T) { + if ib, _ := nodeRow(t, contractTag); ib.Up+ib.Down != 0 { + t.Fatalf("node inbound usage = %d, want 0", ib.Up+ib.Down) + } + }}, + {"ResetAllTraffics zeroes every inbound's counters", []string{"ResetAllTraffics"}, func() error { + seedNodeTraffic(t) + return master.ResetAllTraffics(ctx) + }, func(t *testing.T) { + if ib, _ := nodeRow(t, contractTag); ib.Up+ib.Down != 0 { + t.Fatalf("node inbound usage = %d, want 0", ib.Up+ib.Down) + } + }}, + {"FetchTrafficSnapshot reads every part of the snapshot", []string{"FetchTrafficSnapshot"}, func() error { + _, err := master.FetchTrafficSnapshot(ctx) + return err + }, nil}, + {"PushGlobalClientTraffics is accepted", []string{"PushGlobalClientTraffics"}, func() error { + return master.PushGlobalClientTraffics(ctx, "master-guid", []*xray.ClientTraffic{{Email: "c0", Up: 1, Down: 2}}) + }, nil}, + {"client IP sync is accepted both ways", []string{"FetchAllClientIps", "PushAllClientIps", "FetchClientIpsByGuid"}, func() error { + ips, err := master.FetchAllClientIps(ctx) + if err != nil { + return err + } + if err := master.PushAllClientIps(ctx, ips); err != nil { + return err + } + _, err = master.FetchClientIpsByGuid(ctx) + return err + }, nil}, + {"host groups, descendants and web cert files are readable", []string{"FetchHostGroups", "GetDescendants", "GetWebCertFiles", "ListInboundOptions", "ListRemoteTags"}, func() error { + if _, err := master.FetchHostGroups(ctx); err != nil { + return err + } + if _, err := master.GetDescendants(ctx); err != nil { + return err + } + if _, err := master.GetWebCertFiles(ctx); err != nil { + return err + } + if _, err := master.ListInboundOptions(ctx); err != nil { + return err + } + _, err := master.ListRemoteTags(ctx) + return err + }, nil}, + {"RestartXray is accepted by the node", []string{"RestartXray"}, func() error { + // No core binary here: only the node's own restart failure may come back. + if err := master.RestartXray(ctx); err != nil && !strings.Contains(err.Error(), "rebooting the Xray") { + return err + } + return nil + }, nil}, + {"DeleteUser detaches the client from the inbound", []string{"DeleteUser"}, func() error { + return master.DeleteUser(ctx, masterInbound("updated", false), "c3") + }, func(t *testing.T) { + if got := nodeEmails(t); strings.Contains(strings.Join(got, ","), "c3") { + t.Fatalf("node clients = %v, want c3 gone", got) + } + }}, + {"DeleteClient removes the client everywhere", []string{"DeleteClient"}, func() error { + return master.DeleteClient(ctx, "c2") + }, func(t *testing.T) { + if got := nodeEmails(t); strings.Contains(strings.Join(got, ","), "c2") { + t.Fatalf("node clients = %v, want c2 gone", got) + } + }}, + {"DelInbound removes the inbound", []string{"DelInbound"}, func() error { + return master.DelInbound(ctx, masterInbound("updated", false)) + }, func(t *testing.T) { + if _, ok := nodeRow(t, contractTag); ok { + t.Fatal("node still has the inbound") + } + }}, + } + covered := map[string]bool{} + for _, c := range cells { + for _, m := range c.covers { + covered[m] = true + } + } + assertEveryRemoteCallCovered(t, covered) + for _, c := range cells { + t.Run(c.name, func(t *testing.T) { + node.takeRefused() + if err := c.run(); err != nil { + t.Fatalf("master call failed: %v", err) + } + if refused := node.takeRefused(); len(refused) != 0 { + t.Fatalf("node refused master requests: %v", refused) + } + if c.check != nil { + c.check(t) + } + }) + } + }) + } +} + +// Remote methods that never reach the node, or that this table must not run. +var remoteMethodsOutsideContract = map[string]string{ + "Name": "local label", + "RecordAdoptedInbound": "local fingerprint bookkeeping", + "AdoptInboundAlias": "local alias bookkeeping", + "AdoptedInboundAliases": "local alias bookkeeping", + "AdvancePushedInbound": "local fingerprint bookkeeping", + "UpdatePanel": "replaces the node binary; node-sync is denied it on purpose (#6201)", +} + +// A Remote method with no cell is how activeInbounds and bulkResetTraffic +// drifted out of the node-sync allowlist unnoticed. +func assertEveryRemoteCallCovered(t *testing.T, covered map[string]bool) { + t.Helper() + rt := reflect.TypeOf(&runtime.Remote{}) + for i := 0; i < rt.NumMethod(); i++ { + name := rt.Method(i).Name + if _, skip := remoteMethodsOutsideContract[name]; skip { + continue + } + if !covered[name] { + t.Errorf("runtime.Remote.%s has no cell in TestMasterNodeContract", name) + } + } +} diff --git a/internal/web/runtime/reconcile_skip_test.go b/internal/web/runtime/reconcile_skip_test.go index 5be854c6b..7bff47238 100644 --- a/internal/web/runtime/reconcile_skip_test.go +++ b/internal/web/runtime/reconcile_skip_test.go @@ -17,7 +17,8 @@ import ( func TestReconcileInbound_SkipsUnchanged(t *testing.T) { var pushes atomic.Int32 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - if r.Method == http.MethodPost && strings.Contains(r.URL.Path, "/panel/api/inbounds/update/") { + if r.Method == http.MethodPost && (strings.Contains(r.URL.Path, "/panel/api/inbounds/update/") || + strings.Contains(r.URL.Path, "/panel/api/inbounds/add")) { pushes.Add(1) } w.Header().Set("Content-Type", "application/json") @@ -283,7 +284,7 @@ func TestDelInboundDropsReconcileFingerprint(t *testing.T) { ib := &model.Inbound{Tag: "in-del", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`} r.cacheSet(ib.Tag, 7) - if pushed, err := r.ReconcileInbound(context.Background(), ib, false); err != nil || !pushed { + if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || !pushed { t.Fatalf("initial reconcile: pushed=%v err=%v, want push", pushed, err) } if err := r.DelInbound(context.Background(), ib); err != nil { @@ -317,3 +318,36 @@ func TestUpdateInboundFallbackAddSeedsReconcileFingerprint(t *testing.T) { t.Fatalf("reconcile sent %d full inbound updates, want 0", got) } } + +// An inbound deleted on the node must be re-created by the next reconcile; a +// cached tag→id from before the delete used to send update/ forever. +func TestReconcileInbound_RecreatesInboundTheNodeLost(t *testing.T) { + var adds, staleUpdates atomic.Int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + switch { + case strings.Contains(r.URL.Path, "/panel/api/inbounds/list"): + _, _ = w.Write([]byte(`{"success":true,"obj":[]}`)) + case strings.Contains(r.URL.Path, "/panel/api/inbounds/update/"): + staleUpdates.Add(1) + _, _ = w.Write([]byte(`{"success":false,"msg":"record not found"}`)) + case strings.Contains(r.URL.Path, "/panel/api/inbounds/add"): + adds.Add(1) + _, _ = w.Write([]byte(`{"success":true,"obj":{"id":9,"tag":"in-1"}}`)) + default: + _, _ = w.Write([]byte(`{"success":true}`)) + } + })) + defer srv.Close() + + r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil) + ib := &model.Inbound{Tag: "n1-in-1", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`} + r.cacheSet("in-1", 7) + + if pushed, err := r.ReconcileInbound(context.Background(), ib, false); err != nil || !pushed { + t.Fatalf("reconcile of a lost inbound: pushed=%v err=%v, want a re-create", pushed, err) + } + if staleUpdates.Load() != 0 || adds.Load() != 1 { + t.Fatalf("updates to the stale id=%d adds=%d, want 0 and 1", staleUpdates.Load(), adds.Load()) + } +} diff --git a/internal/web/runtime/remote.go b/internal/web/runtime/remote.go index 4d93011c3..9044c73e4 100644 --- a/internal/web/runtime/remote.go +++ b/internal/web/runtime/remote.go @@ -237,6 +237,7 @@ func (r *Remote) do(ctx context.Context, method, path string, body any) (*envelo req.Header.Set("Authorization", "Bearer "+token) } req.Header.Set("Accept", "application/json") + req.Header.Set(wirecodec.MasterPushHeader, "1") if contentType != "" { req.Header.Set("Content-Type", contentType) } @@ -359,6 +360,15 @@ func (r *Remote) cacheDel(tag string) { delete(r.pushedFP, tag) } +// forgetTag drops every tag form cacheGetTag would match, once the node reports +// none of them, so the next resolve re-reads the node instead of a deleted id. +func (r *Remote) forgetTag(tag string) { + prefix := nodeInboundTagPrefix(r.node.Id) + bare := strings.TrimPrefix(tag, prefix) + r.cacheDel(bare) + r.cacheDel(prefix + bare) +} + func (r *Remote) ListRemoteTags(ctx context.Context) ([]string, error) { if err := r.refreshRemoteIDs(ctx); err != nil { return nil, err @@ -494,6 +504,8 @@ func (r *Remote) ReconcileInbound(ctx context.Context, ib *model.Inbound, exists if ok && prev == fp { return false, nil } + } else { + r.forgetTag(ib.Tag) } if err := r.UpdateInbound(ctx, ib, ib); err != nil { return false, err