Files
3x-ui/internal/tuic/socks_bridge.go
T
Egor 0054e671f8 feat(tuic): implement native in-process Go TUIC v5 server (#6577)
* feat(tuic): implement native in-process Go TUIC v5 server

- Implement native TUIC v5 protocol server on pure Go using quic-go
- Bridge decrypted TCP/UDP traffic into Xray-core via loopback SOCKS5 inbound
- Support full Xray routing rules (geosite/geoip) and cascading outbounds
- Implement atomic per-client traffic accounting with TotalGB and ExpiryTime
- Add automatic legacy cleanup for older Rust tuic-server binaries, configs, and orphaned processes
- Eliminate external Rust tuic-server downloads from install/CI scripts

* fix(tuic): address traffic accounting, client reload, and socket lifecycle issues

* fix(service): update checkTuicSocksReverseConflict to use bindAddr for listenOverlaps

* fix(tuic): resolve traffic double-accounting, UDP fragmentation, and socket lifecycle issues

* feat(tuic): complete native Go integration and address audit findings

- Integrate an isolated QUIC fork pinned to a specific commit
- Preserve original QUIC dependencies for Xray, Hysteria and Gin
- Apply BBR, CUBIC and Reno to server connections and exported client profiles
- Bridge Xray BBR with correct monotonic time and congestion type conversions
- Handle congestion sender recreation after PMTU changes
- Update congestion control for new connections without restarting the listener
- Preserve existing connections and their selected congestion controller
- Apply per-inbound log levels through the shared panel logger
- Add lifecycle, authentication and TCP/UDP relay events without exposing secrets
- Rate-limit repeated authentication and relay warnings
- Support native and QUIC UDP relay modes on the same listener
- Recover UDP associations after relay worker failures
- Fix TCP relay cancellation, idle shutdown and half-close handling
- Close active sessions when client credentials are revoked or disabled
- Track traffic by immutable client statistics IDs across email and UUID changes
- Prevent ambiguous accounting and duplicate UUIDs within TUIC inbounds
- Persist pending traffic in a durable shutdown journal
- Replay journal batches transactionally without duplicate accounting
- Report server shutdown failures through the shared logger
- Preserve legacy flat and nested TUIC settings compatibility
- Normalize congestion controller values consistently across backend and frontend
- Preserve controller, UDP mode and SNI in client links and subscriptions
- Separate client profile options from server settings in the TUIC form
- Keep certificate path autofill explicit when changing client SNI
- Align UDP packet size validation with protocol limits
- Simplify and localize TUIC field hints and certificate autofill messages
- Add controller, TCP/UDP, logging and live settings update tests
- Add accounting identity, journal replay and shutdown regression tests
- Add relay recovery, session revocation and legacy frontend form tests

* fix(service): alias the TUIC duplicate-UUID subquery for PostgreSQL < 16

syncInboundClients runs a COUNT(*) FROM (subquery) for every client sync,
whatever the protocol. PostgreSQL before 16 rejects a FROM subquery with
no alias, so on the distro PostgreSQL install.sh provisions (14 on Ubuntu
22.04, 15 on Debian 12) every client add or edit failed with SQLSTATE
42601. Reproduced against postgres:15 with the new env-gated test.

* fix(database): create tuic_traffic_receipts through the model migration

AddTuicTrafficBatch issued CREATE TABLE IF NOT EXISTS at runtime, a schema
change outside db.go. The table was invisible to allModels and
migrationModels, so x-ui migrate-db dropped the receipts and a retained
journal could be counted twice after a SQLite to PostgreSQL move. It is
now a GORM model in both lists, and the insert uses OnConflict DoNothing.

* chore(tuic): skip the ICMP-dependent relay test on Windows, drop dead collectors

Go disables SIO_UDP_CONNRESET on Windows, so a dead UDP bridge never fails
a read there and TestAudit3UDPAssociationMustRecoverAfterBridgeReadFailure
was red on every Windows run. Server.CollectTotalTraffic and
Manager.CollectTraffic had no caller.

* refactor(tuic): serve TUIC on apernet/quic-go instead of a personal fork

The native server depended on github.com/poise52/quic-go, a personal fork
of apernet/quic-go patched only to pick the congestion controller before
the handshake. That put a second QUIC/TLS stack in the binary that no
upstream security fix reaches. apernet/quic-go is already in the graph
through xray-core and exposes SetCongestionControl, so BBR is now installed
on each accepted connection with Xray's own congestion.UseBBR; the
cross-module BBR adapter is gone.

apernet ships New Reno as its only built-in sender, so a cubic setting is
served as new_reno server-side (clients still get cubic in their profile).
The test inspectors now read the sender under congestionMutex, which the
post-handshake install writes under.

Linux loopback, single stream through Xray: 2428 -> 3383 Mbit/s (bbr).

---------

Co-authored-by: Sanaei <ho3ein.sanaei@gmail.com>
2026-10-03 02:56:02 +02:00

579 lines
15 KiB
Go

package tuic
import (
"context"
"crypto/rand"
"encoding/base64"
"encoding/binary"
"encoding/json"
"errors"
"fmt"
"io"
"net"
"net/netip"
"sync"
"sync/atomic"
"time"
)
var ErrUdpPayloadTooLarge = errors.New("tuic socks: UDP packet exceeds the maximum SOCKS datagram size")
const maxSocksUdpDatagramSize = 65507
// SocksRelay describes the loopback SOCKS5 endpoint where decrypted TUIC traffic is forwarded.
type SocksRelay struct {
Addr string // e.g. "127.0.0.1:63201"
Password string // internal shared password for the SOCKS inbound
}
// CountingConn wraps a net.Conn and tracks bytes read and written atomically.
type CountingConn struct {
net.Conn
bytesRead *atomic.Int64
bytesWritten *atomic.Int64
}
func (c *CountingConn) Read(p []byte) (int, error) {
n, err := c.Conn.Read(p)
if n > 0 && c.bytesRead != nil {
c.bytesRead.Add(int64(n))
}
return n, err
}
func (c *CountingConn) Write(p []byte) (int, error) {
n, err := c.Conn.Write(p)
if n > 0 && c.bytesWritten != nil {
c.bytesWritten.Add(int64(n))
}
return n, err
}
// DialTCP establishes a SOCKS5 CONNECT tunnel to the target address on behalf of user.
func (r *SocksRelay) DialTCP(ctx context.Context, user string, target *Address) (net.Conn, error) {
dialer := net.Dialer{Timeout: 10 * time.Second}
conn, err := dialer.DialContext(ctx, "tcp", r.Addr)
if err != nil {
return nil, fmt.Errorf("tuic socks: dial relay %s: %w", r.Addr, err)
}
_ = conn.SetDeadline(time.Now().Add(10 * time.Second))
if err := socks5Handshake(conn, user, r.Password); err != nil {
conn.Close()
return nil, err
}
// Send SOCKS5 CONNECT request
req := buildSocks5ConnectRequest(target)
if req == nil {
conn.Close()
return nil, ErrInvalidAddr
}
if _, err := conn.Write(req); err != nil {
conn.Close()
return nil, fmt.Errorf("tuic socks: send CONNECT request: %w", err)
}
if _, err := readSocks5Reply(conn); err != nil {
conn.Close()
return nil, err
}
_ = conn.SetDeadline(time.Time{})
return conn, nil
}
func buildSocks5ConnectRequest(target *Address) []byte {
if target == nil {
return nil
}
var req []byte
switch target.Type {
case AddrTypeIPv4:
ip4 := target.IP.To4()
if len(ip4) != 4 {
return nil
}
req = make([]byte, 4+4+2)
req[0] = 0x05 // SOCKS5
req[1] = 0x01 // CONNECT
req[2] = 0x00 // RSV
req[3] = 0x01 // ATYP IPv4
copy(req[4:8], ip4)
binary.BigEndian.PutUint16(req[8:10], target.Port)
case AddrTypeIPv6:
ip16 := target.IP.To16()
if len(ip16) != 16 {
return nil
}
req = make([]byte, 4+16+2)
req[0] = 0x05
req[1] = 0x01
req[2] = 0x00
req[3] = 0x04 // ATYP IPv6
copy(req[4:20], ip16)
binary.BigEndian.PutUint16(req[20:22], target.Port)
case AddrTypeDomain:
dLen := len(target.Host)
if dLen == 0 || dLen > 255 {
return nil
}
req = make([]byte, 4+1+dLen+2)
req[0] = 0x05
req[1] = 0x01
req[2] = 0x00
req[3] = 0x03 // ATYP Domain
req[4] = byte(dLen)
copy(req[5:5+dLen], []byte(target.Host))
binary.BigEndian.PutUint16(req[5+dLen:7+dLen], target.Port)
default:
return nil
}
return req
}
// SocksUDPSession manages a SOCKS5 UDP ASSOCIATE tunnel to Xray.
type SocksUDPSession struct {
ctrl net.Conn
udpConn *net.UDPConn
targetEP net.Addr
user string
closed atomic.Bool
}
// DialUDP establishes a SOCKS5 UDP ASSOCIATE tunnel to Xray.
func (r *SocksRelay) DialUDP(ctx context.Context, user string) (*SocksUDPSession, error) {
dialer := net.Dialer{Timeout: 10 * time.Second}
ctrl, err := dialer.DialContext(ctx, "tcp", r.Addr)
if err != nil {
return nil, fmt.Errorf("tuic socks: dial UDP control connection: %w", err)
}
_ = ctrl.SetDeadline(time.Now().Add(10 * time.Second))
if err := socks5Handshake(ctrl, user, r.Password); err != nil {
ctrl.Close()
return nil, err
}
// SOCKS5 UDP ASSOCIATE (0x03), dst 0.0.0.0:0
if _, err := ctrl.Write([]byte{0x05, 0x03, 0x00, 0x01, 0, 0, 0, 0, 0, 0}); err != nil {
ctrl.Close()
return nil, fmt.Errorf("tuic socks: send UDP ASSOCIATE request: %w", err)
}
bind, err := readSocks5Reply(ctrl)
if err != nil {
ctrl.Close()
return nil, err
}
_ = ctrl.SetDeadline(time.Time{})
udpConn, err := net.DialUDP("udp", nil, net.UDPAddrFromAddrPort(bind))
if err != nil {
ctrl.Close()
return nil, fmt.Errorf("tuic socks: dial UDP relay endpoint %s: %w", bind, err)
}
return &SocksUDPSession{
ctrl: ctrl,
udpConn: udpConn,
targetEP: udpConn.RemoteAddr(),
user: user,
}, nil
}
// Send sends a UDP payload to target via the SOCKS5 UDP ASSOCIATE relay.
func (s *SocksUDPSession) Send(target *Address, payload []byte) (int, error) {
if s.closed.Load() {
return 0, net.ErrClosed
}
packet, err := buildSocks5UDPRequest(target, payload)
if err != nil {
return 0, err
}
return s.udpConn.Write(packet)
}
func buildSocks5UDPRequest(target *Address, payload []byte) ([]byte, error) {
hdr := buildSocks5UDPHeader(target)
if hdr == nil {
return nil, ErrInvalidAddr
}
if len(hdr)+len(payload) > maxSocksUdpDatagramSize {
return nil, ErrUdpPayloadTooLarge
}
packet := make([]byte, len(hdr)+len(payload))
copy(packet, hdr)
copy(packet[len(hdr):], payload)
return packet, nil
}
// Receive reads a relayed UDP payload and extracts its original source address.
func (s *SocksUDPSession) Receive(buf []byte) (*Address, []byte, error) {
if s.closed.Load() {
return nil, nil, net.ErrClosed
}
n, err := s.udpConn.Read(buf)
if err != nil {
return nil, nil, err
}
if n < 4 {
return nil, nil, fmt.Errorf("tuic socks: UDP packet too short (%d bytes)", n)
}
// SOCKS5 UDP header: [RSV(2)][FRAG(1)][ATYP(1)]
atyp := buf[3]
var addr *Address
var offset int
switch atyp {
case 0x01: // IPv4
if n < 10 {
return nil, nil, fmt.Errorf("tuic socks: truncated IPv4 UDP reply")
}
ip := net.IP(buf[4:8])
port := binary.BigEndian.Uint16(buf[8:10])
addr = &Address{Type: AddrTypeIPv4, IP: ip, Host: ip.String(), Port: port}
offset = 10
case 0x04: // IPv6
if n < 22 {
return nil, nil, fmt.Errorf("tuic socks: truncated IPv6 UDP reply")
}
ip := net.IP(buf[4:20])
port := binary.BigEndian.Uint16(buf[20:22])
addr = &Address{Type: AddrTypeIPv6, IP: ip, Host: ip.String(), Port: port}
offset = 22
case 0x03: // Domain
dLen := int(buf[4])
if n < 5+dLen+2 {
return nil, nil, fmt.Errorf("tuic socks: truncated domain UDP reply")
}
host := string(buf[5 : 5+dLen])
port := binary.BigEndian.Uint16(buf[5+dLen : 7+dLen])
addr = &Address{Type: AddrTypeDomain, Host: host, Port: port}
offset = 7 + dLen
default:
return nil, nil, fmt.Errorf("tuic socks: unsupported reply ATYP 0x%02x", atyp)
}
return addr, buf[offset:n], nil
}
// Close closes the SOCKS5 UDP session.
func (s *SocksUDPSession) Close() error {
if s.closed.Swap(true) {
return nil
}
_ = s.udpConn.Close()
return s.ctrl.Close()
}
func buildSocks5UDPHeader(target *Address) []byte {
if target == nil {
return nil
}
var hdr []byte
switch target.Type {
case AddrTypeIPv4:
ip4 := target.IP.To4()
if len(ip4) != 4 {
return nil
}
hdr = make([]byte, 10)
hdr[0] = 0x00 // RSV
hdr[1] = 0x00 // RSV
hdr[2] = 0x00 // FRAG
hdr[3] = 0x01 // ATYP IPv4
copy(hdr[4:8], ip4)
binary.BigEndian.PutUint16(hdr[8:10], target.Port)
case AddrTypeIPv6:
ip16 := target.IP.To16()
if len(ip16) != 16 {
return nil
}
hdr = make([]byte, 22)
hdr[0] = 0x00
hdr[1] = 0x00
hdr[2] = 0x00
hdr[3] = 0x04 // ATYP IPv6
copy(hdr[4:20], ip16)
binary.BigEndian.PutUint16(hdr[20:22], target.Port)
case AddrTypeDomain:
dLen := len(target.Host)
if dLen == 0 || dLen > 255 {
return nil
}
hdr = make([]byte, 4+1+dLen+2)
hdr[0] = 0x00
hdr[1] = 0x00
hdr[2] = 0x00
hdr[3] = 0x03 // ATYP Domain
hdr[4] = byte(dLen)
copy(hdr[5:5+dLen], []byte(target.Host))
binary.BigEndian.PutUint16(hdr[5+dLen:7+dLen], target.Port)
default:
return nil
}
return hdr
}
func socks5Handshake(conn net.Conn, user, password string) error {
if _, err := conn.Write([]byte{0x05, 0x02, 0x00, 0x02}); err != nil {
return fmt.Errorf("tuic socks: send greeting: %w", err)
}
var resp [2]byte
if _, err := io.ReadFull(conn, resp[:]); err != nil {
return fmt.Errorf("tuic socks: read greeting reply: %w", err)
}
if resp[0] != 0x05 {
return fmt.Errorf("tuic socks: unexpected SOCKS version %d", resp[0])
}
switch resp[1] {
case 0x00: // no auth
return nil
case 0x02: // username/password
req := make([]byte, 0, 3+len(user)+len(password))
req = append(req, 0x01, byte(len(user)))
req = append(req, user...)
req = append(req, byte(len(password)))
req = append(req, password...)
if _, err := conn.Write(req); err != nil {
return fmt.Errorf("tuic socks: send auth: %w", err)
}
var authResp [2]byte
if _, err := io.ReadFull(conn, authResp[:]); err != nil {
return fmt.Errorf("tuic socks: read auth reply: %w", err)
}
if authResp[1] != 0x00 {
return fmt.Errorf("tuic socks: auth rejected (code %d)", authResp[1])
}
return nil
default:
return fmt.Errorf("tuic socks: unsupported auth method %d", resp[1])
}
}
func readSocks5Reply(r io.Reader) (netip.AddrPort, error) {
var hdr [4]byte
if _, err := io.ReadFull(r, hdr[:]); err != nil {
return netip.AddrPort{}, fmt.Errorf("tuic socks: read reply header: %w", err)
}
if hdr[0] != 0x05 {
return netip.AddrPort{}, fmt.Errorf("tuic socks: unexpected SOCKS version %d", hdr[0])
}
if hdr[1] != 0x00 {
return netip.AddrPort{}, fmt.Errorf("tuic socks: request rejected (code %d)", hdr[1])
}
addr, err := readSocks5Addr(r, hdr[3])
if err != nil {
return netip.AddrPort{}, err
}
var portBytes [2]byte
if _, err := io.ReadFull(r, portBytes[:]); err != nil {
return netip.AddrPort{}, fmt.Errorf("tuic socks: read reply port: %w", err)
}
return netip.AddrPortFrom(addr, binary.BigEndian.Uint16(portBytes[:])), nil
}
func readSocks5Addr(r io.Reader, atyp byte) (netip.Addr, error) {
switch atyp {
case 0x01:
var b [4]byte
if _, err := io.ReadFull(r, b[:]); err != nil {
return netip.Addr{}, err
}
return netip.AddrFrom4(b), nil
case 0x04:
var b [16]byte
if _, err := io.ReadFull(r, b[:]); err != nil {
return netip.Addr{}, err
}
return netip.AddrFrom16(b), nil
case 0x03:
var l [1]byte
if _, err := io.ReadFull(r, l[:]); err != nil {
return netip.Addr{}, err
}
name := make([]byte, l[0])
if _, err := io.ReadFull(r, name); err != nil {
return netip.Addr{}, err
}
resolved, err := net.ResolveIPAddr("ip", string(name))
if err != nil {
return netip.Addr{}, fmt.Errorf("tuic socks: resolve domain reply %q: %w", name, err)
}
addr, ok := netip.AddrFromSlice(resolved.IP)
if !ok {
return netip.Addr{}, fmt.Errorf("tuic socks: unparseable domain reply address")
}
return addr, nil
default:
return netip.Addr{}, fmt.Errorf("tuic socks: unsupported SOCKS5 address type %d", atyp)
}
}
// halfCloseIdle bounds how long the surviving direction of a half-closed pair
// may sit idle, so a peer that vanished mid-transfer cannot pin it forever.
const halfCloseIdle = 2 * time.Minute
type closeWriter interface {
CloseWrite() error
}
type readDeadliner interface {
SetReadDeadline(t time.Time) error
}
type guardedReader struct {
r io.Reader
dl readDeadliner
armed atomic.Bool
}
func newGuardedReader(r io.Reader) *guardedReader {
gr := &guardedReader{r: r}
if dl, ok := r.(readDeadliner); ok {
gr.dl = dl
}
return gr
}
func (g *guardedReader) Read(p []byte) (int, error) {
if g.armed.Load() && g.dl != nil {
_ = g.dl.SetReadDeadline(time.Now().Add(halfCloseIdle))
}
return g.r.Read(p)
}
func (g *guardedReader) arm() {
g.armed.Store(true)
if g.dl != nil {
_ = g.dl.SetReadDeadline(time.Now().Add(halfCloseIdle))
}
}
// PipeBiDirectional pipes data between two connections and tracks byte counts in each direction.
func PipeBiDirectional(a, b io.ReadWriteCloser, upCounter, downCounter *atomic.Int64) {
PipeBiDirectionalContext(context.Background(), a, b, upCounter, downCounter)
}
func PipeBiDirectionalContext(ctx context.Context, a, b io.ReadWriteCloser, upCounter, downCounter *atomic.Int64) {
closeBoth := func() { _ = a.Close(); _ = b.Close() }
stop := context.AfterFunc(ctx, closeBoth)
defer stop()
ga := newGuardedReader(a)
gb := newGuardedReader(b)
var wg sync.WaitGroup
wg.Add(2)
pipe := func(dst io.Writer, dstGuard *guardedReader, src *guardedReader, counter *atomic.Int64) {
defer wg.Done()
buf := make([]byte, 32*1024)
for {
n, err := src.Read(buf)
if n > 0 {
if counter != nil {
counter.Add(int64(n))
}
if _, werr := dst.Write(buf[:n]); werr != nil {
closeBoth()
break
}
}
if err != nil {
if !errors.Is(err, io.EOF) {
closeBoth()
}
break
}
}
if cw, ok := dst.(closeWriter); ok {
_ = cw.CloseWrite()
} else if closer, ok := dst.(io.Closer); ok {
_ = closer.Close()
}
dstGuard.arm()
}
// a -> b (upload: client to upstream)
go pipe(b, gb, ga, upCounter)
// b -> a (download: upstream to client)
go pipe(a, ga, gb, downCounter)
wg.Wait()
_ = a.Close()
_ = b.Close()
}
// SOCKSBasePort is the first loopback port used for a TUIC inbound's
// internal Xray SOCKS5 relay inbound.
const SOCKSBasePort = 64000
// relayPortSlots is how many ids fit in the TUIC relay port window (64001..65000).
const relayPortSlots = 1000
// SOCKSPortForInbound derives one inbound's loopback SOCKS5 relay port from
// its id, bounded within the dedicated range 64001..65000 so it never collides
// with AmneziaWG (65101..65535), API (62789), or node egress (62800..63800).
func SOCKSPortForInbound(inboundID int) int {
if inboundID <= 0 {
return SOCKSBasePort + 1
}
return SOCKSBasePort + 1 + (inboundID-1)%relayPortSlots
}
var (
socksPasswordOnce sync.Once
socksPassword string
)
// SocksPassword returns the process-wide password used to authenticate into
// every TUIC SOCKS5 relay inbound.
func SocksPassword() string {
socksPasswordOnce.Do(func() {
var b [24]byte
if _, err := rand.Read(b[:]); err != nil {
socksPassword = fmt.Sprintf("tuic-fallback-%x", b)
return
}
socksPassword = base64.RawURLEncoding.EncodeToString(b[:])
})
return socksPassword
}
// SocksInboundSettings builds the JSON `settings` block for a stock Xray
// SOCKS5 inbound with one username/password account per email, all sharing
// password. UDP is enabled for UDP ASSOCIATE proxying.
func SocksInboundSettings(emails []string, password string) ([]byte, error) {
type account struct {
User string `json:"user"`
Pass string `json:"pass"`
}
settings := struct {
Auth string `json:"auth"`
UDP bool `json:"udp"`
Accounts []account `json:"accounts"`
}{Auth: "password", UDP: true}
for _, email := range emails {
settings.Accounts = append(settings.Accounts, account{User: email, Pass: password})
}
return json.Marshal(settings)
}