Browse Source

fix(amneziawg): fall back to a free egress port when 64900 is refused

The panel's SOCKS5 egress for AmneziaWG outbounds bound the fixed
127.0.0.1:64900, and every generated socks bridge dialed that constant.
64900 sits inside Windows' dynamic port range, where the OS can reserve
whole blocks (this host excludes 64885-64984), so on the Windows builds
release.yml ships the listener could stay down and every AmneziaWG
outbound with it. Listen now tries 64900 first and falls back to any
free loopback port; bridges, the outbound probe and the port-conflict
check use EgressPort(), the port actually held. Bridges are generated
apart from the listener, so BuildSocksBridge records the port it wrote
and the AmneziaWG job requests an Xray restart while the listener holds
a different one. Where 64900 is free nothing changes.

The job's restart request is two lines of wiring no test reaches; the
staleness it acts on is pinned by
TestBridgesStaleUntilRegeneratedForTheBoundPort.
MHSanaei 6 hours ago
parent
commit
75f3702dd3

+ 39 - 6
internal/amneziawgnet/egress.go

@@ -4,6 +4,7 @@ import (
 	"context"
 	"context"
 	"crypto/hmac"
 	"crypto/hmac"
 	"encoding/binary"
 	"encoding/binary"
+	"errors"
 	"fmt"
 	"fmt"
 	"io"
 	"io"
 	"net"
 	"net"
@@ -18,8 +19,8 @@ import (
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 )
 )
 
 
-// EgressBasePort is the fixed loopback port of the panel's SOCKS5 egress
-// server; it appears in every generated amneziawg socks bridge.
+// EgressBasePort is the loopback port the panel's SOCKS5 egress server tries
+// first; generated amneziawg socks bridges dial whichever port it holds.
 const EgressBasePort = 64900
 const EgressBasePort = 64900
 
 
 // socks5EgressServer is a minimal loopback SOCKS5 server routing Xray's
 // socks5EgressServer is a minimal loopback SOCKS5 server routing Xray's
@@ -105,17 +106,22 @@ func (s *socks5EgressServer) DeleteStack(tag string) {
 	flushTunnelDNSCacheForTag(tag)
 	flushTunnelDNSCacheForTag(tag)
 }
 }
 
 
