feat(inbounds): deploy AmneziaWG, TUIC and MTProto inbounds to nodes

The node gate assumed a node-assigned sidecar row would never converge,
because the master's reconcile loops only read node_id IS NULL rows. But
a pushed row is local on the node's own panel, whose AmneziaWG, TUIC and
mtg loops run it like any other; node-adopted rows of these protocols
already worked. Only creating or cloning them from the master was blocked.

Open the three protocols on both lists and fix what assumed the master's
host for a node row:
- A node older than the release that introduced the protocol (MTProto
  v3.5.0, AmneziaWG v3.7.0, TUIC v3.8.0) would hand it to Xray as-is, so
  add and protocol-changing update refuse it, and a node that has not
  reported its version yet. Dev builds ("dev+<sha>") track main and pass.
- MTProto's routeXrayPort is a loopback port on the host running mtg.
  The master no longer allocates one, nor forwards a cloned source's, for
  a node row; the node allocates its own and node sync adopts it back.
- AmneziaWG forwardedPorts were checked against the master's inbounds,
  web port and id-derived relay ports. A node row is now checked against
  its own node's inbounds; the node re-checks what only it knows.
  UpdateInbound restores the stored nodeId before that check.

Verified on a docker master+node pair: AWG, routed MTProto and TUIC
created and cloned from the master start on the node (interface, mtg,
tuic-server); an AWG client added and an MTProto client edited on the
master apply on the node; the node-chosen egress port survives edits.

