1
0

relay_e2e_test.go 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616
  1. package amneziawgnet
  2. import (
  3. "encoding/json"
  4. "fmt"
  5. "net"
  6. "net/netip"
  7. "os"
  8. "os/exec"
  9. "path/filepath"
  10. "strings"
  11. "sync"
  12. "testing"
  13. "time"
  14. "github.com/amnezia-vpn/amneziawg-go/v3/device"
  15. "github.com/amnezia-vpn/amneziawg-go/v3/tun/netstack"
  16. "gvisor.dev/gvisor/pkg/tcpip/adapters/gonet"
  17. "github.com/mhsanaei/3x-ui/v3/internal/amneziawg"
  18. "github.com/mhsanaei/3x-ui/v3/internal/util/wireguard"
  19. )
  20. // TestSocksRelayAgainstRealXray is Phase 2's real end-to-end proof: a
  21. // genuine amneziawg-go client completes a real handshake against a Device
  22. // built by NewDevice, dials a real TCP echo server and sends a real UDP
  23. // echo datagram, and this package's own AttachTCPForwarder/AttachUDPHandler
  24. // handlers relay both through RelayTCP/UDPRelay into an *actual xray-core
  25. // process* (not a mock) running a SOCKS5 inbound built by
  26. // SocksInboundSettings. Verifies real data round-trips on both protocols,
  27. // then greps the real process's own debug log for
  28. // "user>>>{email}>>>traffic>>>{up,down}link" -- the same proof Finding 3 of
  29. // the migration plan established manually in Phase 0, now permanent,
  30. // repo-owned test infrastructure. The UDP half in particular is the first
  31. // real test of this package's hand-rolled SOCKS5 UDP ASSOCIATE client
  32. // (relay.go) against an independent, authoritative implementation of the
  33. // protocol rather than a mock this same session wrote.
  34. //
  35. // Skipped unless XRAY_E2E_BINARY points at an xray executable built from
  36. // the same xray-core version as go.mod, matching internal/xray's own
  37. // TestXrayAPI_E2E convention:
  38. //
  39. // go install github.com/xtls/xray-core/main@<version from go.mod>
  40. // XRAY_E2E_BINARY=$GOBIN/main go test ./internal/amneziawgnet -run TestSocksRelayAgainstRealXray -v
  41. func TestSocksRelayAgainstRealXray(t *testing.T) {
  42. bin := os.Getenv("XRAY_E2E_BINARY")
  43. if bin == "" {
  44. t.Skip("set XRAY_E2E_BINARY to an xray binary to run this test")
  45. }
  46. localIP, ok := firstNonLoopbackIPv4()
  47. if !ok {
  48. t.Skip("no non-loopback IPv4 address available on this host")
  49. }
  50. const wantEmail = "[email protected]"
  51. const socksPassword = "loopback-only-not-a-real-secret"
  52. // --- real TCP + UDP echo servers on a real, non-loopback address ---
  53. // (dialing 127.0.0.1 as a tunnel-internal destination hangs -- gVisor
  54. // won't route loopback out an arbitrary NIC -- so the client dials
  55. // localIP instead; it must still be a *real* address since the actual
  56. // relay leg is a genuine OS-level dial from the xray-core process, not
  57. // anything inside the tunnel's virtual netstack.)
  58. tcpEcho, tcpEchoAddr := startTCPEcho(t, localIP)
  59. defer tcpEcho.Close()
  60. udpEcho, udpEchoAddr := startUDPEcho(t, localIP)
  61. defer udpEcho.Close()
  62. // --- real embedded AmneziaWG server + client, same shape as Phase 1's tests ---
  63. serverPriv, serverPub, err := wireguard.GenerateWireguardKeypair()
  64. if err != nil {
  65. t.Fatalf("generate server keypair: %v", err)
  66. }
  67. clientPriv, clientPub, err := wireguard.GenerateWireguardKeypair()
  68. if err != nil {
  69. t.Fatalf("generate client keypair: %v", err)
  70. }
  71. const listenPort = 58715
  72. inst := amneziawg.Instance{
  73. Id: 4,
  74. InterfaceName: "awgtest4",
  75. ListenPort: listenPort,
  76. PrivateKey: serverPriv,
  77. PublicKey: serverPub,
  78. Address: []string{"10.204.0.1/24"},
  79. MTU: 1420,
  80. Obfuscation: amneziawg.Obfuscation31{
  81. Jc: 4, Jmin: 40, Jmax: 70,
  82. S1: 20, S2: 30, S3: 20, S4: 20,
  83. },
  84. Peers: []amneziawg.Peer{{
  85. Email: wantEmail,
  86. PublicKey: clientPub,
  87. AllowedIPs: []string{"10.204.0.2/32"},
  88. }},
  89. }
  90. dev, err := newUnconfiguredDevice(inst, DeviceOptions{})
  91. if err != nil {
  92. t.Fatalf("newUnconfiguredDevice: %v", err)
  93. }
  94. defer dev.Close()
  95. idx := NewPeerIndex(inst.Peers)
  96. // --- real xray-core process with a SOCKS5 inbound built by this package ---
  97. socksPort := freePort(t)
  98. settingsJSON, err := SocksInboundSettings([]string{wantEmail}, socksPassword)
  99. if err != nil {
  100. t.Fatalf("SocksInboundSettings: %v", err)
  101. }
  102. var rawSettings any
  103. if err := json.Unmarshal(settingsJSON, &rawSettings); err != nil {
  104. t.Fatalf("unmarshal generated SOCKS5 settings: %v", err)
  105. }
  106. xrayCfg := map[string]any{
  107. "log": map[string]any{"loglevel": "debug"},
  108. "inbounds": []any{
  109. map[string]any{
  110. "listen": "127.0.0.1",
  111. "port": socksPort,
  112. "protocol": "socks",
  113. "settings": rawSettings,
  114. "tag": "awg-e2e-socks",
  115. },
  116. },
  117. "outbounds": []any{
  118. map[string]any{"protocol": "freedom", "settings": map[string]any{}, "tag": "direct"},
  119. },
  120. "policy": map[string]any{
  121. "levels": map[string]any{
  122. "0": map[string]any{"statsUserUplink": true, "statsUserDownlink": true},
  123. },
  124. },
  125. "stats": map[string]any{},
  126. }
  127. cfgBytes, err := json.MarshalIndent(xrayCfg, "", " ")
  128. if err != nil {
  129. t.Fatalf("marshal xray config: %v", err)
  130. }
  131. cfgPath := filepath.Join(t.TempDir(), "config.json")
  132. if err := os.WriteFile(cfgPath, cfgBytes, 0o644); err != nil {
  133. t.Fatalf("write xray config: %v", err)
  134. }
  135. var xrayLog syncBuffer
  136. cmd := exec.Command(bin, "-c", cfgPath)
  137. cmd.Stdout = &xrayLog
  138. cmd.Stderr = &xrayLog
  139. if err := cmd.Start(); err != nil {
  140. t.Fatalf("start xray: %v", err)
  141. }
  142. defer func() {
  143. _ = cmd.Process.Kill()
  144. _, _ = cmd.Process.Wait()
  145. }()
  146. waitForPort(t, socksPort)
  147. socksAddr := fmt.Sprintf("127.0.0.1:%d", socksPort)
  148. relay := SocksRelay{Addr: socksAddr, Password: socksPassword}
  149. udpRelay := NewUDPRelay(relay, dev.Stack)
  150. defer udpRelay.Close()
  151. AttachTCPForwarder(dev.Stack, func(conn *gonet.TCPConn, dest netip.AddrPort) {
  152. srcAddrPort, err := netip.ParseAddrPort(conn.RemoteAddr().String())
  153. if err != nil {
  154. conn.Close()
  155. return
  156. }
  157. peer, ok := idx.Lookup(srcAddrPort.Addr().Unmap())
  158. if !ok {
  159. conn.Close()
  160. return
  161. }
  162. relay.RelayTCP(conn, peer.Email, dest)
  163. })
  164. AttachUDPHandler(dev.Stack, func(src, dst netip.AddrPort, payload []byte) {
  165. peer, ok := idx.Lookup(src.Addr())
  166. if !ok {
  167. return
  168. }
  169. udpRelay.Handle(src, dst, peer.Email, payload)
  170. })
  171. // Configure (IpcSet) must come after both attaches -- see
  172. // newUnconfiguredDevice's doc comment.
  173. if err := dev.Configure(inst, DeviceOptions{}); err != nil {
  174. t.Fatalf("Configure: %v", err)
  175. }
  176. // --- real client, real handshake, real traffic through the whole chain ---
  177. clientTun, clientNet, err := netstack.CreateNetTUN(
  178. []netip.Addr{netip.MustParseAddr("10.204.0.2")},
  179. []netip.Addr{netip.MustParseAddr("1.1.1.1")}, 1420)
  180. if err != nil {
  181. t.Fatalf("client CreateNetTUN: %v", err)
  182. }
  183. clientDev := device.NewDevice(clientTun, newListenBind(""), device.NewLogger(device.LogLevelSilent, ""))
  184. defer clientDev.Close()
  185. clientPrivHex, err := wireguard.KeyToHex(clientPriv)
  186. if err != nil {
  187. t.Fatalf("client key to hex: %v", err)
  188. }
  189. serverPubHex, err := wireguard.KeyToHex(serverPub)
  190. if err != nil {
  191. t.Fatalf("server key to hex: %v", err)
  192. }
  193. clientConf := fmt.Sprintf(
  194. "private_key=%s\njc=4\njmin=40\njmax=70\ns1=20\ns2=30\ns3=20\ns4=20\npublic_key=%s\nendpoint=127.0.0.1:%d\nallowed_ip=0.0.0.0/0\n",
  195. clientPrivHex, serverPubHex, listenPort)
  196. if err := clientDev.IpcSet(clientConf); err != nil {
  197. t.Fatalf("client IpcSet: %v", err)
  198. }
  199. if err := clientDev.Up(); err != nil {
  200. t.Fatalf("client Up: %v", err)
  201. }
  202. // TCP round trip.
  203. const tcpMsg = "hello over amneziawgnet+socks5+xray"
  204. dialDeadline := time.Now().Add(10 * time.Second)
  205. var tcpConn interface {
  206. Write([]byte) (int, error)
  207. Read([]byte) (int, error)
  208. Close() error
  209. }
  210. for {
  211. c, dialErr := clientNet.DialContext(t.Context(), "tcp", tcpEchoAddr.String())
  212. if dialErr == nil {
  213. tcpConn = c
  214. break
  215. }
  216. if time.Now().After(dialDeadline) {
  217. t.Fatalf("client TCP dial via tunnel never succeeded: %v", dialErr)
  218. }
  219. time.Sleep(150 * time.Millisecond)
  220. }
  221. defer tcpConn.Close()
  222. if _, err := tcpConn.Write([]byte(tcpMsg)); err != nil {
  223. t.Fatalf("client TCP write: %v", err)
  224. }
  225. tcpBuf := make([]byte, len(tcpMsg))
  226. if _, err := readFull(tcpConn, tcpBuf, 10*time.Second); err != nil {
  227. t.Fatalf("client TCP read: %v", err)
  228. }
  229. if string(tcpBuf) != tcpMsg {
  230. t.Errorf("TCP echo = %q, want %q", tcpBuf, tcpMsg)
  231. }
  232. // UDP round trip.
  233. const udpMsg = "hello-udp-over-socks5"
  234. uconn, err := clientNet.DialUDPAddrPort(netip.AddrPort{}, udpEchoAddr)
  235. if err != nil {
  236. t.Fatalf("client DialUDPAddrPort: %v", err)
  237. }
  238. defer uconn.Close()
  239. udpDeadline := time.Now().Add(10 * time.Second)
  240. var udpBuf [256]byte
  241. var gotUDP string
  242. for time.Now().Before(udpDeadline) {
  243. _ = uconn.SetWriteDeadline(time.Now().Add(300 * time.Millisecond))
  244. if _, err := uconn.Write([]byte(udpMsg)); err != nil {
  245. continue
  246. }
  247. _ = uconn.SetReadDeadline(time.Now().Add(300 * time.Millisecond))
  248. n, err := uconn.Read(udpBuf[:])
  249. if err == nil {
  250. gotUDP = string(udpBuf[:n])
  251. break
  252. }
  253. }
  254. if gotUDP != udpMsg {
  255. t.Fatalf("UDP echo = %q, want %q (xray log follows)\n%s", gotUDP, udpMsg, xrayLog.String())
  256. }
  257. // Real per-peer stats attribution: stop xray so its log is complete, then
  258. // look for both directions' counters keyed by the peer's real email --
  259. // the exact proof Finding 3 established manually in Phase 0.
  260. _ = cmd.Process.Kill()
  261. _, _ = cmd.Process.Wait()
  262. log := xrayLog.String()
  263. wantUp := fmt.Sprintf("user>>>%s>>>traffic>>>uplink", wantEmail)
  264. wantDown := fmt.Sprintf("user>>>%s>>>traffic>>>downlink", wantEmail)
  265. if !strings.Contains(log, wantUp) {
  266. t.Errorf("xray log missing uplink stats counter %q\nfull log:\n%s", wantUp, log)
  267. }
  268. if !strings.Contains(log, wantDown) {
  269. t.Errorf("xray log missing downlink stats counter %q\nfull log:\n%s", wantDown, log)
  270. }
  271. }
  272. // TestManagerEnsureAutomaticallyWiresRelay is Phase 3's own real proof: unlike
  273. // TestSocksRelayAgainstRealXray above (which builds a Device and attaches
  274. // RelayTCP/UDPRelay by hand), this drives everything through the public
  275. // Manager.Ensure entry point the real app actually calls -- confirming
  276. // ensureLocked's own forwarder/UDP-handler attachment (added this phase)
  277. // really does relay a fresh Device's traffic into Xray with zero manual
  278. // wiring from the caller. Uses the exact port/password
  279. // (SOCKSPortForInbound/SocksPassword) the Manager computes internally, so
  280. // this only passes if that internal derivation and the externally-visible
  281. // contract genuinely agree.
  282. func TestManagerEnsureAutomaticallyWiresRelay(t *testing.T) {
  283. bin := os.Getenv("XRAY_E2E_BINARY")
  284. if bin == "" {
  285. t.Skip("set XRAY_E2E_BINARY to an xray binary to run this test")
  286. }
  287. localIP, ok := firstNonLoopbackIPv4()
  288. if !ok {
  289. t.Skip("no non-loopback IPv4 address available on this host")
  290. }
  291. const wantEmail = "[email protected]"
  292. const listenPort = 58716
  293. const inboundID = 5
  294. tcpEcho, tcpEchoAddr := startTCPEcho(t, localIP)
  295. defer tcpEcho.Close()
  296. serverPriv, serverPub, err := wireguard.GenerateWireguardKeypair()
  297. if err != nil {
  298. t.Fatalf("generate server keypair: %v", err)
  299. }
  300. clientPriv, clientPub, err := wireguard.GenerateWireguardKeypair()
  301. if err != nil {
  302. t.Fatalf("generate client keypair: %v", err)
  303. }
  304. inst := amneziawg.Instance{
  305. Id: inboundID,
  306. InterfaceName: "awgtest5",
  307. ListenPort: listenPort,
  308. PrivateKey: serverPriv,
  309. PublicKey: serverPub,
  310. Address: []string{"10.205.0.1/24"},
  311. MTU: 1420,
  312. Obfuscation: amneziawg.Obfuscation31{
  313. Jc: 4, Jmin: 40, Jmax: 70,
  314. S1: 20, S2: 30, S3: 20, S4: 20,
  315. },
  316. Peers: []amneziawg.Peer{{
  317. Email: wantEmail,
  318. PublicKey: clientPub,
  319. AllowedIPs: []string{"10.205.0.2/32"},
  320. }},
  321. }
  322. // A real xray-core process with a SOCKS5 inbound at exactly the port and
  323. // password ensureLocked will derive on its own for this instance --
  324. // SocksPassword() is cached (sync.Once), so calling it here first and
  325. // again inside Manager.Ensure below returns the identical value.
  326. socksPort := SOCKSPortForInbound(inboundID)
  327. password := SocksPassword()
  328. settingsJSON, err := SocksInboundSettings([]string{wantEmail}, password)
  329. if err != nil {
  330. t.Fatalf("SocksInboundSettings: %v", err)
  331. }
  332. var rawSettings any
  333. if err := json.Unmarshal(settingsJSON, &rawSettings); err != nil {
  334. t.Fatalf("unmarshal generated SOCKS5 settings: %v", err)
  335. }
  336. xrayCfg := map[string]any{
  337. "log": map[string]any{"loglevel": "debug"},
  338. "inbounds": []any{
  339. map[string]any{
  340. "listen": "127.0.0.1",
  341. "port": socksPort,
  342. "protocol": "socks",
  343. "settings": rawSettings,
  344. "tag": "awg-e2e-manager",
  345. },
  346. },
  347. "outbounds": []any{
  348. map[string]any{"protocol": "freedom", "settings": map[string]any{}, "tag": "direct"},
  349. },
  350. "policy": map[string]any{
  351. "levels": map[string]any{
  352. "0": map[string]any{"statsUserUplink": true, "statsUserDownlink": true},
  353. },
  354. },
  355. "stats": map[string]any{},
  356. }
  357. cfgBytes, err := json.MarshalIndent(xrayCfg, "", " ")
  358. if err != nil {
  359. t.Fatalf("marshal xray config: %v", err)
  360. }
  361. cfgPath := filepath.Join(t.TempDir(), "config.json")
  362. if err := os.WriteFile(cfgPath, cfgBytes, 0o644); err != nil {
  363. t.Fatalf("write xray config: %v", err)
  364. }
  365. var xrayLog syncBuffer
  366. cmd := exec.Command(bin, "-c", cfgPath)
  367. cmd.Stdout = &xrayLog
  368. cmd.Stderr = &xrayLog
  369. if err := cmd.Start(); err != nil {
  370. t.Fatalf("start xray: %v", err)
  371. }
  372. defer func() {
  373. _ = cmd.Process.Kill()
  374. _, _ = cmd.Process.Wait()
  375. }()
  376. waitForPort(t, socksPort)
  377. // A throwaway Manager, not the process-wide singleton, so this test
  378. // doesn't interact with any other test's state.
  379. m := &Manager{ifaces: map[int]*managed{}}
  380. defer m.StopAll()
  381. if err := m.Ensure(Desired{Instance: inst}); err != nil {
  382. t.Fatalf("Manager.Ensure: %v", err)
  383. }
  384. dev, _, ok := m.Lookup(inboundID)
  385. if !ok {
  386. t.Fatal("Lookup after Ensure: not found")
  387. }
  388. defer dev.Close() // StopAll would also do this; explicit for clarity
  389. clientTun, clientNet, err := netstack.CreateNetTUN(
  390. []netip.Addr{netip.MustParseAddr("10.205.0.2")},
  391. []netip.Addr{netip.MustParseAddr("1.1.1.1")}, 1420)
  392. if err != nil {
  393. t.Fatalf("client CreateNetTUN: %v", err)
  394. }
  395. clientDev := device.NewDevice(clientTun, newListenBind(""), device.NewLogger(device.LogLevelSilent, ""))
  396. defer clientDev.Close()
  397. clientPrivHex, err := wireguard.KeyToHex(clientPriv)
  398. if err != nil {
  399. t.Fatalf("client key to hex: %v", err)
  400. }
  401. serverPubHex, err := wireguard.KeyToHex(serverPub)
  402. if err != nil {
  403. t.Fatalf("server key to hex: %v", err)
  404. }
  405. clientConf := fmt.Sprintf(
  406. "private_key=%s\njc=4\njmin=40\njmax=70\ns1=20\ns2=30\ns3=20\ns4=20\npublic_key=%s\nendpoint=127.0.0.1:%d\nallowed_ip=0.0.0.0/0\n",
  407. clientPrivHex, serverPubHex, listenPort)
  408. if err := clientDev.IpcSet(clientConf); err != nil {
  409. t.Fatalf("client IpcSet: %v", err)
  410. }
  411. if err := clientDev.Up(); err != nil {
  412. t.Fatalf("client Up: %v", err)
  413. }
  414. const tcpMsg = "hello via Manager.Ensure's automatic relay wiring"
  415. dialDeadline := time.Now().Add(10 * time.Second)
  416. var conn net.Conn
  417. for {
  418. c, dialErr := clientNet.DialContext(t.Context(), "tcp", tcpEchoAddr.String())
  419. if dialErr == nil {
  420. conn = c
  421. break
  422. }
  423. if time.Now().After(dialDeadline) {
  424. t.Fatalf("client TCP dial via tunnel never succeeded: %v", dialErr)
  425. }
  426. time.Sleep(150 * time.Millisecond)
  427. }
  428. defer conn.Close()
  429. if _, err := conn.Write([]byte(tcpMsg)); err != nil {
  430. t.Fatalf("client TCP write: %v", err)
  431. }
  432. buf := make([]byte, len(tcpMsg))
  433. if _, err := readFull(conn, buf, 10*time.Second); err != nil {
  434. t.Fatalf("client TCP read: %v", err)
  435. }
  436. if string(buf) != tcpMsg {
  437. t.Errorf("TCP echo = %q, want %q", buf, tcpMsg)
  438. }
  439. _ = cmd.Process.Kill()
  440. _, _ = cmd.Process.Wait()
  441. log := xrayLog.String()
  442. wantUp := fmt.Sprintf("user>>>%s>>>traffic>>>uplink", wantEmail)
  443. if !strings.Contains(log, wantUp) {
  444. t.Errorf("xray log missing uplink stats counter %q (Manager.Ensure's automatic relay wiring may not be attributing traffic correctly)\nfull log:\n%s", wantUp, log)
  445. }
  446. }
  447. // firstNonLoopbackIPv4 finds a real, locally-bound IPv4 address suitable as
  448. // a relay-reachable test destination.
  449. func firstNonLoopbackIPv4() (netip.Addr, bool) {
  450. addrs, err := net.InterfaceAddrs()
  451. if err != nil {
  452. return netip.Addr{}, false
  453. }
  454. for _, a := range addrs {
  455. ipNet, ok := a.(*net.IPNet)
  456. if !ok || ipNet.IP.IsLoopback() {
  457. continue
  458. }
  459. if v4 := ipNet.IP.To4(); v4 != nil {
  460. addr, ok := netip.AddrFromSlice(v4)
  461. if ok {
  462. return addr, true
  463. }
  464. }
  465. }
  466. return netip.Addr{}, false
  467. }
  468. func startTCPEcho(t *testing.T, addr netip.Addr) (io interface{ Close() error }, ap netip.AddrPort) {
  469. t.Helper()
  470. ln, err := net.Listen("tcp", net.JoinHostPort(addr.String(), "0"))
  471. if err != nil {
  472. t.Fatalf("start TCP echo listener: %v", err)
  473. }
  474. go func() {
  475. for {
  476. c, err := ln.Accept()
  477. if err != nil {
  478. return
  479. }
  480. go func() {
  481. defer c.Close()
  482. buf := make([]byte, 4096)
  483. for {
  484. n, err := c.Read(buf)
  485. if n > 0 {
  486. if _, werr := c.Write(buf[:n]); werr != nil {
  487. return
  488. }
  489. }
  490. if err != nil {
  491. return
  492. }
  493. }
  494. }()
  495. }
  496. }()
  497. port := ln.Addr().(*net.TCPAddr).Port
  498. return ln, netip.AddrPortFrom(addr, uint16(port))
  499. }
  500. func startUDPEcho(t *testing.T, addr netip.Addr) (io interface{ Close() error }, ap netip.AddrPort) {
  501. t.Helper()
  502. pc, err := net.ListenPacket("udp", net.JoinHostPort(addr.String(), "0"))
  503. if err != nil {
  504. t.Fatalf("start UDP echo listener: %v", err)
  505. }
  506. go func() {
  507. buf := make([]byte, 4096)
  508. for {
  509. n, raddr, err := pc.ReadFrom(buf)
  510. if err != nil {
  511. return
  512. }
  513. if _, err := pc.WriteTo(buf[:n], raddr); err != nil {
  514. return
  515. }
  516. }
  517. }()
  518. port := pc.LocalAddr().(*net.UDPAddr).Port
  519. return pc, netip.AddrPortFrom(addr, uint16(port))
  520. }
  521. // readFull reads exactly len(buf) bytes or fails after timeout, since
  522. // gonet.TCPConn (and net.Conn generally) may return short reads.
  523. func readFull(r interface{ Read([]byte) (int, error) }, buf []byte, timeout time.Duration) (int, error) {
  524. deadline := time.Now().Add(timeout)
  525. total := 0
  526. for total < len(buf) {
  527. if time.Now().After(deadline) {
  528. return total, fmt.Errorf("timed out after reading %d/%d bytes", total, len(buf))
  529. }
  530. n, err := r.Read(buf[total:])
  531. total += n
  532. if err != nil {
  533. return total, err
  534. }
  535. }
  536. return total, nil
  537. }
  538. // syncBuffer is a concurrency-safe bytes buffer for capturing a subprocess's
  539. // combined stdout/stderr while the test may read it from another goroutine.
  540. type syncBuffer struct {
  541. mu sync.Mutex
  542. buf strings.Builder
  543. }
  544. func (s *syncBuffer) Write(p []byte) (int, error) {
  545. s.mu.Lock()
  546. defer s.mu.Unlock()
  547. return s.buf.Write(p)
  548. }
  549. func (s *syncBuffer) String() string {
  550. s.mu.Lock()
  551. defer s.mu.Unlock()
  552. return s.buf.String()
  553. }
  554. func freePort(t *testing.T) int {
  555. t.Helper()
  556. l, err := net.Listen("tcp", "127.0.0.1:0")
  557. if err != nil {
  558. t.Fatal(err)
  559. }
  560. defer l.Close()
  561. return l.Addr().(*net.TCPAddr).Port
  562. }
  563. func waitForPort(t *testing.T, port int) {
  564. t.Helper()
  565. deadline := time.Now().Add(15 * time.Second)
  566. addr := fmt.Sprintf("127.0.0.1:%d", port)
  567. for time.Now().Before(deadline) {
  568. conn, err := net.DialTimeout("tcp", addr, time.Second)
  569. if err == nil {
  570. conn.Close()
  571. return
  572. }
  573. time.Sleep(200 * time.Millisecond)
  574. }
  575. t.Fatalf("xray port %d did not open in time", port)
  576. }