-// Listen starts accepting on the loopback listener. Idempotent; a bind
-// failure is returned and retried by the caller's reconcile tick.
+// Listen starts accepting on EgressBasePort, or on any free loopback port when
+// the OS refuses it. Idempotent; the caller's reconcile tick retries a failure.
 func (s *socks5EgressServer) Listen() error {
 func (s *socks5EgressServer) Listen() error {
 	s.mu.Lock()
 	s.mu.Lock()
 	defer s.mu.Unlock()
 	defer s.mu.Unlock()
 	if s.listener != nil {
 	if s.listener != nil {
 		return nil
 		return nil
 	}
 	}
-	ln, err := (&net.ListenConfig{}).Listen(context.Background(), "tcp", fmt.Sprintf("127.0.0.1:%d", EgressBasePort))
+	lc := &net.ListenConfig{}
+	ln, err := lc.Listen(context.Background(), "tcp", fmt.Sprintf("127.0.0.1:%d", EgressBasePort))
 	if err != nil {
 	if err != nil {
-		return fmt.Errorf("amneziawgnet: egress listen: %w", err)
+		var fallbackErr error
+		if ln, fallbackErr = lc.Listen(context.Background(), "tcp", "127.0.0.1:0"); fallbackErr != nil {
+			return fmt.Errorf("amneziawgnet: egress listen: %w", errors.Join(err, fallbackErr))
+		}
+		logger.Warningf("amneziawgnet: egress port %d unavailable, using a free port: %v", EgressBasePort, err)
 	}
 	}
 	s.listener = ln
 	s.listener = ln
 	s.closing = make(chan struct{})
 	s.closing = make(chan struct{})
@@ -125,6 +131,33 @@ func (s *socks5EgressServer) Listen() error {
 	return nil
 	return nil
 }
 }
 
 
+// Port is the port generated socks bridges must dial: the bound one, or
+// EgressBasePort while nothing is bound, since Listen tries it first.
+func (s *socks5EgressServer) Port() int {
+	if port, ok := s.boundPort(); ok {
+		return port
+	}
+	return EgressBasePort
+}
+
+func (s *socks5EgressServer) boundPort() (int, bool) {
+	s.mu.Lock()
+	defer s.mu.Unlock()
+	if s.listener == nil {
+		return 0, false
+	}
+	addr, ok := s.listener.Addr().(*net.TCPAddr)
+	if !ok {
+		return 0, false
+	}
+	return addr.Port, true
+}
+
+// EgressPort is the process-wide egress server's Port.
+func EgressPort() int {
+	return GetEgressServer().Port()
+}
+
 // Close stops the listener and in-flight handlers; signal first so an accept
 // Close stops the listener and in-flight handlers; signal first so an accept
 // error always observes closing.
 // error always observes closing.
 func (s *socks5EgressServer) Close() {
 func (s *socks5EgressServer) Close() {

+ 5 - 5
internal/amneziawgnet/egress_domain_test.go

@@ -358,7 +358,7 @@ func TestEgressGreetingRejectsNoAuthClient(t *testing.T) {
 	tun := newPairedTunnelForTest(t)
 	tun := newPairedTunnelForTest(t)
 	registerEgressDeviceForTest(t, tun.client)
 	registerEgressDeviceForTest(t, tun.client)
 
 
-	ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(EgressBasePort)))
+	ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(GetEgressServer().Port())))
 	if err != nil {
 	if err != nil {
 		t.Fatal(err)
 		t.Fatal(err)
 	}
 	}
@@ -384,7 +384,7 @@ func TestEgressConnectDomainResolvesThroughTunnel(t *testing.T) {
 	// fail fast (nothing listens on :80), while proving resolution happened.
 	// fail fast (nothing listens on :80), while proving resolution happened.
 	gotQuery := tun.overrideDNS(t, tun.serverIP)
 	gotQuery := tun.overrideDNS(t, tun.serverIP)
 
 
-	ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(EgressBasePort)))
+	ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(GetEgressServer().Port())))
 	if err != nil {
 	if err != nil {
 		t.Fatal(err)
 		t.Fatal(err)
 	}
 	}
@@ -436,7 +436,7 @@ func TestEgressConnectDomainIPv6OnlyTunnelResolvesThroughTunnel(t *testing.T) {
 	resetTunnelDNSCacheForTest()
 	resetTunnelDNSCacheForTest()
 	gotQuery := tun.startDNS(t, tun.serverIP)
 	gotQuery := tun.startDNS(t, tun.serverIP)
 
 
-	ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(EgressBasePort)))
+	ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(GetEgressServer().Port())))
 	if err != nil {
 	if err != nil {
 		t.Fatal(err)
 		t.Fatal(err)
 	}
 	}
@@ -482,7 +482,7 @@ func TestEgressUDPDatagramDomainForwardedIntoTunnel(t *testing.T) {
 	}
 	}
 	defer in.Close()
 	defer in.Close()
 
 
-	ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(EgressBasePort)))
+	ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(GetEgressServer().Port())))
 	if err != nil {
 	if err != nil {
 		t.Fatal(err)
 		t.Fatal(err)
 	}
 	}
@@ -585,7 +585,7 @@ func TestEgressUDPDatagramDomainInterleavedClients(t *testing.T) {
 	}()
 	}()
 
 
 	dialUDP := func() *net.UDPConn {
 	dialUDP := func() *net.UDPConn {
-		ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(EgressBasePort)))
+		ctl, err := (&net.Dialer{Timeout: egressTestDialTimeout}).Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(GetEgressServer().Port())))
 		if err != nil {
 		if err != nil {
 			t.Fatal(err)
 			t.Fatal(err)
 		}
 		}