Closes #6306
This commit is contained in:
MHSanaei
2026-10-03 00:26:46 +02:00
parent 9b957b969b
commit ed31ee432c
12 changed files with 387 additions and 75 deletions
+4 -4
View File
@@ -433,7 +433,7 @@ func (s *ClientService) AddInboundClient(inboundSvc *InboundService, data *model
var portCtx portConflictContext
if oldInbound.Protocol == model.AmneziaWG {
portCtx, err = inboundSvc.loadPortConflictContext(database.GetDB())
portCtx, err = inboundSvc.loadPortConflictContext(database.GetDB(), oldInbound.NodeID)
if err != nil {
return false, err
}
@@ -547,7 +547,7 @@ func (s *ClientService) AddInboundClient(inboundSvc *InboundService, data *model
}
}
if oldInbound.Protocol == model.AmneziaWG {
txPortCtx, pErr := inboundSvc.loadPortConflictContext(tx)
txPortCtx, pErr := inboundSvc.loadPortConflictContext(tx, oldInbound.NodeID)
if pErr != nil {
return pErr
}
@@ -784,7 +784,7 @@ func (s *ClientService) UpdateInboundClient(inboundSvc *InboundService, data *mo
}
}
if oldInbound.Protocol == model.AmneziaWG {
portCtx, err := inboundSvc.loadPortConflictContext(database.GetDB())
portCtx, err := inboundSvc.loadPortConflictContext(database.GetDB(), oldInbound.NodeID)
if err != nil {
return false, err
}
@@ -921,7 +921,7 @@ func (s *ClientService) UpdateInboundClient(inboundSvc *InboundService, data *mo
// Same re-check-inside-the-writer rule as AddInboundClient (#6225):
// the pre-tx pass can race a concurrent writer on another inbound.
if oldInbound.Protocol == model.AmneziaWG {
txPortCtx, pErr := inboundSvc.loadPortConflictContext(tx)
txPortCtx, pErr := inboundSvc.loadPortConflictContext(tx, oldInbound.NodeID)
if pErr != nil {
return pErr
}
+26 -7
View File
@@ -1086,6 +1086,21 @@ func (s *InboundService) normalizeMtprotoXrayPort(inbound *model.Inbound, oldSet
// Prefer the already-stored port (carried across edits), then any value the
// client sent, then allocate a fresh one.
port := parseRouteXrayPort(oldSettings)
if inbound.NodeID != nil {
// The port is free or taken on the node's host, not here: the node's own
// panel allocates it, and node sync brings its choice back as oldSettings.
if port <= 0 {
delete(parsed, "routeXrayPort")
} else {
parsed["routeXrayPort"] = port
}
bs, err := json.MarshalIndent(parsed, "", " ")
if err != nil {
return common.NewError("mtproto: could not persist the Xray egress port:", err)
}
inbound.Settings = string(bs)
return nil
}
if port <= 0 {
port = settingsRouteXrayPort(parsed)
}
@@ -1136,8 +1151,10 @@ func (s *InboundService) AddInbound(inbound *model.Inbound) (*model.Inbound, boo
if err := s.normalizeAmneziaWGSettings(inbound, ""); err != nil {
return inbound, false, err
}
if inbound.NodeID != nil && !isNodeEligibleProtocol(inbound.Protocol) {
return inbound, false, common.NewErrorf("%s inbounds cannot be assigned to a node", inbound.Protocol)
if inbound.NodeID != nil {
if err := checkNodeCanHostProtocol(database.GetDB(), *inbound.NodeID, inbound.Protocol); err != nil {
return inbound, false, err
}
}
inbound.SubSortIndex = normalizeSubSortIndex(inbound.SubSortIndex)
if err := normalizeInboundShareAddressStrict(inbound); err != nil {
@@ -1752,6 +1769,9 @@ func (s *InboundService) UpdateInbound(inbound *model.Inbound) (*model.Inbound,
if err != nil {
return inbound, false, err
}
// Restore the stored NodeID before any host-scoped check so a node inbound
// stays scoped to its own node (the payload's nodeId is unreliable, often absent).
inbound.NodeID = oldInbound.NodeID
if err := s.normalizeAmneziaWGSettings(inbound, oldInbound.Settings); err != nil {
return inbound, false, err
}
@@ -1766,13 +1786,12 @@ func (s *InboundService) UpdateInbound(inbound *model.Inbound) (*model.Inbound,
}
}
}
// Restore the stored NodeID before the port-conflict check so a node inbound
// stays scoped to its own node (the payload's nodeId is unreliable, often absent).
inbound.NodeID = oldInbound.NodeID
// The node assignment is the stored one, so only a protocol change can
// introduce one; a row adopted from a node keeps the protocol it arrived with.
if inbound.NodeID != nil && inbound.Protocol != oldInbound.Protocol && !isNodeEligibleProtocol(inbound.Protocol) {
return inbound, false, common.NewErrorf("%s inbounds cannot be assigned to a node", inbound.Protocol)
if inbound.NodeID != nil && inbound.Protocol != oldInbound.Protocol {
if err := checkNodeCanHostProtocol(database.GetDB(), *inbound.NodeID, inbound.Protocol); err != nil {
return inbound, false, err
}
}
// Capture the pre-edit protocol and routing state before oldInbound is
+21 -13
View File
@@ -262,7 +262,7 @@ func (s *InboundService) normalizeAmneziaWGSettings(inbound *model.Inbound, oldS
}
}
portCtx, err := s.loadPortConflictContext(database.GetDB())
portCtx, err := s.loadPortConflictContext(database.GetDB(), inbound.NodeID)
if err != nil {
return err
}
@@ -303,23 +303,31 @@ func (s *InboundService) normalizeAmneziaWGSettings(inbound *model.Inbound, oldS
return nil
}
// portConflictContext caches what checkForwardedPortsConflict needs — the panel's
// own port and this host's enabled rows — so one save costs one query, not N.
// portConflictContext caches what checkForwardedPortsConflict needs about the
// host a forward listener binds on, so one save costs one query, not N.
type portConflictContext struct {
webPort int
inbounds []*model.Inbound
// onNode: the host is a node, whose web port and relay ports (derived from
// its own inbound ids) this panel does not know; the node re-checks both.
onNode bool
}
// loadPortConflictContext loads the panel's own port and every enabled inbound
// hosted on THIS panel: a node-hosted one listens on that node's host, not here.
func (s *InboundService) loadPortConflictContext(db *gorm.DB) (portConflictContext, error) {
// loadPortConflictContext loads every enabled inbound hosted where nodeID's rows
// run -- this panel for nil, else that node -- plus this panel's own port.
func (s *InboundService) loadPortConflictContext(db *gorm.DB, nodeID *int) (portConflictContext, error) {
var ctx portConflictContext
if webPort, err := (&SettingService{}).GetPort(); err == nil {
ctx.webPort = webPort
q := db.Model(model.Inbound{}).Where("enable = ?", true)
if nodeID != nil {
ctx.onNode = true
q = q.Where("node_id = ?", *nodeID)
} else {
if webPort, err := (&SettingService{}).GetPort(); err == nil {
ctx.webPort = webPort
}
q = q.Where("node_id IS NULL")
}
err := db.Model(model.Inbound{}).
Where("enable = ? AND node_id IS NULL", true).
Find(&ctx.inbounds).Error
err := q.Find(&ctx.inbounds).Error
return ctx, err
}
@@ -340,7 +348,7 @@ func (s *InboundService) checkAmneziaWGForwardedPorts(db *gorm.DB, settings stri
if err := json.Unmarshal([]byte(settings), &parsed); err != nil {
return nil
}
ctx, err := s.loadPortConflictContext(db)
ctx, err := s.loadPortConflictContext(db, nil)
if err != nil {
return err
}
@@ -372,7 +380,7 @@ func (s *InboundService) checkForwardedPortsConflict(ctx portConflictContext, fo
}
return fmt.Sprintf("inbound '%s' (#%d, port %d)", name, ib.Id, ib.Port)
}
if ib.Protocol != model.AmneziaWG {
if ctx.onNode || ib.Protocol != model.AmneziaWG {
continue
}
socksPort := amneziawgnet.SOCKSPortForInbound(ib.Id)
@@ -29,7 +29,7 @@ var awgTestPrivateKey, awgTestPublicKey = func() (string, string) {
func TestCheckForwardedPortsConflict_EmptySpecNoConflict(t *testing.T) {
setupConflictDB(t)
svc := &InboundService{}
ctx, err := svc.loadPortConflictContext(database.GetDB())
ctx, err := svc.loadPortConflictContext(database.GetDB(), nil)
if err != nil {
t.Fatalf("loadPortConflictContext: %v", err)
}
@@ -41,7 +41,7 @@ func TestCheckForwardedPortsConflict_EmptySpecNoConflict(t *testing.T) {
func TestCheckForwardedPortsConflict_CollidesWithPanelPort(t *testing.T) {
setupConflictDB(t)
svc := &InboundService{}
ctx, err := svc.loadPortConflictContext(database.GetDB())
ctx, err := svc.loadPortConflictContext(database.GetDB(), nil)
if err != nil {
t.Fatalf("loadPortConflictContext: %v", err)
}
@@ -58,7 +58,7 @@ func TestCheckForwardedPortsConflict_CollidesWithEnabledInboundPort(t *testing.T
seedInboundConflict(t, "vless-8080", "0.0.0.0", 8080, model.VLESS, `{"network":"tcp"}`, `{}`)
svc := &InboundService{}
ctx, err := svc.loadPortConflictContext(database.GetDB())
ctx, err := svc.loadPortConflictContext(database.GetDB(), nil)
if err != nil {
t.Fatalf("loadPortConflictContext: %v", err)
}
@@ -76,7 +76,7 @@ func TestCheckForwardedPortsConflict_IgnoresDisabledInboundPort(t *testing.T) {
}
svc := &InboundService{}
ctx, err := svc.loadPortConflictContext(database.GetDB())
ctx, err := svc.loadPortConflictContext(database.GetDB(), nil)
if err != nil {
t.Fatalf("loadPortConflictContext: %v", err)
}
@@ -90,7 +90,7 @@ func TestCheckForwardedPortsConflict_NoCollisionWhenPortsDontOverlap(t *testing.
seedInboundConflict(t, "vless-8080", "0.0.0.0", 8080, model.VLESS, `{"network":"tcp"}`, `{}`)
svc := &InboundService{}
ctx, err := svc.loadPortConflictContext(database.GetDB())
ctx, err := svc.loadPortConflictContext(database.GetDB(), nil)
if err != nil {
t.Fatalf("loadPortConflictContext: %v", err)
}
@@ -110,7 +110,7 @@ func TestCheckForwardedPortsConflict_IgnoresPortOnDifferentNode(t *testing.T) {
seedInboundConflictNode(t, "node1-8080", "0.0.0.0", 8080, model.VLESS, `{"network":"tcp"}`, `{}`, new(1))
svc := &InboundService{}
ctx, err := svc.loadPortConflictContext(database.GetDB())
ctx, err := svc.loadPortConflictContext(database.GetDB(), nil)
if err != nil {
t.Fatalf("loadPortConflictContext: %v", err)
}
@@ -320,7 +320,7 @@ func TestGetAmneziaWGLogs_ClampsCountAndFiltersEvents(t *testing.T) {
func TestCheckForwardedPortsConflict_RejectsSpecOverCap(t *testing.T) {
setupConflictDB(t)
svc := &InboundService{}
ctx, err := svc.loadPortConflictContext(database.GetDB())
ctx, err := svc.loadPortConflictContext(database.GetDB(), nil)
if err != nil {
t.Fatalf("loadPortConflictContext: %v", err)
}
@@ -337,7 +337,7 @@ func TestCheckForwardedPortsConflict_RejectsSpecOverCap(t *testing.T) {
func TestCheckForwardedPortsConflict_AcceptsSpecExactlyAtCap(t *testing.T) {
setupConflictDB(t)
svc := &InboundService{}
ctx, err := svc.loadPortConflictContext(database.GetDB())
ctx, err := svc.loadPortConflictContext(database.GetDB(), nil)
if err != nil {
t.Fatalf("loadPortConflictContext: %v", err)
}
@@ -361,7 +361,7 @@ func TestCheckForwardedPortsConflict_CollidesWithAmneziawgnetSocksPort(t *testin
relayPort := amneziawgnet.SOCKSPortForInbound(awgInbound.Id)
svc := &InboundService{}
ctx, err := svc.loadPortConflictContext(database.GetDB())
ctx, err := svc.loadPortConflictContext(database.GetDB(), nil)
if err != nil {
t.Fatalf("loadPortConflictContext: %v", err)
}
@@ -42,12 +42,12 @@ func TestUpdateInbound_NodeMtprotoShareAddrIsEditable(t *testing.T) {
}
}
// Converting a node inbound to a protocol the master's sidecars only reconcile
// for local rows is still refused: those loops query node_id IS NULL.
// Converting a node inbound to a protocol that never lives on a node is still
// refused, and so is one the node's panel is too old to run.
func TestUpdateInbound_RejectsProtocolChangeToNodeIneligible(t *testing.T) {
setupConflictDB(t)
nodeID := 6
seedNodeRow(t, database.GetDB(), &model.Node{Id: nodeID, Name: "n6", Address: "127.0.0.1", Port: 2096, ApiToken: "tok", Enable: true})
seedNodeRow(t, database.GetDB(), &model.Node{Id: nodeID, Name: "n6", Address: "127.0.0.1", Port: 2096, ApiToken: "tok", Enable: true, PanelVersion: "v3.4.1"})
seedInboundConflictNode(t, "vless-node", "127.0.0.1", 4064, model.VLESS,
`{"network":"tcp","security":"none"}`,
`{"clients":[{"id":"11111111-2222-4333-8444-555555555555","email":"vn-c","enable":true}],"decryption":"none"}`, &nodeID)
@@ -57,11 +57,27 @@ func TestUpdateInbound_RejectsProtocolChangeToNodeIneligible(t *testing.T) {
t.Fatalf("read seeded row: %v", err)
}
update := existing
update.Protocol = model.MTProto
update.Settings = `{"clients":[{"email":"vn-c","enable":true,"secret":"ee0123456789abcdef0123456789abcdef"}]}`
if _, _, err := (&InboundService{}).UpdateInbound(&update); err == nil ||
!strings.Contains(err.Error(), "cannot be assigned to a node") {
t.Fatalf("err = %v, want a node-eligibility refusal", err)
cases := []struct {
name string
protocol model.Protocol
settings string
wantErr string
}{
{"panel-local protocol", model.HTTP, `{"accounts":[]}`, "cannot be assigned to a node"},
{
"protocol newer than the node", model.MTProto,
`{"clients":[{"email":"vn-c","enable":true,"secret":"ee0123456789abcdef0123456789abcdef"}]}`, "need v3.5.0 or newer",
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
update := existing
update.Protocol = tc.protocol
update.Settings = tc.settings
if _, _, err := (&InboundService{}).UpdateInbound(&update); err == nil ||
!strings.Contains(err.Error(), tc.wantErr) {
t.Fatalf("err = %v, want a refusal mentioning %q", err, tc.wantErr)
}
})
}
}
@@ -0,0 +1,244 @@
package service
import (
"encoding/json"
"fmt"
"strings"
"testing"
"github.com/mhsanaei/3x-ui/v3/internal/amneziawgnet"
"github.com/mhsanaei/3x-ui/v3/internal/database"
"github.com/mhsanaei/3x-ui/v3/internal/database/model"
)
func seedVersionedNode(t *testing.T, panelVersion string) *model.Node {
t.Helper()
node := &model.Node{
Name: "n-" + panelVersion, Address: "127.0.0.1", Port: 2096, Scheme: "https",
Enable: true, Status: "online", PanelVersion: panelVersion,
}
seedNodeRow(t, database.GetDB(), node)
return node
}
func sidecarNodeInbound(t *testing.T, protocol model.Protocol, tag string, nodeID int) *model.Inbound {
t.Helper()
ib := &model.Inbound{Tag: tag, Enable: true, Listen: "0.0.0.0", Port: 44300, Protocol: protocol, NodeID: &nodeID}
switch protocol {
case model.AmneziaWG:
ib.Settings = awgRelayWindowSettings(t, tag)
case model.TUIC:
ib.Settings = `{"clients":[{"id":"8a4f0c7e-1d2b-4c3a-9e5f-6a7b8c9d0e1f","password":"pass","email":"` + tag + `@tuic","enable":true}]}`
case model.MTProto:
ib.Settings = `{"clients":[{"email":"` + tag + `@mt","enable":true,"secret":"ee0123456789abcdef0123456789abcdef"}]}`
default:
t.Fatalf("no fixture for %s", protocol)
}
return ib
}
// The node's own panel runs the sidecar for a pushed row, so a sidecar protocol
// is deployable to any node new enough to know it -- the release that added it.
func TestAddInbound_SidecarProtocolDeploysToANodeAtItsFirstRelease(t *testing.T) {
cases := []struct {
protocol model.Protocol
firstRelease string
}{
{model.MTProto, "v3.5.0"},
{model.AmneziaWG, "v3.7.0"},
{model.TUIC, "v3.8.0"},
}
for _, tc := range cases {
t.Run(string(tc.protocol), func(t *testing.T) {
setupConflictDB(t)
node := seedVersionedNode(t, tc.firstRelease)
created, _, err := (&InboundService{}).AddInbound(sidecarNodeInbound(t, tc.protocol, "side-"+string(tc.protocol), node.Id))
if err != nil {
t.Fatalf("AddInbound(%s on a %s node): %v", tc.protocol, tc.firstRelease, err)
}
var stored model.Inbound
if err := database.GetDB().First(&stored, created.Id).Error; err != nil {
t.Fatalf("read created row: %v", err)
}
if stored.NodeID == nil || *stored.NodeID != node.Id {
t.Fatalf("stored nodeId = %v, want %d", stored.NodeID, node.Id)
}
})
}
}
// A node that predates a protocol stores it as an Xray inbound Xray cannot load,
// so the assignment is refused while the version is too old or still unknown.
func TestAddInbound_SidecarProtocolNodeVersionGate(t *testing.T) {
cases := []struct {
name string
panelVersion string
wantErr string
}{
{"one release too old", "v3.6.9", "v3.7.0"},
{"version not reported yet", "", "has not reported its panel version"},
{"dev build", "dev+1d1128cf", ""},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
setupConflictDB(t)
node := seedVersionedNode(t, tc.panelVersion)
_, _, err := (&InboundService{}).AddInbound(sidecarNodeInbound(t, model.AmneziaWG, "awg-gate", node.Id))
if tc.wantErr == "" {
if err != nil {
t.Fatalf("a dev build tracks main and knows every protocol; got %v", err)
}
return
}
if err == nil || !strings.Contains(err.Error(), tc.wantErr) {
t.Fatalf("err = %v, want a refusal mentioning %q", err, tc.wantErr)
}
})
}
}
// The egress port is a loopback port on the host running mtg; one picked or
// copied on the master can be taken on the node, so the node must allocate it.
func TestAddInbound_NodeMtprotoLeavesTheEgressPortToTheNode(t *testing.T) {
setupConflictDB(t)
node := seedVersionedNode(t, "v3.8.0")
ib := sidecarNodeInbound(t, model.MTProto, "mt-routed", node.Id)
ib.Settings = `{"routeThroughXray":true,"routeXrayPort":4444,` + strings.TrimPrefix(ib.Settings, "{")
created, _, err := (&InboundService{}).AddInbound(ib)
if err != nil {
t.Fatalf("AddInbound: %v", err)
}
var stored model.Inbound
if err := database.GetDB().First(&stored, created.Id).Error; err != nil {
t.Fatalf("read created row: %v", err)
}
var settings map[string]any
if err := json.Unmarshal([]byte(stored.Settings), &settings); err != nil {
t.Fatalf("decode settings: %v", err)
}
if port, ok := settings["routeXrayPort"]; ok {
t.Fatalf("routeXrayPort = %v reached the node row; the node must allocate its own", port)
}
if settings["routeThroughXray"] != true {
t.Fatalf("routeThroughXray = %v, want true kept", settings["routeThroughXray"])
}
}
// A peer's forward listener binds on the host running the AmneziaWG row, so a
// node row is checked against that node's inbounds, never the master's.
func TestAddInbound_NodeAmneziaWGForwardedPortsCheckTheNodesHost(t *testing.T) {
cases := []struct {
name string
onNode bool
wantErr bool
}{
{"port of a master-local inbound", false, false},
{"port of an inbound on the same node", true, true},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
setupConflictDB(t)
node := seedVersionedNode(t, "v3.8.0")
var owner *int
if tc.onNode {
owner = &node.Id
}
seedInboundConflictNode(t, "holder", "0.0.0.0", 8443, model.VLESS, `{"network":"tcp"}`, `{"clients":[]}`, owner)
ib := sidecarNodeInbound(t, model.AmneziaWG, "awg-fwd", node.Id)
ib.Settings = awgRelayWindowSettingsWithForward(t, "awg-fwd", "8443")
_, _, err := (&InboundService{}).AddInbound(ib)
if !tc.wantErr {
if err != nil {
t.Fatalf("8443 is bound on the master, not on the node; the create must be allowed: %v", err)
}
return
}
if err == nil || !strings.Contains(err.Error(), "forwardedPorts collides with inbound 'holder'") {
t.Fatalf("err = %v, want a forwardedPorts refusal naming 'holder'", err)
}
})
}
}
// awgClientsPayload is the clients-only body the client endpoints take: one
// fresh peer of tag at address, forwarding forwardedPorts.
func awgClientsPayload(t *testing.T, tag, address, forwardedPorts string) string {
t.Helper()
settings := replaceFirst(t, awgRelayWindowSettingsWithForward(t, tag, forwardedPorts),
`"allowedIPs":["10.8.1.2/32"]`, `"allowedIPs":["`+address+`"]`)
var parsed map[string]json.RawMessage
if err := json.Unmarshal([]byte(settings), &parsed); err != nil {
t.Fatalf("decode fixture: %v", err)
}
return `{"clients":` + string(parsed["clients"]) + `}`
}
// seedNodeAmneziaWGBesideLocalHolder seeds a node AmneziaWG row and a local
// inbound on 8443, the port a forward on the node is free to take.
func seedNodeAmneziaWGBesideLocalHolder(t *testing.T) *model.Inbound {
t.Helper()
setupConflictDB(t)
node := seedVersionedNode(t, "v3.8.0")
seedInboundConflict(t, "holder", "0.0.0.0", 8443, model.VLESS, `{"network":"tcp"}`, `{"clients":[]}`)
seedInboundConflictNode(t, "awg-node", "0.0.0.0", 51820, model.AmneziaWG, ``, awgRelayWindowSettings(t, "awg-node"), &node.Id)
var row model.Inbound
if err := database.GetDB().Where("tag = ?", "awg-node").First(&row).Error; err != nil {
t.Fatalf("read seeded row: %v", err)
}
return &row
}
// The edit payload carries no reliable nodeId, so the forward check must use
// the stored one or it judges a node row against the master's own ports.
func TestUpdateInbound_NodeAmneziaWGForwardIgnoresMasterPorts(t *testing.T) {
row := seedNodeAmneziaWGBesideLocalHolder(t)
update := *row
update.NodeID = nil
update.Settings = replaceFirst(t, row.Settings, `"enable":true`, `"enable":true,"forwardedPorts":"8443"`)
if _, _, err := (&InboundService{}).UpdateInbound(&update); err != nil {
t.Fatalf("8443 is bound on the master, not on the node; the edit must be allowed: %v", err)
}
}
func TestAddInboundClient_NodeAmneziaWGForwardIgnoresMasterPorts(t *testing.T) {
row := seedNodeAmneziaWGBesideLocalHolder(t)
data := &model.Inbound{Id: row.Id, Settings: awgClientsPayload(t, "awg-new", "10.8.1.3/32", "8443")}
if _, err := (&ClientService{}).AddInboundClient(&InboundService{}, data); err != nil {
t.Fatalf("8443 is bound on the master, not on the node; the add must be allowed: %v", err)
}
}
func TestUpdateInboundClient_NodeAmneziaWGForwardIgnoresMasterPorts(t *testing.T) {
row := seedNodeAmneziaWGBesideLocalHolder(t)
var stored struct {
Clients []json.RawMessage `json:"clients"`
}
if err := json.Unmarshal([]byte(row.Settings), &stored); err != nil {
t.Fatalf("decode seeded settings: %v", err)
}
edited := replaceFirst(t, string(stored.Clients[0]), `"enable":true`, `"enable":true,"forwardedPorts":"8443"`)
data := &model.Inbound{Id: row.Id, Settings: `{"clients":[` + edited + `]}`}
if _, err := (&ClientService{}).UpdateInboundClient(&InboundService{}, data, "awg-node@relay-window"); err != nil {
t.Fatalf("8443 is bound on the master, not on the node; the edit must be allowed: %v", err)
}
}
// A relay port derives from the inbound id on the panel running it; the node's
// id differs from the master's, so a master-side derivation names a wrong port.
func TestAddInbound_NodeAmneziaWGForwardIgnoresMasterDerivedRelayPorts(t *testing.T) {
row := seedNodeAmneziaWGBesideLocalHolder(t)
ib := sidecarNodeInbound(t, model.AmneziaWG, "awg-fwd", *row.NodeID)
ib.Port = 51821
ib.Settings = awgRelayWindowSettingsWithForward(t, "awg-fwd", fmt.Sprintf("%d", amneziawgnet.SOCKSPortForInbound(row.Id)))
if _, _, err := (&InboundService{}).AddInbound(ib); err != nil {
t.Fatalf("the master's id for a node row derives no port on the node; the create must be allowed: %v", err)
}
}
+42 -7
View File
@@ -2,8 +2,13 @@ package service
import (
"encoding/json"
"strings"
"github.com/mhsanaei/3x-ui/v3/internal/database/model"
"github.com/mhsanaei/3x-ui/v3/internal/util/common"
"github.com/mhsanaei/3x-ui/v3/internal/util/version"
"gorm.io/gorm"
)
// inboundShadowsocksMethod extracts settings.method for Shadowsocks inbounds so
@@ -53,10 +58,8 @@ func inboundCanEnableTlsFlow(protocol, streamSettings, settings string) bool {
}
}
// nodeEligibleProtocols mirrors the frontend's NODE_ELIGIBLE_PROTOCOLS. The
// sidecar-managed protocols are absent because their reconcile loops only query
// NodeID IS NULL rows, so a node-assigned one would never be reconciled at all.
// A new protocol defaults to ineligible until added here, as on the frontend.
// nodeEligibleProtocols mirrors the frontend's NODE_ELIGIBLE_PROTOCOLS. A sidecar
// protocol's row is local on the node it is pushed to, so that panel runs it.
var nodeEligibleProtocols = map[model.Protocol]bool{
model.VLESS: true,
model.VMESS: true,
@@ -64,11 +67,43 @@ var nodeEligibleProtocols = map[model.Protocol]bool{
model.Shadowsocks: true,
model.Hysteria: true,
model.WireGuard: true,
model.MTProto: true,
model.AmneziaWG: true,
model.TUIC: true,
}
// isNodeEligibleProtocol reports whether protocol may be assigned to a node.
func isNodeEligibleProtocol(protocol model.Protocol) bool {
return nodeEligibleProtocols[protocol]
// nodeProtocolFirstRelease is the panel release that introduced each protocol
// newer than node support itself; an older node would hand it to Xray as-is.
var nodeProtocolFirstRelease = map[model.Protocol]string{
model.MTProto: "v3.5.0",
model.AmneziaWG: "v3.7.0",
model.TUIC: "v3.8.0",
}
// checkNodeCanHostProtocol refuses assigning protocol to nodeID unless the
// protocol may live on a node and that node's panel is new enough to run it.
func checkNodeCanHostProtocol(db *gorm.DB, nodeID int, protocol model.Protocol) error {
if !nodeEligibleProtocols[protocol] {
return common.NewErrorf("%s inbounds cannot be assigned to a node", protocol)
}
firstRelease, ok := nodeProtocolFirstRelease[protocol]
if !ok {
return nil
}
var node model.Node
if err := db.Select("id", "name", "panel_version").First(&node, nodeID).Error; err != nil {
return err
}
if strings.TrimSpace(node.PanelVersion) == "" {
return common.NewErrorf("node %q has not reported its panel version yet; %s inbounds need %s or newer",
node.Name, protocol, firstRelease)
}
// A dev build reports "dev+<sha>": it tracks main, which carries every protocol.
if cmp, ok := version.Compare(node.PanelVersion, firstRelease); ok && cmp < 0 {
return common.NewErrorf("node %q runs panel %s; %s inbounds need %s or newer",
node.Name, node.PanelVersion, protocol, firstRelease)
}
return nil
}
// vlessEncryptionEnabled reports whether a VLESS inbound has VLESS-level
@@ -88,21 +88,3 @@ func TestInboundCanHostFallbacks_StaysTcpOnly(t *testing.T) {
t.Errorf("inboundCanHostFallbacks(nil) = true, want false")
}
}
// Mirrors NODE_ELIGIBLE_PROTOCOLS in
// frontend/src/pages/inbounds/form/InboundFormModal.tsx -- keep both lists
// in sync if a protocol's node-eligibility ever changes.
func TestIsNodeEligibleProtocol(t *testing.T) {
eligible := []model.Protocol{model.VLESS, model.VMESS, model.Trojan, model.Shadowsocks, model.Hysteria, model.WireGuard}
for _, p := range eligible {
if !isNodeEligibleProtocol(p) {
t.Errorf("isNodeEligibleProtocol(%q) = false, want true", p)
}
}
ineligible := []model.Protocol{model.MTProto, model.AmneziaWG, model.Mixed, model.HTTP, model.Tunnel}
for _, p := range ineligible {
if isNodeEligibleProtocol(p) {
t.Errorf("isNodeEligibleProtocol(%q) = true, want false", p)
}
}
}