outbound_manager_test.go 3.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145
  1. package amneziawgnet
  2. import (
  3. "net"
  4. "sync"
  5. "testing"
  6. "time"
  7. "github.com/mhsanaei/3x-ui/v3/internal/amneziawg"
  8. "github.com/mhsanaei/3x-ui/v3/internal/util/wireguard"
  9. )
  10. // egressPortBound reports whether 127.0.0.1:<EgressBasePort> accepts TCP.
  11. func egressPortBound(t *testing.T) bool {
  12. t.Helper()
  13. conn, err := net.DialTimeout("tcp", net.JoinHostPort("127.0.0.1", itoa(int(EgressBasePort))), 500*time.Millisecond)
  14. if err != nil {
  15. return false
  16. }
  17. conn.Close()
  18. return true
  19. }
  20. func itoa(n int) string {
  21. if n == 0 {
  22. return "0"
  23. }
  24. var b [8]byte
  25. i := len(b)
  26. for n > 0 {
  27. i--
  28. b[i] = byte('0' + n%10)
  29. n /= 10
  30. }
  31. return string(b[i:])
  32. }
  33. // newTestOutboundDesired builds one runnable outbound desired state.
  34. func newTestOutboundDesired(t *testing.T, tag string) OutboundDesired {
  35. t.Helper()
  36. priv, _, err := wireguard.GenerateWireguardKeypair()
  37. if err != nil {
  38. t.Fatalf("generate keypair: %v", err)
  39. }
  40. return OutboundDesired{
  41. Instance: amneziawg.OutboundInstance{
  42. Tag: tag,
  43. Address: []string{"10.204.0.1/24"},
  44. MTU: 1420,
  45. PrivateKey: priv,
  46. ListenPort: 0,
  47. },
  48. }
  49. }
  50. // TestOutboundManagerReconcileEmptyDesiredClosesEgress verifies that an empty
  51. // desired set tears down interfaces and releases 127.0.0.1:64900.
  52. func TestOutboundManagerReconcileEmptyDesiredClosesEgress(t *testing.T) {
  53. m := &OutboundManager{iface: map[string]*managedOutbound{}}
  54. defer m.Reconcile(nil)
  55. // Other tests in this package may leave the process-wide egress
  56. // singleton bound; converge to a known-free state before pinning.
  57. GetEgressServer().Close()
  58. if egressPortBound(t) {
  59. t.Fatal("egress port still bound after Close; Close() failed to release it")
  60. }
  61. // Non-empty: listener must come up.
  62. d := newTestOutboundDesired(t, "t1")
  63. m.Reconcile([]OutboundDesired{d})
  64. if !egressPortBound(t) {
  65. t.Fatal("egress port not bound after Reconcile with a desired outbound")
  66. }
  67. // Empty: listener must be released so other listeners can take the port.
  68. m.Reconcile(nil)
  69. if egressPortBound(t) {
  70. t.Fatal("egress port still bound after Reconcile(nil)")
  71. }
  72. ln, err := net.Listen("tcp", net.JoinHostPort("127.0.0.1", itoa(int(EgressBasePort))))
  73. if err != nil {
  74. t.Fatalf("egress port must be free after Reconcile(nil): %v", err)
  75. }
  76. ln.Close()
  77. // Back to non-empty and empty again: Close/Listen must be repeatable.
  78. m.Reconcile([]OutboundDesired{d})
  79. if !egressPortBound(t) {
  80. t.Fatal("egress port not re-bound after a second non-empty Reconcile")
  81. }
  82. m.Reconcile(nil)
  83. if egressPortBound(t) {
  84. t.Fatal("egress port still bound after a second Reconcile(nil)")
  85. }
  86. }
  87. // TestEgressServerCloseDuringConcurrentAccepts ensures Close during
  88. // concurrent accepts shuts down cleanly without hanging wg.Wait().
  89. func TestEgressServerCloseDuringConcurrentAccepts(t *testing.T) {
  90. srv := GetEgressServer()
  91. if err := srv.Listen(); err != nil {
  92. t.Fatal(err)
  93. }
  94. stop := make(chan struct{})
  95. done := make(chan struct{})
  96. var clientWg sync.WaitGroup
  97. go func() {
  98. defer close(done)
  99. for {
  100. select {
  101. case <-stop:
  102. return
  103. default:
  104. c, err := net.DialTimeout("tcp", net.JoinHostPort("127.0.0.1", itoa(int(EgressBasePort))), 50*time.Millisecond)
  105. if err == nil {
  106. clientWg.Add(1)
  107. go func(conn net.Conn) {
  108. defer clientWg.Done()
  109. time.Sleep(20 * time.Millisecond)
  110. conn.Close()
  111. }(c)
  112. }
  113. time.Sleep(2 * time.Millisecond)
  114. }
  115. }
  116. }()
  117. time.Sleep(20 * time.Millisecond)
  118. closeChan := make(chan struct{})
  119. go func() {
  120. srv.Close()
  121. close(closeChan)
  122. }()
  123. select {
  124. case <-closeChan:
  125. case <-time.After(3 * time.Second):
  126. t.Fatal("srv.Close() hung waiting for connection handlers to exit")
  127. }
  128. close(stop)
  129. <-done
  130. clientWg.Wait()
  131. }