diff --git a/internal/amneziawgnet/egress.go b/internal/amneziawgnet/egress.go index 572900cd7..821ad533d 100644 --- a/internal/amneziawgnet/egress.go +++ b/internal/amneziawgnet/egress.go @@ -4,6 +4,7 @@ import ( "context" "crypto/hmac" "encoding/binary" + "errors" "fmt" "io" "net" @@ -18,8 +19,8 @@ import ( "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 // socks5EgressServer is a minimal loopback SOCKS5 server routing Xray's @@ -105,17 +106,22 @@ func (s *socks5EgressServer) DeleteStack(tag string) { 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 { s.mu.Lock() defer s.mu.Unlock() if s.listener != 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 { - 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.closing = make(chan struct{}) @@ -125,6 +131,33 @@ func (s *socks5EgressServer) Listen() error { 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 // error always observes closing. func (s *socks5EgressServer) Close() { diff --git a/internal/amneziawgnet/egress_domain_test.go b/internal/amneziawgnet/egress_domain_test.go index cf1471966..f88dcab07 100644 --- a/internal/amneziawgnet/egress_domain_test.go +++ b/internal/amneziawgnet/egress_domain_test.go @@ -358,7 +358,7 @@ func TestEgressGreetingRejectsNoAuthClient(t *testing.T) { tun := newPairedTunnelForTest(t) 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 { t.Fatal(err) } @@ -384,7 +384,7 @@ func TestEgressConnectDomainResolvesThroughTunnel(t *testing.T) { // fail fast (nothing listens on :80), while proving resolution happened. 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 { t.Fatal(err) } @@ -436,7 +436,7 @@ func TestEgressConnectDomainIPv6OnlyTunnelResolvesThroughTunnel(t *testing.T) { resetTunnelDNSCacheForTest() 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 { t.Fatal(err) } @@ -482,7 +482,7 @@ func TestEgressUDPDatagramDomainForwardedIntoTunnel(t *testing.T) { } 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 { t.Fatal(err) } @@ -585,7 +585,7 @@ func TestEgressUDPDatagramDomainInterleavedClients(t *testing.T) { }() 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 { t.Fatal(err) } diff --git a/internal/amneziawgnet/egress_port_test.go b/internal/amneziawgnet/egress_port_test.go new file mode 100644 index 000000000..5030f0b6a --- /dev/null +++ b/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()) + } +} diff --git a/internal/amneziawgnet/outbound_manager.go b/internal/amneziawgnet/outbound_manager.go index ea17d26d0..0b367ea86 100644 --- a/internal/amneziawgnet/outbound_manager.go +++ b/internal/amneziawgnet/outbound_manager.go @@ -59,7 +59,7 @@ func (m *OutboundManager) Reconcile(desired []OutboundDesired) { defer m.mu.Unlock() // 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 { for tag, cur := range m.iface { cur.dev.Close() diff --git a/internal/amneziawgnet/outbound_manager_test.go b/internal/amneziawgnet/outbound_manager_test.go index 119a34f53..1ea45d2eb 100644 --- a/internal/amneziawgnet/outbound_manager_test.go +++ b/internal/amneziawgnet/outbound_manager_test.go @@ -10,10 +10,10 @@ import ( "github.com/mhsanaei/3x-ui/v3/internal/util/wireguard" ) -// egressPortBound reports whether 127.0.0.1: accepts TCP. +// egressPortBound reports whether the egress server's Port accepts TCP. func egressPortBound(t *testing.T) bool { 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 { return false } @@ -54,7 +54,7 @@ func newTestOutboundDesired(t *testing.T, tag string) OutboundDesired { } // 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) { m := &OutboundManager{iface: map[string]*managedOutbound{}} defer m.Reconcile(nil) @@ -72,13 +72,14 @@ func TestOutboundManagerReconcileEmptyDesiredClosesEgress(t *testing.T) { if !egressPortBound(t) { 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. m.Reconcile(nil) if egressPortBound(t) { 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 { 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 { t.Fatal(err) } + addr := net.JoinHostPort("127.0.0.1", itoa(srv.Port())) stop := make(chan struct{}) done := make(chan struct{}) @@ -113,7 +115,7 @@ func TestEgressServerCloseDuringConcurrentAccepts(t *testing.T) { case <-stop: return 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 { clientWg.Add(1) go func(conn net.Conn) { diff --git a/internal/amneziawgnet/socks_bridge.go b/internal/amneziawgnet/socks_bridge.go index 54da7fd4e..9995459cb 100644 --- a/internal/amneziawgnet/socks_bridge.go +++ b/internal/amneziawgnet/socks_bridge.go @@ -1,6 +1,12 @@ 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 // form, preserving sibling keys; false = unbridgeable, fail loudly upstream. @@ -13,9 +19,10 @@ func BuildSocksBridge(raw []byte) ([]byte, bool) { if tag == "" { return nil, false } + port := EgressPort() settings := map[string]any{ "address": "127.0.0.1", - "port": EgressBasePort, + "port": port, "user": tag, "pass": SocksPassword(), } @@ -29,5 +36,14 @@ func BuildSocksBridge(raw []byte) ([]byte, bool) { if err != nil { return nil, false } + bridgedPort.Store(int64(port)) 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 +} diff --git a/internal/web/job/amneziawg_job.go b/internal/web/job/amneziawg_job.go index d7bb3d6fb..452c4b4ac 100644 --- a/internal/web/job/amneziawg_job.go +++ b/internal/web/job/amneziawg_job.go @@ -15,6 +15,7 @@ import ( type AmneziaWGJob struct { inboundService service.InboundService settingService service.SettingService + xrayService service.XrayService } // NewAmneziaWGJob creates a new AmneziaWG reconcile job instance. @@ -55,6 +56,10 @@ func (j *AmneziaWGJob) Run() { return } 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. diff --git a/internal/web/service/port_conflict.go b/internal/web/service/port_conflict.go index 772616493..38ba7696f 100644 --- a/internal/web/service/port_conflict.go +++ b/internal/web/service/port_conflict.go @@ -256,9 +256,9 @@ func checkPortConflictTx(db *gorm.DB, inbound *model.Inbound, ignoreId int) (*po }, 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. - if inbound.NodeID == nil && inbound.Port == int(amneziawgnet.EgressBasePort) && + if inbound.NodeID == nil && inbound.Port == amneziawgnet.EgressPort() && newBits&transportTCP != 0 && listenOverlaps(loopbackBind, inboundBindAddr(inbound)) { return &portConflictDetail{ Tag: "amneziawg-egress", diff --git a/internal/web/service/port_conflict_test.go b/internal/web/service/port_conflict_test.go index 085092436..c3a1fe356 100644 --- a/internal/web/service/port_conflict_test.go +++ b/internal/web/service/port_conflict_test.go @@ -1,7 +1,9 @@ package service import ( + "net" "path/filepath" + "strconv" "strings" "sync" "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) { setupConflictDB(t) seedInboundConflict(t, "awg-1", "0.0.0.0", 51820, model.AmneziaWG, ``, amneziawgRoutedSettings) diff --git a/internal/web/service/xray_amneziawg_outbound_test.go b/internal/web/service/xray_amneziawg_outbound_test.go index 849ac993f..2f5cc7d72 100644 --- a/internal/web/service/xray_amneziawg_outbound_test.go +++ b/internal/web/service/xray_amneziawg_outbound_test.go @@ -10,7 +10,7 @@ import ( "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) { return wgutil.GenerateWireguardKeypair()