mirror of
https://github.com/MHSanaei/3x-ui.git
synced 2026-10-06 22:22:08 +03:00
Merge branch 'main' into fix/node-empty-snapshot-client-deletion
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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`.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 (
|
||||
<ConfigProvider theme={antdThemeConfig}>
|
||||
{messageContextHolder}
|
||||
<div
|
||||
className={`command-palette-backdrop ${themeModeClass}`}
|
||||
role="presentation"
|
||||
|
||||
@@ -21,6 +21,7 @@ import { HttpUtil } from '@/utils';
|
||||
|
||||
export default function TuicFields() {
|
||||
const { t } = useTranslation();
|
||||
const [messageApi, messageContextHolder] = message.useMessage();
|
||||
const { control, setValue } = useFormContext();
|
||||
const [loadingPanelCert, setLoadingPanelCert] = useState(false);
|
||||
|
||||
@@ -36,7 +37,7 @@ export default function TuicFields() {
|
||||
const autofillFromSni = () => {
|
||||
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}
|
||||
<Form.Item label={t('pages.inbounds.publicKey')}>
|
||||
<AutoComplete
|
||||
value={certificate}
|
||||
|
||||
@@ -28,6 +28,7 @@ export default function HappSettingsContent({
|
||||
remoteSourceBadge,
|
||||
}: HappSettingsContentProps) {
|
||||
const { t } = useTranslation();
|
||||
const [messageApi, messageContextHolder] = message.useMessage();
|
||||
// Generator choices stay local until Apply updates the draft; page Save persists it.
|
||||
const [selectedPreset, setSelectedPreset] = useState<string>('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}
|
||||
<SettingListItem
|
||||
paddings="small"
|
||||
title={t('pages.settings.subHappAutoDetect')}
|
||||
|
||||
@@ -45,6 +45,7 @@ export default function RoutingTab({
|
||||
isMobile,
|
||||
}: RoutingTabProps) {
|
||||
const { t } = useTranslation();
|
||||
const [messageApi, messageContextHolder] = message.useMessage();
|
||||
const [modal, modalContextHolder] = Modal.useModal();
|
||||
const [ruleModalOpen, setRuleModalOpen] = useState(false);
|
||||
const [editingRule, setEditingRule] = useState<RoutingRule | null>(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}
|
||||
<Tabs
|
||||
defaultActiveKey="basic"
|
||||
items={[
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { useState } from 'react';
|
||||
import { describe, expect, it, vi } from 'vitest';
|
||||
import { act, fireEvent, screen, within } from '@testing-library/react';
|
||||
import { act, cleanup, fireEvent, screen, within } from '@testing-library/react';
|
||||
import { EditorView } from 'codemirror';
|
||||
|
||||
import { AllSetting } from '@/models/setting';
|
||||
@@ -319,4 +319,19 @@ describe('Happ routing editor', () => {
|
||||
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',
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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',
|
||||
|
||||
@@ -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 = /(?<![\w.$])message\.(success|error|warning|info|loading|open)\(/;
|
||||
|
||||
function sourceFiles(dir: string): string[] {
|
||||
return readdirSync(dir).flatMap((name) => {
|
||||
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([]);
|
||||
});
|
||||
});
|
||||
@@ -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])
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -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.
|
||||
|
||||
@@ -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: {}},
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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/<gone id> 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())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user