| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196 |
- package tuic
- import (
- "errors"
- "net"
- "sync"
- "sync/atomic"
- "time"
- )
- // A QUIC flow the sidecar has not touched for this long is forgotten; QUIC's
- // own max_idle_time (15s by default) closes the session well before that.
- const relayFlowIdle = 2 * time.Minute
- const (
- relaySocketBuffer = 4 << 20
- maxRelayFlows = 4096
- )
- // udpRelay owns an inbound's public UDP port and forwards each client's
- // datagrams to the sidecar on loopback, which is the only place the panel can
- // count the inbound's bytes: upstream tuic-server exposes no stats API and
- // its socket syscalls never reach /proc/<pid>/io. Per-client attribution stays
- // impossible because QUIC payloads are opaque.
- type udpRelay struct {
- public *net.UDPConn
- upstream *net.UDPAddr
- idle time.Duration
- up atomic.Int64
- down atomic.Int64
- mu sync.Mutex
- flows map[string]*relayFlow
- done chan struct{}
- closeOnce sync.Once
- wg sync.WaitGroup
- }
- type relayFlow struct {
- conn *net.UDPConn
- client *net.UDPAddr
- lastSeen atomic.Int64
- }
- func startUDPRelay(bind string, upstream *net.UDPAddr, idle time.Duration) (*udpRelay, error) {
- addr, err := net.ResolveUDPAddr("udp", bind)
- if err != nil {
- return nil, err
- }
- public, err := net.ListenUDP("udp", addr)
- if err != nil {
- return nil, err
- }
- _ = public.SetReadBuffer(relaySocketBuffer)
- _ = public.SetWriteBuffer(relaySocketBuffer)
- r := &udpRelay{
- public: public,
- upstream: upstream,
- idle: idle,
- flows: make(map[string]*relayFlow),
- done: make(chan struct{}),
- }
- r.wg.Add(2)
- go r.serve()
- go r.sweep()
- return r, nil
- }
- func freeLoopbackUDPPort() (int, error) {
- c, err := net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)})
- if err != nil {
- return 0, err
- }
- defer c.Close()
- return c.LocalAddr().(*net.UDPAddr).Port, nil
- }
- func (r *udpRelay) LocalAddr() net.Addr {
- return r.public.LocalAddr()
- }
- // CollectTraffic returns the client-to-sidecar and sidecar-to-client bytes
- // relayed since the previous call.
- func (r *udpRelay) CollectTraffic() (up, down int64) {
- return r.up.Swap(0), r.down.Swap(0)
- }
- func (r *udpRelay) Close() {
- if r == nil {
- return
- }
- r.closeOnce.Do(func() {
- close(r.done)
- _ = r.public.Close()
- r.mu.Lock()
- for key, f := range r.flows {
- _ = f.conn.Close()
- delete(r.flows, key)
- }
- r.mu.Unlock()
- r.wg.Wait()
- })
- }
- func (r *udpRelay) serve() {
- defer r.wg.Done()
- buf := make([]byte, 65535)
- for {
- n, client, err := r.public.ReadFromUDP(buf)
- if err != nil {
- if errors.Is(err, net.ErrClosed) {
- return
- }
- continue
- }
- flow, err := r.flowFor(client)
- if err != nil {
- continue
- }
- if _, err := flow.conn.Write(buf[:n]); err == nil {
- r.up.Add(int64(n))
- }
- }
- }
- func (r *udpRelay) flowFor(client *net.UDPAddr) (*relayFlow, error) {
- key := client.String()
- now := time.Now().UnixMilli()
- r.mu.Lock()
- defer r.mu.Unlock()
- select {
- case <-r.done:
- return nil, net.ErrClosed
- default:
- }
- if f, ok := r.flows[key]; ok {
- f.lastSeen.Store(now)
- return f, nil
- }
- if len(r.flows) >= maxRelayFlows {
- return nil, errors.New("tuic: max relay flows reached")
- }
- conn, err := net.DialUDP("udp", nil, r.upstream)
- if err != nil {
- return nil, err
- }
- _ = conn.SetReadBuffer(relaySocketBuffer)
- _ = conn.SetWriteBuffer(relaySocketBuffer)
- f := &relayFlow{conn: conn, client: client}
- f.lastSeen.Store(now)
- r.flows[key] = f
- r.wg.Add(1)
- go r.pump(f)
- return f, nil
- }
- func (r *udpRelay) pump(f *relayFlow) {
- defer r.wg.Done()
- buf := make([]byte, 65535)
- for {
- n, err := f.conn.Read(buf)
- if err != nil {
- if errors.Is(err, net.ErrClosed) {
- return
- }
- // ICMP unreachable while the sidecar restarts: drop it, keep the flow.
- time.Sleep(20 * time.Millisecond)
- continue
- }
- if _, err := r.public.WriteToUDP(buf[:n], f.client); err == nil {
- r.down.Add(int64(n))
- }
- f.lastSeen.Store(time.Now().UnixMilli())
- }
- }
- func (r *udpRelay) sweep() {
- defer r.wg.Done()
- ticker := time.NewTicker(r.idle / 2)
- defer ticker.Stop()
- for {
- select {
- case <-r.done:
- return
- case <-ticker.C:
- cutoff := time.Now().Add(-r.idle).UnixMilli()
- r.mu.Lock()
- for key, f := range r.flows {
- if f.lastSeen.Load() < cutoff {
- _ = f.conn.Close()
- delete(r.flows, key)
- }
- }
- r.mu.Unlock()
- }
- }
- }
|