relay.go 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196
  1. package tuic
  2. import (
  3. "errors"
  4. "net"
  5. "sync"
  6. "sync/atomic"
  7. "time"
  8. )
  9. // A QUIC flow the sidecar has not touched for this long is forgotten; QUIC's
  10. // own max_idle_time (15s by default) closes the session well before that.
  11. const relayFlowIdle = 2 * time.Minute
  12. const (
  13. relaySocketBuffer = 4 << 20
  14. maxRelayFlows = 4096
  15. )
  16. // udpRelay owns an inbound's public UDP port and forwards each client's
  17. // datagrams to the sidecar on loopback, which is the only place the panel can
  18. // count the inbound's bytes: upstream tuic-server exposes no stats API and
  19. // its socket syscalls never reach /proc/<pid>/io. Per-client attribution stays
  20. // impossible because QUIC payloads are opaque.
  21. type udpRelay struct {
  22. public *net.UDPConn
  23. upstream *net.UDPAddr
  24. idle time.Duration
  25. up atomic.Int64
  26. down atomic.Int64
  27. mu sync.Mutex
  28. flows map[string]*relayFlow
  29. done chan struct{}
  30. closeOnce sync.Once
  31. wg sync.WaitGroup
  32. }
  33. type relayFlow struct {
  34. conn *net.UDPConn
  35. client *net.UDPAddr
  36. lastSeen atomic.Int64
  37. }
  38. func startUDPRelay(bind string, upstream *net.UDPAddr, idle time.Duration) (*udpRelay, error) {
  39. addr, err := net.ResolveUDPAddr("udp", bind)
  40. if err != nil {
  41. return nil, err
  42. }
  43. public, err := net.ListenUDP("udp", addr)
  44. if err != nil {
  45. return nil, err
  46. }
  47. _ = public.SetReadBuffer(relaySocketBuffer)
  48. _ = public.SetWriteBuffer(relaySocketBuffer)
  49. r := &udpRelay{
  50. public: public,
  51. upstream: upstream,
  52. idle: idle,
  53. flows: make(map[string]*relayFlow),
  54. done: make(chan struct{}),
  55. }
  56. r.wg.Add(2)
  57. go r.serve()
  58. go r.sweep()
  59. return r, nil
  60. }
  61. func freeLoopbackUDPPort() (int, error) {
  62. c, err := net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)})
  63. if err != nil {
  64. return 0, err
  65. }
  66. defer c.Close()
  67. return c.LocalAddr().(*net.UDPAddr).Port, nil
  68. }
  69. func (r *udpRelay) LocalAddr() net.Addr {
  70. return r.public.LocalAddr()
  71. }
  72. // CollectTraffic returns the client-to-sidecar and sidecar-to-client bytes
  73. // relayed since the previous call.
  74. func (r *udpRelay) CollectTraffic() (up, down int64) {
  75. return r.up.Swap(0), r.down.Swap(0)
  76. }
  77. func (r *udpRelay) Close() {
  78. if r == nil {
  79. return
  80. }
  81. r.closeOnce.Do(func() {
  82. close(r.done)
  83. _ = r.public.Close()
  84. r.mu.Lock()
  85. for key, f := range r.flows {
  86. _ = f.conn.Close()
  87. delete(r.flows, key)
  88. }
  89. r.mu.Unlock()
  90. r.wg.Wait()
  91. })
  92. }
  93. func (r *udpRelay) serve() {
  94. defer r.wg.Done()
  95. buf := make([]byte, 65535)
  96. for {
  97. n, client, err := r.public.ReadFromUDP(buf)
  98. if err != nil {
  99. if errors.Is(err, net.ErrClosed) {
  100. return
  101. }
  102. continue
  103. }
  104. flow, err := r.flowFor(client)
  105. if err != nil {
  106. continue
  107. }
  108. if _, err := flow.conn.Write(buf[:n]); err == nil {
  109. r.up.Add(int64(n))
  110. }
  111. }
  112. }
  113. func (r *udpRelay) flowFor(client *net.UDPAddr) (*relayFlow, error) {
  114. key := client.String()
  115. now := time.Now().UnixMilli()
  116. r.mu.Lock()
  117. defer r.mu.Unlock()
  118. select {
  119. case <-r.done:
  120. return nil, net.ErrClosed
  121. default:
  122. }
  123. if f, ok := r.flows[key]; ok {
  124. f.lastSeen.Store(now)
  125. return f, nil
  126. }
  127. if len(r.flows) >= maxRelayFlows {
  128. return nil, errors.New("tuic: max relay flows reached")
  129. }
  130. conn, err := net.DialUDP("udp", nil, r.upstream)
  131. if err != nil {
  132. return nil, err
  133. }
  134. _ = conn.SetReadBuffer(relaySocketBuffer)
  135. _ = conn.SetWriteBuffer(relaySocketBuffer)
  136. f := &relayFlow{conn: conn, client: client}
  137. f.lastSeen.Store(now)
  138. r.flows[key] = f
  139. r.wg.Add(1)
  140. go r.pump(f)
  141. return f, nil
  142. }
  143. func (r *udpRelay) pump(f *relayFlow) {
  144. defer r.wg.Done()
  145. buf := make([]byte, 65535)
  146. for {
  147. n, err := f.conn.Read(buf)
  148. if err != nil {
  149. if errors.Is(err, net.ErrClosed) {
  150. return
  151. }
  152. // ICMP unreachable while the sidecar restarts: drop it, keep the flow.
  153. time.Sleep(20 * time.Millisecond)
  154. continue
  155. }
  156. if _, err := r.public.WriteToUDP(buf[:n], f.client); err == nil {
  157. r.down.Add(int64(n))
  158. }
  159. f.lastSeen.Store(time.Now().UnixMilli())
  160. }
  161. }
  162. func (r *udpRelay) sweep() {
  163. defer r.wg.Done()
  164. ticker := time.NewTicker(r.idle / 2)
  165. defer ticker.Stop()
  166. for {
  167. select {
  168. case <-r.done:
  169. return
  170. case <-ticker.C:
  171. cutoff := time.Now().Add(-r.idle).UnixMilli()
  172. r.mu.Lock()
  173. for key, f := range r.flows {
  174. if f.lastSeen.Load() < cutoff {
  175. _ = f.conn.Close()
  176. delete(r.flows, key)
  177. }
  178. }
  179. r.mu.Unlock()
  180. }
  181. }
  182. }