+ 81 - 0
internal/amneziawgnet/egress_port_test.go

@@ -0,0 +1,81 @@
+package amneziawgnet
+
+import (
+	"encoding/json"
+	"net"
+	"strconv"
+	"testing"
+)
+
+// holdEgressBasePort occupies EgressBasePort the way another service would; a
+// port the OS already refuses, such as a Windows reservation, needs no holder.
+func holdEgressBasePort(t *testing.T) {
+	t.Helper()
+	ln, err := net.Listen("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(EgressBasePort)))
+	if err != nil {
+		return
+	}
+	t.Cleanup(func() { ln.Close() })
+}
+
+// bridgePort is the port a socks bridge generated right now dials.
+func bridgePort(t *testing.T) int {
+	t.Helper()
+	out, ok := BuildSocksBridge([]byte(`{"protocol":"amneziawg","tag":"awg-hop","settings":{}}`))
+	if !ok {
+		t.Fatal("bridge rejected")
+	}
+	var got struct {
+		Settings struct {
+			Port int `json:"port"`
+		} `json:"settings"`
+	}
+	if err := json.Unmarshal(out, &got); err != nil {
+		t.Fatal(err)
+	}
+	return got.Settings.Port
+}
+
+// Windows can reserve a port range covering EgressBasePort, and any host can run
+// another service on it; the egress must still come up and report where.
+func TestEgressListenFallsBackWhenBasePortIsTaken(t *testing.T) {
+	srv := GetEgressServer()
+	srv.Close()
+	t.Cleanup(srv.Close)
+	holdEgressBasePort(t)
+
+	if err := srv.Listen(); err != nil {
+		t.Fatalf("Listen with EgressBasePort taken: %v", err)
+	}
+	port := srv.Port()
+	if port == EgressBasePort {
+		t.Fatalf("Port() = %d, the taken EgressBasePort", port)
+	}
+	conn, err := net.Dial("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(port)))
+	if err != nil {
+		t.Fatalf("egress not accepting on Port() %d: %v", port, err)
+	}
+	conn.Close()
+}
+
+// Xray's bridges are generated apart from the listener, so a listener that came
+// up elsewhere must be reported until the bridges are regenerated for it.
+func TestBridgesStaleUntilRegeneratedForTheBoundPort(t *testing.T) {
+	srv := GetEgressServer()
+	srv.Close()
+	t.Cleanup(srv.Close)
+	holdEgressBasePort(t)
+
+	if got := bridgePort(t); got != EgressBasePort {
+		t.Fatalf("bridge generated before Listen dials %d, want %d", got, EgressBasePort)
+	}
+	if err := srv.Listen(); err != nil {
+		t.Fatal(err)
+	}
+	if !BridgesStale() {
+		t.Fatalf("bridges dial %d while the egress listens on %d, but BridgesStale() = false", EgressBasePort, srv.Port())
+	}
+	if got := bridgePort(t); got != srv.Port() || BridgesStale() {
+		t.Fatalf("regenerated bridge dials %d with BridgesStale() = %v, want %d and false", got, BridgesStale(), srv.Port())
+	}
+}

+ 1 - 1
internal/amneziawgnet/outbound_manager.go

@@ -59,7 +59,7 @@ func (m *OutboundManager) Reconcile(desired []OutboundDesired) {
 	defer m.mu.Unlock()
 	defer m.mu.Unlock()
 
 
 	// Empty desired converges to "no tunnels": close egress listener so
 	// Empty desired converges to "no tunnels": close egress listener so
-	// 127.0.0.1:64900 stays free on installs without AWG outbounds.
+	// the egress port stays free on installs without AWG outbounds.
 	if len(desired) == 0 {
 	if len(desired) == 0 {
 		for tag, cur := range m.iface {
 		for tag, cur := range m.iface {
 			cur.dev.Close()
 			cur.dev.Close()

+ 7 - 5
internal/amneziawgnet/outbound_manager_test.go

@@ -10,10 +10,10 @@ import (
 	"github.com/mhsanaei/3x-ui/v3/internal/util/wireguard"
 	"github.com/mhsanaei/3x-ui/v3/internal/util/wireguard"
 )
 )
 
 
-// egressPortBound reports whether 127.0.0.1:<EgressBasePort> accepts TCP.
+// egressPortBound reports whether the egress server's Port accepts TCP.
 func egressPortBound(t *testing.T) bool {
 func egressPortBound(t *testing.T) bool {
 	t.Helper()
 	t.Helper()
-	conn, err := net.DialTimeout("tcp", net.JoinHostPort("127.0.0.1", itoa(int(EgressBasePort))), 500*time.Millisecond)
+	conn, err := net.DialTimeout("tcp", net.JoinHostPort("127.0.0.1", itoa(GetEgressServer().Port())), 500*time.Millisecond)
 	if err != nil {
 	if err != nil {
 		return false
 		return false
 	}
 	}
@@ -54,7 +54,7 @@ func newTestOutboundDesired(t *testing.T, tag string) OutboundDesired {
 }
 }
 
 
 // TestOutboundManagerReconcileEmptyDesiredClosesEgress verifies that an empty
 // TestOutboundManagerReconcileEmptyDesiredClosesEgress verifies that an empty
-// desired set tears down interfaces and releases 127.0.0.1:64900.
+// desired set tears down interfaces and releases the egress port.
 func TestOutboundManagerReconcileEmptyDesiredClosesEgress(t *testing.T) {
 func TestOutboundManagerReconcileEmptyDesiredClosesEgress(t *testing.T) {
 	m := &OutboundManager{iface: map[string]*managedOutbound{}}
 	m := &OutboundManager{iface: map[string]*managedOutbound{}}
 	defer m.Reconcile(nil)
 	defer m.Reconcile(nil)
@@ -72,13 +72,14 @@ func TestOutboundManagerReconcileEmptyDesiredClosesEgress(t *testing.T) {
 	if !egressPortBound(t) {
 	if !egressPortBound(t) {
 		t.Fatal("egress port not bound after Reconcile with a desired outbound")
 		t.Fatal("egress port not bound after Reconcile with a desired outbound")
 	}
 	}
+	held := GetEgressServer().Port()
 
 
 	// Empty: listener must be released so other listeners can take the port.
 	// Empty: listener must be released so other listeners can take the port.
 	m.Reconcile(nil)
 	m.Reconcile(nil)
 	if egressPortBound(t) {
 	if egressPortBound(t) {
 		t.Fatal("egress port still bound after Reconcile(nil)")
 		t.Fatal("egress port still bound after Reconcile(nil)")
 	}
 	}
-	ln, err := net.Listen("tcp", net.JoinHostPort("127.0.0.1", itoa(int(EgressBasePort))))
+	ln, err := net.Listen("tcp", net.JoinHostPort("127.0.0.1", itoa(held)))
 	if err != nil {
 	if err != nil {
 		t.Fatalf("egress port must be free after Reconcile(nil): %v", err)
 		t.Fatalf("egress port must be free after Reconcile(nil): %v", err)
 	}
 	}
@@ -102,6 +103,7 @@ func TestEgressServerCloseDuringConcurrentAccepts(t *testing.T) {
 	if err := srv.Listen(); err != nil {
 	if err := srv.Listen(); err != nil {
 		t.Fatal(err)
 		t.Fatal(err)
 	}
 	}
+	addr := net.JoinHostPort("127.0.0.1", itoa(srv.Port()))
 
 
 	stop := make(chan struct{})
 	stop := make(chan struct{})
 	done := make(chan struct{})
 	done := make(chan struct{})
@@ -113,7 +115,7 @@ func TestEgressServerCloseDuringConcurrentAccepts(t *testing.T) {
 			case <-stop:
 			case <-stop:
 				return
 				return
 			default:
 			default:
-				c, err := net.DialTimeout("tcp", net.JoinHostPort("127.0.0.1", itoa(int(EgressBasePort))), 50*time.Millisecond)
+				c, err := net.DialTimeout("tcp", addr, 50*time.Millisecond)
 				if err == nil {
 				if err == nil {
 					clientWg.Add(1)
 					clientWg.Add(1)
 					go func(conn net.Conn) {
 					go func(conn net.Conn) {

+ 18 - 2
internal/amneziawgnet/socks_bridge.go

@@ -1,6 +1,12 @@
 package amneziawgnet
 package amneziawgnet
 
 
-import "encoding/json"
+import (
+	"encoding/json"
+	"sync/atomic"
+)
+
+// bridgedPort is the port the last generated socks bridge dials; 0 before any.
+var bridgedPort atomic.Int64
 
 
 // BuildSocksBridge swaps an "amneziawg" outbound for its loopback socks
 // BuildSocksBridge swaps an "amneziawg" outbound for its loopback socks
 // form, preserving sibling keys; false = unbridgeable, fail loudly upstream.
 // form, preserving sibling keys; false = unbridgeable, fail loudly upstream.
@@ -13,9 +19,10 @@ func BuildSocksBridge(raw []byte) ([]byte, bool) {
 	if tag == "" {
 	if tag == "" {
 		return nil, false
 		return nil, false
 	}
 	}
+	port := EgressPort()
 	settings := map[string]any{
 	settings := map[string]any{
 		"address": "127.0.0.1",
 		"address": "127.0.0.1",
-		"port":    EgressBasePort,
+		"port":    port,
 		"user":    tag,
 		"user":    tag,
 		"pass":    SocksPassword(),
 		"pass":    SocksPassword(),
 	}
 	}
@@ -29,5 +36,14 @@ func BuildSocksBridge(raw []byte) ([]byte, bool) {
 	if err != nil {
 	if err != nil {
 		return nil, false
 		return nil, false
 	}
 	}
+	bridgedPort.Store(int64(port))
 	return out, true
 	return out, true
 }
 }
+
+// BridgesStale reports that the egress listener holds another port than the last
+// generated socks bridge dials, so Xray has to regenerate its config.
+func BridgesStale() bool {
+	bound, listening := GetEgressServer().boundPort()
+	bridged := int(bridgedPort.Load())
+	return listening && bridged != 0 && bridged != bound
+}

+ 5 - 0
internal/web/job/amneziawg_job.go

@@ -15,6 +15,7 @@ import (
 type AmneziaWGJob struct {
 type AmneziaWGJob struct {
 	inboundService service.InboundService
 	inboundService service.InboundService
 	settingService service.SettingService
 	settingService service.SettingService
+	xrayService    service.XrayService
 }
 }
 
 
 // NewAmneziaWGJob creates a new AmneziaWG reconcile job instance.
 // NewAmneziaWGJob creates a new AmneziaWG reconcile job instance.
@@ -55,6 +56,10 @@ func (j *AmneziaWGJob) Run() {
 		return
 		return
 	}
 	}
 	amneziawgnet.GetOutboundManager().Reconcile(outboundDesired)
 	amneziawgnet.GetOutboundManager().Reconcile(outboundDesired)
+	// Xray's bridges are generated apart from the listener; one that moved needs them regenerated.
+	if amneziawgnet.BridgesStale() {
+		j.xrayService.SetToNeedRestart()
+	}
 }
 }
 
 
 // desiredOutboundInstances derives client instances per template "amneziawg" outbound.
 // desiredOutboundInstances derives client instances per template "amneziawg" outbound.

+ 2 - 2
internal/web/service/port_conflict.go

@@ -256,9 +256,9 @@ func checkPortConflictTx(db *gorm.DB, inbound *model.Inbound, ignoreId int) (*po
 		}, nil
 		}, nil
 	}
 	}
 
 
-	// Egress SOCKS server holds loopback EgressBasePort when AWG outbounds are
+	// Egress SOCKS server holds loopback EgressPort when AWG outbounds are
 	// active; conflict check prevents inbounds from colliding with it.
 	// active; conflict check prevents inbounds from colliding with it.
-	if inbound.NodeID == nil && inbound.Port == int(amneziawgnet.EgressBasePort) &&
+	if inbound.NodeID == nil && inbound.Port == amneziawgnet.EgressPort() &&
 		newBits&transportTCP != 0 && listenOverlaps(loopbackBind, inboundBindAddr(inbound)) {
 		newBits&transportTCP != 0 && listenOverlaps(loopbackBind, inboundBindAddr(inbound)) {
 		return &portConflictDetail{
 		return &portConflictDetail{
 			Tag:        "amneziawg-egress",
 			Tag:        "amneziawg-egress",

+ 26 - 0
internal/web/service/port_conflict_test.go

@@ -1,7 +1,9 @@
 package service
 package service
 
 
 import (
 import (
+	"net"
 	"path/filepath"
 	"path/filepath"
+	"strconv"
 	"strings"
 	"strings"
 	"sync"
 	"sync"
 	"testing"
 	"testing"
@@ -754,6 +756,30 @@ func TestCheckPortConflict_EgressPortBlockedLocal(t *testing.T) {
 	}
 	}
 }
 }
 
 
+// Where EgressBasePort is taken the egress listens on another port, and that
+// is the port an inbound must not collide with.
+func TestCheckPortConflict_EgressPortFollowsTheListener(t *testing.T) {
+	setupConflictDB(t)
+	if ln, err := net.Listen("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(amneziawgnet.EgressBasePort))); err == nil {
+		t.Cleanup(func() { ln.Close() })
+	}
+	egress := amneziawgnet.GetEgressServer()
+	if err := egress.Listen(); err != nil {
+		t.Fatal(err)
+	}
+	t.Cleanup(egress.Close)
+
+	svc := &InboundService{}
+	candidate := &model.Inbound{Tag: "vless-bridge", Listen: "0.0.0.0", Port: egress.Port(), Protocol: model.VLESS}
+	got, err := svc.checkPortConflict(candidate, 0)
+	if err != nil {
+		t.Fatalf("checkPortConflict: %v", err)
+	}
+	if got == nil || got.Tag != "amneziawg-egress" {
+		t.Fatalf("an inbound on the egress's port %d must conflict with amneziawg-egress, got %+v", egress.Port(), got)
+	}
+}
+
 func TestCheckPortConflict_AmneziawgnetSocksRelayBlockedLocal(t *testing.T) {
 func TestCheckPortConflict_AmneziawgnetSocksRelayBlockedLocal(t *testing.T) {
 	setupConflictDB(t)
 	setupConflictDB(t)
 	seedInboundConflict(t, "awg-1", "0.0.0.0", 51820, model.AmneziaWG, ``, amneziawgRoutedSettings)
 	seedInboundConflict(t, "awg-1", "0.0.0.0", 51820, model.AmneziaWG, ``, amneziawgRoutedSettings)

+ 1 - 1
internal/web/service/xray_amneziawg_outbound_test.go

@@ -10,7 +10,7 @@ import (
 	"github.com/mhsanaei/3x-ui/v3/internal/xray"
 	"github.com/mhsanaei/3x-ui/v3/internal/xray"
 )
 )
 
 
-func amneziawgnetEgressPortForTest() int { return amneziawgnet.EgressBasePort }
+func amneziawgnetEgressPortForTest() int { return amneziawgnet.EgressPort() }
 
 
 func wgKeypairForTest() (priv, pub string, err error) {
 func wgKeypairForTest() (priv, pub string, err error) {
 	return wgutil.GenerateWireguardKeypair()
 	return wgutil.GenerateWireguardKeypair()