manager_test.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399
  1. package tuic
  2. import (
  3. "context"
  4. "crypto/tls"
  5. "fmt"
  6. "net"
  7. "reflect"
  8. "strings"
  9. "sync"
  10. "testing"
  11. "time"
  12. "unsafe"
  13. "github.com/apernet/quic-go"
  14. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  15. )
  16. func TestEnsureStartsServerAndReconciles(t *testing.T) {
  17. certPEM, keyPEM := generateTestCert(t)
  18. // Find free UDP port
  19. pc, err := net.ListenPacket("udp", "127.0.0.1:0")
  20. if err != nil {
  21. t.Fatal(err)
  22. }
  23. port := pc.LocalAddr().(*net.UDPAddr).Port
  24. _ = pc.Close()
  25. inst := Instance{
  26. Id: 11,
  27. Tag: "tuic-11",
  28. Listen: "127.0.0.1",
  29. Port: port,
  30. Certificate: string(certPEM),
  31. PrivateKey: string(keyPEM),
  32. Clients: []TuicClientSettings{{UUID: "a0000000-0000-0000-0000-000000000001", Password: "p", Email: "e1"}},
  33. }
  34. m := &Manager{servers: map[int]*managed{}, lastStartErr: map[int]string{}}
  35. t.Cleanup(m.StopAll)
  36. if err := m.Ensure(inst); err != nil {
  37. t.Fatalf("Ensure failed: %v", err)
  38. }
  39. if !m.HasRunning() {
  40. t.Fatal("expected manager to have running server")
  41. }
  42. // Port should be taken by the server
  43. if c, err := net.ListenPacket("udp", inst.BindTo()); err == nil {
  44. _ = c.Close()
  45. t.Fatal("expected server to be listening on port")
  46. }
  47. // Reconcile with empty list should remove it
  48. m.Reconcile([]Instance{})
  49. if m.HasRunning() {
  50. t.Fatal("expected no running servers after reconcile empty")
  51. }
  52. // Port should now be released
  53. c, err := net.ListenPacket("udp", inst.BindTo())
  54. if err != nil {
  55. t.Fatalf("expected port to be free after reconcile: %v", err)
  56. }
  57. _ = c.Close()
  58. }
  59. func TestEnsureHotUpdatesUsersWithoutRestart(t *testing.T) {
  60. certPEM, keyPEM := generateTestCert(t)
  61. pc, err := net.ListenPacket("udp", "127.0.0.1:0")
  62. if err != nil {
  63. t.Fatal(err)
  64. }
  65. port := pc.LocalAddr().(*net.UDPAddr).Port
  66. _ = pc.Close()
  67. inst := Instance{
  68. Id: 12,
  69. Tag: "tuic-12",
  70. Listen: "127.0.0.1",
  71. Port: port,
  72. Certificate: string(certPEM),
  73. PrivateKey: string(keyPEM),
  74. Clients: []TuicClientSettings{{UUID: "a0000000-0000-0000-0000-000000000001", Password: "p1", Email: "e1"}},
  75. }
  76. m := &Manager{servers: map[int]*managed{}, lastStartErr: map[int]string{}}
  77. t.Cleanup(m.StopAll)
  78. if err := m.Ensure(inst); err != nil {
  79. t.Fatalf("Ensure failed: %v", err)
  80. }
  81. server1 := m.servers[12].server
  82. // Update user list without changing port or certs
  83. inst.Clients = append(inst.Clients, TuicClientSettings{
  84. UUID: "a0000000-0000-0000-0000-000000000002",
  85. Password: "p2",
  86. Email: "e2",
  87. })
  88. if err := m.Ensure(inst); err != nil {
  89. t.Fatalf("Ensure with updated clients failed: %v", err)
  90. }
  91. server2 := m.servers[12].server
  92. if server1 != server2 {
  93. t.Fatal("expected server instance to be reused across user updates (zero-downtime hot update)")
  94. }
  95. // Verify both users are now in user registry
  96. if len(server2.users.users) != 2 {
  97. t.Fatalf("expected 2 users in registry, got %d", len(server2.users.users))
  98. }
  99. }
  100. func TestEnsureUpdatesTagWithoutRestart(t *testing.T) {
  101. certPEM, keyPEM := generateTestCert(t)
  102. pc, err := net.ListenPacket("udp", "127.0.0.1:0")
  103. if err != nil {
  104. t.Fatal(err)
  105. }
  106. port := pc.LocalAddr().(*net.UDPAddr).Port
  107. _ = pc.Close()
  108. inst := Instance{
  109. Id: 13,
  110. Tag: "old-tag",
  111. Listen: "127.0.0.1",
  112. Port: port,
  113. Certificate: string(certPEM),
  114. PrivateKey: string(keyPEM),
  115. Clients: []TuicClientSettings{{UUID: "a0000000-0000-0000-0000-000000000001", Password: "p", Email: "e"}},
  116. }
  117. m := &Manager{servers: map[int]*managed{}, lastStartErr: map[int]string{}}
  118. t.Cleanup(m.StopAll)
  119. if err := m.Ensure(inst); err != nil {
  120. t.Fatalf("Ensure failed: %v", err)
  121. }
  122. inst.Tag = "new-tag"
  123. if err := m.Ensure(inst); err != nil {
  124. t.Fatalf("Ensure updated tag failed: %v", err)
  125. }
  126. m.mu.Lock()
  127. gotTag := m.servers[13].tag
  128. m.mu.Unlock()
  129. if gotTag != "new-tag" {
  130. t.Fatalf("manager tag = %q, want %q", gotTag, "new-tag")
  131. }
  132. }
  133. func TestEnsureUpdatesControllerWithoutRestartForNewConnections(t *testing.T) {
  134. certPEM, keyPEM := generateTestCert(t)
  135. pc, err := net.ListenPacket("udp", "127.0.0.1:0")
  136. if err != nil {
  137. t.Fatal(err)
  138. }
  139. port := pc.LocalAddr().(*net.UDPAddr).Port
  140. _ = pc.Close()
  141. inst := Instance{
  142. Id: 14,
  143. Tag: "tuic-controller-reload",
  144. Listen: "127.0.0.1",
  145. Port: port,
  146. Certificate: string(certPEM),
  147. PrivateKey: string(keyPEM),
  148. CongestionControl: "bbr",
  149. LogLevel: "debug",
  150. MaxIdleTime: 30,
  151. AuthenticationTimeout: 30,
  152. Clients: []TuicClientSettings{{UUID: "a0000000-0000-0000-0000-000000000001", Password: "p", Email: "e"}},
  153. }
  154. m := &Manager{servers: map[int]*managed{}, lastStartErr: map[int]string{}}
  155. t.Cleanup(m.StopAll)
  156. if err := m.Ensure(inst); err != nil {
  157. t.Fatalf("Ensure failed: %v", err)
  158. }
  159. server := m.servers[inst.Id].server
  160. dial := func() *quic.Conn {
  161. t.Helper()
  162. ctx, cancel := context.WithTimeout(context.Background(), 4*time.Second)
  163. defer cancel()
  164. conn, err := quic.DialAddr(ctx, server.packetConn.LocalAddr().String(), &tls.Config{
  165. InsecureSkipVerify: true,
  166. NextProtos: []string{"h3"},
  167. }, &quic.Config{EnableDatagrams: true, MaxIdleTimeout: 30 * time.Second})
  168. if err != nil {
  169. t.Fatalf("QUIC dial failed: %v", err)
  170. }
  171. return conn
  172. }
  173. connBBR := dial()
  174. defer connBBR.CloseWithError(0, "")
  175. waitForTuicLog(t, "applied bbr congestion controller")
  176. bbrSender := waitForClientCongestionSender(t, server, connBBR, "bbr")
  177. inst.CongestionControl = "cubic"
  178. if err := m.Ensure(inst); err != nil {
  179. t.Fatalf("Ensure after CUBIC update failed: %v", err)
  180. }
  181. if m.servers[inst.Id].server != server {
  182. t.Fatal("changing congestion control restarted the listener")
  183. }
  184. if err := connBBR.Context().Err(); err != nil {
  185. t.Fatalf("existing BBR connection closed after controller update: %v", err)
  186. }
  187. if sender := congestionSenderForClient(t, server, connBBR); sender != bbrSender {
  188. t.Fatal("existing connection's BBR sender changed after hot update")
  189. }
  190. connCubic := dial()
  191. defer connCubic.CloseWithError(0, "")
  192. waitForTuicLog(t, "cubic is not available; applied new_reno congestion controller")
  193. cubicSender := waitForClientCongestionSender(t, server, connCubic, "new_reno")
  194. inst.CongestionControl = "new_reno"
  195. if err := m.Ensure(inst); err != nil {
  196. t.Fatalf("Ensure after New Reno update failed: %v", err)
  197. }
  198. if err := connBBR.Context().Err(); err != nil {
  199. t.Fatalf("existing BBR connection closed after second update: %v", err)
  200. }
  201. if err := connCubic.Context().Err(); err != nil {
  202. t.Fatalf("existing CUBIC connection closed after second update: %v", err)
  203. }
  204. if sender := congestionSenderForClient(t, server, connBBR); sender != bbrSender {
  205. t.Fatal("existing connection's BBR sender changed after second hot update")
  206. }
  207. if sender := congestionSenderForClient(t, server, connCubic); sender != cubicSender {
  208. t.Fatal("existing connection's CUBIC sender changed after second hot update")
  209. }
  210. connReno := dial()
  211. defer connReno.CloseWithError(0, "")
  212. waitForTuicLog(t, "applied new_reno congestion controller")
  213. waitForClientCongestionSender(t, server, connReno, "new_reno")
  214. }
  215. func waitForClientCongestionSender(t *testing.T, server *Server, client *quic.Conn, want string) uintptr {
  216. t.Helper()
  217. deadline := time.Now().Add(4 * time.Second)
  218. var observed []string
  219. seen := make(map[string]struct{})
  220. for time.Now().Before(deadline) {
  221. server.connectionsMu.Lock()
  222. for conn := range server.connections {
  223. remote := conn.RemoteAddr().String()
  224. actual, sender := inspectCongestionSender(conn)
  225. description := fmt.Sprintf("remote=%s controller=%s sender=%x", remote, actual, sender)
  226. if _, exists := seen[description]; !exists {
  227. seen[description] = struct{}{}
  228. observed = append(observed, description)
  229. }
  230. if !matchesClientSocket(conn, client) {
  231. continue
  232. }
  233. if actual == want {
  234. server.connectionsMu.Unlock()
  235. return sender
  236. }
  237. }
  238. server.connectionsMu.Unlock()
  239. time.Sleep(5 * time.Millisecond)
  240. }
  241. t.Fatalf("server connection did not install %s congestion sender for client %s (observed %v)", want, client.LocalAddr(), observed)
  242. return 0
  243. }
  244. func congestionSenderForClient(t *testing.T, server *Server, client *quic.Conn) uintptr {
  245. t.Helper()
  246. server.connectionsMu.Lock()
  247. defer server.connectionsMu.Unlock()
  248. for conn := range server.connections {
  249. if matchesClientSocket(conn, client) {
  250. _, sender := inspectCongestionSender(conn)
  251. return sender
  252. }
  253. }
  254. t.Fatal("server connection for client is not registered")
  255. return 0
  256. }
  257. func matchesClientSocket(serverConn, clientConn *quic.Conn) bool {
  258. serverAddr, serverOK := serverConn.RemoteAddr().(*net.UDPAddr)
  259. clientAddr, clientOK := clientConn.LocalAddr().(*net.UDPAddr)
  260. return serverOK && clientOK && serverAddr.Port == clientAddr.Port
  261. }
  262. // lockedCongestion reads the sender under the mutex SetCongestionControl writes
  263. // it under, which the BBR install after the handshake races otherwise.
  264. func lockedCongestion(conn *quic.Conn) (reflect.Value, func()) {
  265. handler := reflect.ValueOf(conn).Elem().FieldByName("sentPacketHandler").Elem().Elem()
  266. mu := (*sync.RWMutex)(unsafe.Pointer(handler.FieldByName("congestionMutex").UnsafeAddr()))
  267. mu.RLock()
  268. return handler.FieldByName("congestion").Elem(), mu.RUnlock
  269. }
  270. func inspectCongestionSender(conn *quic.Conn) (string, uintptr) {
  271. controller, unlock := lockedCongestion(conn)
  272. defer unlock()
  273. sender := controller
  274. if controller.Type().String() == "*ackhandler.ccAdapterEx" || controller.Type().String() == "*ackhandler.ccAdapter" {
  275. sender = controller.Elem().FieldByName("CC").Elem()
  276. }
  277. switch {
  278. case strings.Contains(sender.Type().String(), "bbrSender"):
  279. return "bbr", sender.Pointer()
  280. case strings.Contains(sender.Type().String(), "cubicSender"):
  281. if sender.Elem().FieldByName("reno").Bool() {
  282. return "new_reno", sender.Pointer()
  283. }
  284. return "cubic", sender.Pointer()
  285. }
  286. return sender.Type().String(), sender.Pointer()
  287. }
  288. func waitForTuicLog(t *testing.T, message string) {
  289. t.Helper()
  290. marker := "inbound 14 (tuic-controller-reload): " + message
  291. deadline := time.Now().Add(4 * time.Second)
  292. for time.Now().Before(deadline) {
  293. if strings.Contains(strings.Join(logger.GetLogs(10000, "DEBUG"), "\n"), marker) {
  294. return
  295. }
  296. time.Sleep(5 * time.Millisecond)
  297. }
  298. t.Fatalf("timed out waiting for TUIC log %q", marker)
  299. }
  300. func TestCollectAllTraffic(t *testing.T) {
  301. certPEM, keyPEM := generateTestCert(t)
  302. pc, err := net.ListenPacket("udp", "127.0.0.1:0")
  303. if err != nil {
  304. t.Fatal(err)
  305. }
  306. port := pc.LocalAddr().(*net.UDPAddr).Port
  307. _ = pc.Close()
  308. inst := Instance{
  309. Id: 20,
  310. Tag: "tuic-20",
  311. Listen: "127.0.0.1",
  312. Port: port,
  313. Certificate: string(certPEM),
  314. PrivateKey: string(keyPEM),
  315. Clients: []TuicClientSettings{{UUID: "a0000000-0000-0000-0000-000000000001", Password: "p", Email: "e1"}},
  316. }
  317. m := &Manager{servers: map[int]*managed{}, lastStartErr: map[int]string{}}
  318. t.Cleanup(m.StopAll)
  319. if err := m.Ensure(inst); err != nil {
  320. t.Fatalf("Ensure failed: %v", err)
  321. }
  322. if !m.AddTestTraffic(20, "e1", 500, 1000) {
  323. t.Fatal("AddTestTraffic failed")
  324. }
  325. inbounds, clients := m.CollectAllTraffic()
  326. if len(inbounds) != 1 || inbounds[0].Up != 500 || inbounds[0].Down != 1000 {
  327. t.Fatalf("unexpected inbounds: %+v", inbounds)
  328. }
  329. if len(clients) != 1 || clients[0].Up != 500 || clients[0].Down != 1000 || clients[0].Email != "e1" {
  330. t.Fatalf("unexpected clients: %+v", clients)
  331. }
  332. // Subsequent call returns empty deltas
  333. inbounds2, clients2 := m.CollectAllTraffic()
  334. if len(inbounds2) != 0 || len(clients2) != 0 {
  335. t.Fatalf("expected empty deltas after drain, got %+v, %+v", inbounds2, clients2)
  336. }
  337. }
  338. func TestManagerRequeuesClientTrafficWithoutLosingDeltas(t *testing.T) {
  339. m := &Manager{servers: make(map[int]*managed)}
  340. m.RequeueClientTraffic([]ClientTrafficDelta{{Email: "[email protected]", Up: 10, Down: 20}})
  341. m.RequeueClientTraffic([]ClientTrafficDelta{{Email: "[email protected]", Up: 30, Down: 40}})
  342. _, got := m.CollectAllTraffic()
  343. if len(got) != 1 || got[0] != (ClientTrafficDelta{Email: "[email protected]", Up: 40, Down: 60}) {
  344. t.Fatalf("requeued client traffic = %+v", got)
  345. }
  346. if _, got = m.CollectAllTraffic(); len(got) != 0 {
  347. t.Fatalf("requeued traffic was collected more than once: %+v", got)
  348. }
  349. }