add go-tun2socks code

This commit is contained in:
Jason
2019-07-16 11:37:52 +08:00
parent 828ba9948d
commit 6d01dec5a4
301 changed files with 69694 additions and 1 deletions
+96
View File
@@ -0,0 +1,96 @@
package d
import (
"io"
"net"
"strconv"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/common/lsof"
"github.com/xjasonlyu/tun2socks/core"
)
// This handler allows you chain another proxy behind tun2socks locally, typically a rule-based proxy client, e.g. V2Ray.
//
// Rule-based proxy clients are very useful, they are able to dispatch requests to different servers based on powerful rule filters.
// By using this setup, you are able to make all your TCP/UDP traffic under control with your favorite rule-based proxy client.
//
// Here's an example setup on macOS:
//
// tun2socks -tunGw 10.255.0.1 -fakeDns -proxyType d -proxyServer 127.0.0.1:1086 -exceptionSendThrough 192.168.1.189:0 -exceptionApps "v2ray"
//
// route delete default
// route add default 10.255.0.1
// route add default 192.168.1.1 -ifscope en0
//
// Where 192.168.1.189 is the default interface address, in my case, it's the WiFi interface and it's en0.
// 192.168.1.1 is the default gateway.
// It's very important to have two default routes, and the default route to TUN should has the highest priority.
//
// Start v2ray (or any other chainable proxy clients) and has SOCKS inbound listen on 127.0.0.1:1086.
//
// Optinally with all outbounds have sendThrough set to 192.168.1.189, if applicable.
// https://v2ray.com/chapter_02/01_overview.html#outboundobject
type tcpHandler struct {
proxyHandler core.TCPConnHandler
exceptionApps []string
sendThrough net.Addr
}
func NewTCPHandler(proxyHandler core.TCPConnHandler, exceptionApps []string, sendThrough net.Addr) core.TCPConnHandler {
return &tcpHandler{
proxyHandler,
exceptionApps,
sendThrough,
}
}
func (h *tcpHandler) isExceptionApp(name string) bool {
for _, app := range h.exceptionApps {
if name == app {
return true
}
}
return false
}
func (h *tcpHandler) relay(lhs, rhs net.Conn) {
cls := func() {
rhs.Close()
lhs.Close()
}
go func() {
io.Copy(rhs, lhs)
cls()
}()
io.Copy(lhs, rhs)
cls()
}
func (h *tcpHandler) Handle(conn net.Conn, target *net.TCPAddr) error {
localHost, localPortStr, _ := net.SplitHostPort(conn.LocalAddr().String())
localPortInt, _ := strconv.Atoi(localPortStr)
cmd, err := lsof.GetCommandNameBySocket("tcp", localHost, uint16(localPortInt))
if err != nil {
cmd = "unknown process"
}
if h.isExceptionApp(cmd) {
dialer := net.Dialer{LocalAddr: h.sendThrough}
rc, err := dialer.Dial("tcp", target.String())
if err != nil {
return err
}
go h.relay(conn, rc)
log.Access(cmd, "direct", target.Network(), conn.LocalAddr().String(), target.String())
return nil
} else {
return h.proxyHandler.Handle(conn, target)
}
}
+121
View File
@@ -0,0 +1,121 @@
package d
import (
"net"
"strconv"
"sync"
"time"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/common/lsof"
"github.com/xjasonlyu/tun2socks/core"
)
type udpHandler struct {
sync.Mutex
proxyHandler core.UDPConnHandler
exceptionApps []string
sendThrough net.Addr
exceptionConns map[core.UDPConn]*net.UDPConn
timeout time.Duration
}
func (h *udpHandler) isExceptionApp(name string) bool {
for _, app := range h.exceptionApps {
if name == app {
return true
}
}
return false
}
func NewUDPHandler(proxyHandler core.UDPConnHandler, exceptionApps []string, sendThrough net.Addr, timeout time.Duration) core.UDPConnHandler {
return &udpHandler{
proxyHandler: proxyHandler,
exceptionApps: exceptionApps,
sendThrough: sendThrough,
exceptionConns: make(map[core.UDPConn]*net.UDPConn),
timeout: timeout,
}
}
func (h *udpHandler) handleInput(conn core.UDPConn, pc *net.UDPConn) {
buf := core.NewBytes(core.BufSize)
defer func() {
h.Close(conn)
core.FreeBytes(buf)
}()
for {
pc.SetDeadline(time.Now().Add(h.timeout))
n, addr, err := pc.ReadFromUDP(buf)
if err != nil {
return
}
_, err = conn.WriteFrom(buf[:n], addr)
if err != nil {
return
}
}
}
func (h *udpHandler) Connect(conn core.UDPConn, target *net.UDPAddr) error {
localHost, localPortStr, _ := net.SplitHostPort(conn.LocalAddr().String())
localPortInt, _ := strconv.Atoi(localPortStr)
cmd, err := lsof.GetCommandNameBySocket("udp", localHost, uint16(localPortInt))
if err != nil {
cmd = "unknown process"
}
if h.isExceptionApp(cmd) {
bindAddr, _ := net.ResolveUDPAddr(
"udp",
h.sendThrough.String(),
)
pc, err := net.ListenUDP("udp", bindAddr)
if err != nil {
return err
}
h.Lock()
h.exceptionConns[conn] = pc
h.Unlock()
go h.handleInput(conn, pc)
log.Access(cmd, "direct", target.Network(), conn.LocalAddr().String(), target.String())
return nil
} else {
return h.proxyHandler.Connect(conn, target)
}
}
func (h *udpHandler) ReceiveTo(conn core.UDPConn, data []byte, addr *net.UDPAddr) error {
h.Lock()
defer h.Unlock()
if pc, found := h.exceptionConns[conn]; found {
_, err := pc.WriteTo(data, addr)
if err != nil {
return err
}
return nil
} else {
return h.proxyHandler.ReceiveTo(conn, data, addr)
}
}
func (h *udpHandler) Close(conn core.UDPConn) {
conn.Close()
h.Lock()
defer h.Unlock()
if pc, ok := h.exceptionConns[conn]; ok {
pc.Close()
delete(h.exceptionConns, conn)
}
}
+46
View File
@@ -0,0 +1,46 @@
package direct
import (
"errors"
"fmt"
"io"
"net"
"sync"
"time"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/core"
)
type tcpHandler struct{}
func NewTCPHandler() core.TCPConnHandler {
return &tcpHandler{}
}
func (h *tcpHandler) handleInput(conn net.Conn, input io.ReadCloser) {
defer func() {
conn.Close()
input.Close()
}()
io.Copy(conn, input)
}
func (h *tcpHandler) handleOutput(conn net.Conn, output io.WriteCloser) {
defer func() {
conn.Close()
output.Close()
}()
io.Copy(output, conn)
}
func (h *tcpHandler) Handle(conn net.Conn, target *net.TCPAddr) error {
c, err := net.DialTCP("tcp", nil, target)
if err != nil {
return err
}
go h.handleInput(conn, c)
go h.handleOutput(conn, c)
log.Infof("new proxy connection for target: %s:%s", target.Network(), target.String())
return nil
}
+94
View File
@@ -0,0 +1,94 @@
package direct
import (
"errors"
"fmt"
"net"
"sync"
"time"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/core"
)
type udpHandler struct {
sync.Mutex
timeout time.Duration
udpConns map[core.UDPConn]*net.UDPConn
}
func NewUDPHandler(timeout time.Duration) core.UDPConnHandler {
return &udpHandler{
timeout: timeout,
udpConns: make(map[core.UDPConn]*net.UDPConn, 8),
}
}
func (h *udpHandler) fetchUDPInput(conn core.UDPConn, pc *net.UDPConn) {
buf := core.NewBytes(core.BufSize)
defer func() {
h.Close(conn)
core.FreeBytes(buf)
}()
for {
pc.SetDeadline(time.Now().Add(h.timeout))
n, addr, err := pc.ReadFromUDP(buf)
if err != nil {
// log.Printf("failed to read UDP data from remote: %v", err)
return
}
_, err = conn.WriteFrom(buf[:n], addr)
if err != nil {
log.Warnf("failed to write UDP data to TUN")
return
}
}
}
func (h *udpHandler) Connect(conn core.UDPConn, target net.Addr) error {
bindAddr := &net.UDPAddr{IP: nil, Port: 0}
pc, err := net.ListenUDP("udp", bindAddr)
if err != nil {
log.Errorf("failed to bind udp address")
return err
}
h.Lock()
h.udpConns[conn] = pc
h.Unlock()
go h.fetchUDPInput(conn, pc)
log.Infof("new proxy connection for target: %s:%s", target.Network(), target.String())
return nil
}
func (h *udpHandler) ReceiveTo(conn core.UDPConn, data []byte, addr net.Addr) error {
h.Lock()
pc, ok1 := h.udpConns[conn]
h.Unlock()
if ok1 {
_, err := pc.WriteToUDP(data, addr)
if err != nil {
log.Warnf("failed to write UDP payload to SOCKS5 server: %v", err)
return errors.New("failed to write UDP data")
}
return nil
} else {
return errors.New(fmt.Sprintf("proxy connection %v->%v does not exists", conn.LocalAddr(), conn.RemoteAddr()))
}
}
func (h *udpHandler) Close(conn core.UDPConn) {
conn.Close()
h.Lock()
defer h.Unlock()
if pc, ok := h.udpConns[conn]; ok {
pc.Close()
delete(h.udpConns, conn)
}
}
+66
View File
@@ -0,0 +1,66 @@
package dnsfallback
import (
"encoding/binary"
"errors"
"net"
"github.com/xjasonlyu/tun2socks/common/dns"
"github.com/xjasonlyu/tun2socks/core"
)
// UDP handler that intercepts DNS queries and replies with a truncated response (TC bit)
// in order for the client to retry over TCP. This DNS/TCP fallback mechanism is
// useful for proxy servers that do not support UDP.
// Note that non-DNS UDP traffic is dropped.
type udpHandler struct{}
const (
dnsHeaderLength = 12
dnsMaskQr = uint8(0x80)
dnsMaskTc = uint8(0x02)
dnsMaskRcode = uint8(0x0F)
)
func NewUDPHandler() core.UDPConnHandler {
return &udpHandler{}
}
func (h *udpHandler) Connect(conn core.UDPConn, udpAddr *net.UDPAddr) error {
if udpAddr.Port != dns.CommonDnsPort {
return errors.New("Cannot handle non-DNS packet")
}
return nil
}
func (h *udpHandler) ReceiveTo(conn core.UDPConn, data []byte, addr *net.UDPAddr) error {
if len(data) < dnsHeaderLength {
return errors.New("Received malformed DNS query")
}
// DNS Header
// 0 1 2 3 4 5 6 7 0 1 2 3 4 5 6 7
// +--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+
// | ID |
// +--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+
// |QR| Opcode |AA|TC|RD|RA| Z | RCODE |
// +--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+
// | QDCOUNT |
// +--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+
// | ANCOUNT |
// +--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+
// | NSCOUNT |
// +--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+
// | ARCOUNT |
// +--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+--+
// Set response and truncated bits
data[2] |= dnsMaskQr | dnsMaskTc
// Set response code to 'no error'.
data[3] &= ^dnsMaskRcode
// Set ANCOUNT to QDCOUNT. This is technically incorrect, since the response does not
// include an answer. However, without it some DNS clients (i.e. Windows 7) do not retry
// over TCP.
var qdcount = binary.BigEndian.Uint16(data[4:6])
binary.BigEndian.PutUint16(data[6:], qdcount)
_, err := conn.WriteFrom(data, addr)
return err
}
+26
View File
@@ -0,0 +1,26 @@
package echo
import (
"io"
"net"
"github.com/xjasonlyu/tun2socks/core"
)
// An echo proxy, do nothing but echo back data to the sender, the handler was
// created for testing purposes, it may causes issues when more than one clients
// are connecting the handler simultaneously.
type tcpHandler struct{}
func NewTCPHandler() core.TCPConnHandler {
return &tcpHandler{}
}
func (h *tcpHandler) echoBack(conn net.Conn) {
io.Copy(conn, conn)
}
func (h *tcpHandler) Handle(conn net.Conn, target *net.TCPAddr) error {
go h.echoBack(conn)
return nil
}
+31
View File
@@ -0,0 +1,31 @@
package echo
import (
"net"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/core"
)
// An echo server, do nothing but echo back data to the sender.
type udpHandler struct{}
func NewUDPHandler() core.UDPConnHandler {
return &udpHandler{}
}
func (h *udpHandler) Connect(conn core.UDPConn, target *net.UDPAddr) error {
return nil
}
func (h *udpHandler) ReceiveTo(conn core.UDPConn, data []byte, addr *net.UDPAddr) error {
// Dispatch to another goroutine, otherwise will result in deadlock.
payload := append([]byte(nil), data...)
go func(b []byte) {
_, err := conn.WriteFrom(b, addr)
if err != nil {
log.Warnf("failed to echo back data: %v", err)
}
}(payload)
return nil
}
+87
View File
@@ -0,0 +1,87 @@
package redirect
import (
"io"
"net"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/core"
)
// To do a benchmark using iperf3 locally, you may follow these steps:
//
// 1. Setup and configure the TUN device and start tun2socks with the
// redirect handler using the following command:
// tun2socks -proxyType redirect -proxyServer 127.0.0.1:1234
// Tun2socks will redirect all traffic to 127.0.0.1:1234.
//
// 2. Route traffic targeting 1.2.3.4 to the TUN interface (240.0.0.1):
// route add 1.2.3.4/32 240.0.0.1
//
// 3. Run iperf3 server locally and listening on 1234 port:
// iperf3 -s -p 1234
//
// 4. Run iperf3 client locally and connect to 1.2.3.4:1234:
// iperf3 -c 1.2.3.4 -p 1234
//
// It works this way:
// iperf3 client -> 1.2.3.4:1234 -> routing table -> TUN (240.0.0.1) -> tun2socks -> tun2socks redirect anything to 127.0.0.1:1234 -> iperf3 server
//
type tcpHandler struct {
target string
}
type duplexConn interface {
net.Conn
CloseWrite() error
CloseRead() error
}
func NewTCPHandler(target string) core.TCPConnHandler {
return &tcpHandler{target: target}
}
func (h *tcpHandler) handleInput(conn net.Conn, input io.ReadCloser) {
defer func() {
if tcpConn, ok := conn.(core.TCPConn); ok {
tcpConn.CloseWrite()
} else {
conn.Close()
}
if tcpInput, ok := input.(duplexConn); ok {
tcpInput.CloseRead()
} else {
input.Close()
}
}()
io.Copy(conn, input)
}
func (h *tcpHandler) handleOutput(conn net.Conn, output io.WriteCloser) {
defer func() {
if tcpConn, ok := conn.(core.TCPConn); ok {
tcpConn.CloseRead()
} else {
conn.Close()
}
if tcpOutput, ok := output.(duplexConn); ok {
tcpOutput.CloseWrite()
} else {
output.Close()
}
}()
io.Copy(output, conn)
}
func (h *tcpHandler) Handle(conn net.Conn, target *net.TCPAddr) error {
c, err := net.Dial("tcp", h.target)
if err != nil {
return err
}
go h.handleInput(conn, c)
go h.handleOutput(conn, c)
log.Infof("new proxy connection for target: %s:%s", target.Network(), target.String())
return nil
}
+104
View File
@@ -0,0 +1,104 @@
package redirect
import (
"errors"
"fmt"
"net"
"sync"
"time"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/core"
)
type udpHandler struct {
sync.Mutex
timeout time.Duration
udpConns map[core.UDPConn]*net.UDPConn
udpTargetAddrs map[core.UDPConn]*net.UDPAddr
target string
}
func NewUDPHandler(target string, timeout time.Duration) core.UDPConnHandler {
return &udpHandler{
timeout: timeout,
udpConns: make(map[core.UDPConn]*net.UDPConn, 8),
udpTargetAddrs: make(map[core.UDPConn]*net.UDPAddr, 8),
target: target,
}
}
func (h *udpHandler) fetchUDPInput(conn core.UDPConn, pc *net.UDPConn) {
buf := core.NewBytes(core.BufSize)
defer func() {
h.Close(conn)
core.FreeBytes(buf)
}()
for {
pc.SetDeadline(time.Now().Add(h.timeout))
n, addr, err := pc.ReadFromUDP(buf)
if err != nil {
// log.Printf("failed to read UDP data from remote: %v", err)
return
}
_, err = conn.WriteFrom(buf[:n], addr)
if err != nil {
log.Warnf("failed to write UDP data to TUN")
return
}
}
}
func (h *udpHandler) Connect(conn core.UDPConn, target *net.UDPAddr) error {
bindAddr := &net.UDPAddr{IP: nil, Port: 0}
pc, err := net.ListenUDP("udp", bindAddr)
if err != nil {
log.Errorf("failed to bind udp address")
return err
}
tgtAddr, _ := net.ResolveUDPAddr("udp", h.target)
h.Lock()
h.udpTargetAddrs[conn] = tgtAddr
h.udpConns[conn] = pc
h.Unlock()
go h.fetchUDPInput(conn, pc)
log.Infof("new proxy connection for target: %s:%s", target.Network(), target.String())
return nil
}
func (h *udpHandler) ReceiveTo(conn core.UDPConn, data []byte, addr *net.UDPAddr) error {
h.Lock()
pc, ok1 := h.udpConns[conn]
tgtAddr, ok2 := h.udpTargetAddrs[conn]
h.Unlock()
if ok1 && ok2 {
_, err := pc.WriteToUDP(data, tgtAddr)
if err != nil {
log.Warnf("failed to write UDP payload to SOCKS5 server: %v", err)
return errors.New("failed to write UDP data")
}
return nil
} else {
return errors.New(fmt.Sprintf("proxy connection %v->%v does not exists", conn.LocalAddr(), addr))
}
}
func (h *udpHandler) Close(conn core.UDPConn) {
conn.Close()
h.Lock()
defer h.Unlock()
if _, ok := h.udpTargetAddrs[conn]; ok {
delete(h.udpTargetAddrs, conn)
}
if pc, ok := h.udpConns[conn]; ok {
pc.Close()
delete(h.udpConns, conn)
}
}
+85
View File
@@ -0,0 +1,85 @@
package shadowsocks
import (
"errors"
"fmt"
"io"
"net"
"strconv"
sscore "github.com/shadowsocks/go-shadowsocks2/core"
sssocks "github.com/shadowsocks/go-shadowsocks2/socks"
"github.com/xjasonlyu/tun2socks/common/dns"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/core"
)
type tcpHandler struct {
cipher sscore.Cipher
server string
fakeDns dns.FakeDns
}
func (h *tcpHandler) handleInput(conn net.Conn, input io.ReadCloser) {
defer func() {
conn.Close()
input.Close()
}()
io.Copy(conn, input)
}
func (h *tcpHandler) handleOutput(conn net.Conn, output io.WriteCloser) {
defer func() {
conn.Close()
output.Close()
}()
io.Copy(output, conn)
}
func NewTCPHandler(server, cipher, password string, fakeDns dns.FakeDns) core.TCPConnHandler {
ciph, err := sscore.PickCipher(cipher, []byte{}, password)
if err != nil {
log.Errorf("failed to pick a cipher: %v", err)
}
return &tcpHandler{
cipher: ciph,
server: server,
fakeDns: fakeDns,
}
}
func (h *tcpHandler) Handle(conn net.Conn, target *net.TCPAddr) error {
if target == nil {
log.Fatalf("unexpected nil target")
}
// Connect the relay server.
rc, err := net.Dial("tcp", h.server)
if err != nil {
return errors.New(fmt.Sprintf("dial remote server failed: %v", err))
}
rc = h.cipher.StreamConn(rc)
// Replace with a domain name if target address IP is a fake IP.
var targetHost string
if h.fakeDns != nil && h.fakeDns.IsFakeIP(target.IP) {
targetHost = h.fakeDns.QueryDomain(target.IP)
} else {
targetHost = target.IP.String()
}
dest := net.JoinHostPort(targetHost, strconv.Itoa(target.Port))
// Write target address.
tgt := sssocks.ParseAddr(dest)
_, err = rc.Write(tgt)
if err != nil {
return fmt.Errorf("send target address failed: %v", err)
}
go h.handleInput(conn, rc)
go h.handleOutput(conn, rc)
log.Infof("new proxy connection for target: %s:%s", target.Network(), dest)
return nil
}
+172
View File
@@ -0,0 +1,172 @@
package shadowsocks
import (
"errors"
"fmt"
"net"
"strconv"
"sync"
"time"
sscore "github.com/shadowsocks/go-shadowsocks2/core"
sssocks "github.com/shadowsocks/go-shadowsocks2/socks"
"github.com/xjasonlyu/tun2socks/common/dns"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/core"
)
type udpHandler struct {
sync.Mutex
cipher sscore.Cipher
remoteAddr net.Addr
conns map[core.UDPConn]net.PacketConn
dnsCache dns.DnsCache
fakeDns dns.FakeDns
timeout time.Duration
}
func NewUDPHandler(server, cipher, password string, timeout time.Duration, dnsCache dns.DnsCache, fakeDns dns.FakeDns) core.UDPConnHandler {
ciph, err := sscore.PickCipher(cipher, []byte{}, password)
if err != nil {
log.Errorf("failed to pick a cipher: %v", err)
}
remoteAddr, err := net.ResolveUDPAddr("udp", server)
if err != nil {
log.Errorf("failed to resolve udp address: %v", err)
}
return &udpHandler{
cipher: ciph,
remoteAddr: remoteAddr,
conns: make(map[core.UDPConn]net.PacketConn, 16),
dnsCache: dnsCache,
fakeDns: fakeDns,
timeout: timeout,
}
}
func (h *udpHandler) fetchUDPInput(conn core.UDPConn, input net.PacketConn) {
buf := core.NewBytes(core.BufSize)
defer func() {
h.Close(conn)
core.FreeBytes(buf)
}()
for {
input.SetDeadline(time.Now().Add(h.timeout))
n, _, err := input.ReadFrom(buf)
if err != nil {
// log.Printf("read remote failed: %v", err)
return
}
addr := sssocks.SplitAddr(buf[:])
resolvedAddr, err := net.ResolveUDPAddr("udp", addr.String())
if err != nil {
return
}
_, err = conn.WriteFrom(buf[int(len(addr)):n], resolvedAddr)
if err != nil {
log.Warnf("write local failed: %v", err)
return
}
if h.dnsCache != nil {
_, port, err := net.SplitHostPort(addr.String())
if err != nil {
panic("impossible error")
}
if port == strconv.Itoa(dns.CommonDnsPort) {
h.dnsCache.Store(buf[int(len(addr)):n])
return // DNS response
}
}
}
}
func (h *udpHandler) Connect(conn core.UDPConn, target *net.UDPAddr) error {
pc, err := net.ListenPacket("udp", "")
if err != nil {
return err
}
pc = h.cipher.PacketConn(pc)
h.Lock()
h.conns[conn] = pc
h.Unlock()
go h.fetchUDPInput(conn, pc)
if target != nil {
log.Infof("new proxy connection for target: %s:%s", target.Network(), target.String())
}
return nil
}
func (h *udpHandler) ReceiveTo(conn core.UDPConn, data []byte, addr *net.UDPAddr) error {
h.Lock()
pc, ok1 := h.conns[conn]
h.Unlock()
if addr.Port == dns.CommonDnsPort {
if h.fakeDns != nil {
resp, err := h.fakeDns.GenerateFakeResponse(data)
if err == nil {
_, err = conn.WriteFrom(resp, addr)
if err != nil {
return errors.New(fmt.Sprintf("write dns answer failed: %v", err))
}
h.Close(conn)
return nil
}
}
if h.dnsCache != nil {
if answer := h.dnsCache.Query(data); answer != nil {
_, err := conn.WriteFrom(answer, addr)
if err != nil {
return errors.New(fmt.Sprintf("cache dns answer failed: %v", err))
}
h.Close(conn)
return nil
}
}
}
if ok1 {
// Replace with a domain name if target address IP is a fake IP.
var targetHost string
if h.fakeDns != nil && h.fakeDns.IsFakeIP(addr.IP) {
targetHost = h.fakeDns.QueryDomain(addr.IP)
} else {
targetHost = addr.IP.String()
}
dest := net.JoinHostPort(targetHost, strconv.Itoa(addr.Port))
buf := append([]byte{0, 0, 0}, sssocks.ParseAddr(dest)...)
buf = append(buf, data[:]...)
_, err := pc.WriteTo(buf[3:], h.remoteAddr)
if err != nil {
h.Close(conn)
return errors.New(fmt.Sprintf("write remote failed: %v", err))
}
return nil
} else {
h.Close(conn)
return errors.New(fmt.Sprintf("proxy connection %v->%v does not exists", conn.LocalAddr(), addr))
}
}
func (h *udpHandler) Close(conn core.UDPConn) {
conn.Close()
h.Lock()
defer h.Unlock()
if pc, ok := h.conns[conn]; ok {
pc.Close()
delete(h.conns, conn)
}
}
+162
View File
@@ -0,0 +1,162 @@
// Code in this file are grabbed from https://github.com/nadoo/glider, which
// is also referencing another repo: https://github.com/shadowsocks/go-shadowsocks2
package socks
import (
"errors"
"io"
"net"
"strconv"
)
// SOCKS request commands as defined in RFC 1928 section 4.
const (
socks5Connect = 1
socks5Bind = 2
socks5UDPAssociate = 3
)
// SOCKS address types as defined in RFC 1928 section 5.
const (
socks5IP4 = 1
socks5Domain = 3
socks5IP6 = 4
)
var socks5Errors = []error{
errors.New(""),
errors.New("general failure"),
errors.New("connection forbidden"),
errors.New("network unreachable"),
errors.New("host unreachable"),
errors.New("connection refused"),
errors.New("TTL expired"),
errors.New("command not supported"),
errors.New("address type not supported"),
errors.New("socks5UDPAssociate"),
}
// MaxAddrLen is the maximum size of SOCKS address in bytes.
const MaxAddrLen = 1 + 1 + 255 + 2
// ATYP return the address type
func ATYP(b byte) int {
return int(b &^ 0x8)
}
// Addr represents a SOCKS address as defined in RFC 1928 section 5.
type Addr []byte
// String serializes SOCKS address a to string form.
func (a Addr) String() string {
var host, port string
switch ATYP(a[0]) { // address type
case socks5Domain:
host = string(a[2 : 2+int(a[1])])
port = strconv.Itoa((int(a[2+int(a[1])]) << 8) | int(a[2+int(a[1])+1]))
case socks5IP4:
host = net.IP(a[1 : 1+net.IPv4len]).String()
port = strconv.Itoa((int(a[1+net.IPv4len]) << 8) | int(a[1+net.IPv4len+1]))
case socks5IP6:
host = net.IP(a[1 : 1+net.IPv6len]).String()
port = strconv.Itoa((int(a[1+net.IPv6len]) << 8) | int(a[1+net.IPv6len+1]))
}
return net.JoinHostPort(host, port)
}
// ParseAddr parses the address in string s. Returns nil if failed.
func ParseAddr(s string) Addr {
var addr Addr
host, port, err := net.SplitHostPort(s)
if err != nil {
return nil
}
if ip := net.ParseIP(host); ip != nil {
if ip4 := ip.To4(); ip4 != nil {
addr = make([]byte, 1+net.IPv4len+2)
addr[0] = socks5IP4
copy(addr[1:], ip4)
} else {
addr = make([]byte, 1+net.IPv6len+2)
addr[0] = socks5IP6
copy(addr[1:], ip)
}
} else {
if len(host) > 255 {
return nil
}
addr = make([]byte, 1+1+len(host)+2)
addr[0] = socks5Domain
addr[1] = byte(len(host))
copy(addr[2:], host)
}
portnum, err := strconv.ParseUint(port, 10, 16)
if err != nil {
return nil
}
addr[len(addr)-2], addr[len(addr)-1] = byte(portnum>>8), byte(portnum)
return addr
}
func readAddr(r io.Reader, b []byte) (Addr, error) {
if len(b) < MaxAddrLen {
return nil, io.ErrShortBuffer
}
_, err := io.ReadFull(r, b[:1]) // read 1st byte for address type
if err != nil {
return nil, err
}
switch ATYP(b[0]) {
case socks5Domain:
_, err = io.ReadFull(r, b[1:2]) // read 2nd byte for domain length
if err != nil {
return nil, err
}
_, err = io.ReadFull(r, b[2:2+int(b[1])+2])
return b[:1+1+int(b[1])+2], err
case socks5IP4:
_, err = io.ReadFull(r, b[1:1+net.IPv4len+2])
return b[:1+net.IPv4len+2], err
case socks5IP6:
_, err = io.ReadFull(r, b[1:1+net.IPv6len+2])
return b[:1+net.IPv6len+2], err
}
return nil, socks5Errors[8]
}
// SplitAddr slices a SOCKS address from beginning of b. Returns nil if failed.
func SplitAddr(b []byte) Addr {
addrLen := 1
if len(b) < addrLen {
return nil
}
switch ATYP(b[0]) {
case socks5Domain:
if len(b) < 2 {
return nil
}
addrLen = 1 + 1 + int(b[1]) + 2
case socks5IP4:
addrLen = 1 + net.IPv4len + 2
case socks5IP6:
addrLen = 1 + net.IPv6len + 2
default:
return nil
}
if len(b) < addrLen {
return nil
}
return b[:addrLen]
}
+192
View File
@@ -0,0 +1,192 @@
package socks
import (
"io"
"net"
"strconv"
"sync"
"time"
"golang.org/x/net/proxy"
"github.com/xjasonlyu/tun2socks/common/dns"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/common/lsof"
"github.com/xjasonlyu/tun2socks/common/stats"
"github.com/xjasonlyu/tun2socks/core"
)
type tcpHandler struct {
sync.Mutex
proxyHost string
proxyPort uint16
fakeDns dns.FakeDns
sessionStater stats.SessionStater
}
func NewTCPHandler(proxyHost string, proxyPort uint16, fakeDns dns.FakeDns, sessionStater stats.SessionStater) core.TCPConnHandler {
return &tcpHandler{
proxyHost: proxyHost,
proxyPort: proxyPort,
fakeDns: fakeDns,
sessionStater: sessionStater,
}
}
type direction byte
const (
dirUplink direction = iota
dirDownlink
)
func statsCopy(dst io.Writer, src io.Reader, sess *stats.Session, dir direction) (written int64, err error) {
buf := make([]byte, 32*1024)
for {
nr, er := src.Read(buf)
if nr > 0 {
nw, ew := dst.Write(buf[0:nr])
if nw > 0 {
switch dir {
case dirUplink:
sess.AddUploadBytes(int64(nw))
case dirDownlink:
sess.AddDownloadBytes(int64(nw))
default:
}
written += int64(nw)
}
if ew != nil {
err = ew
break
}
if nr != nw {
err = io.ErrShortWrite
break
}
}
if er != nil {
if er != io.EOF {
err = er
}
break
}
}
return written, err
}
type duplexConn interface {
net.Conn
CloseRead() error
CloseWrite() error
}
func (h *tcpHandler) relay(lhs, rhs net.Conn, sess *stats.Session) {
var err error
upCh := make(chan struct{})
cls := func(dir direction, interrupt bool) {
lhsDConn, lhsOk := lhs.(duplexConn)
rhsDConn, rhsOk := rhs.(duplexConn)
if !interrupt && lhsOk && rhsOk {
switch dir {
case dirUplink:
_ = lhsDConn.CloseRead()
_ = rhsDConn.CloseWrite()
case dirDownlink:
_ = lhsDConn.CloseWrite()
_ = rhsDConn.CloseRead()
default:
panic("unexpected direction")
}
} else {
_ = lhs.Close()
_ = rhs.Close()
}
}
// Uplink
go func() {
if h.sessionStater != nil && sess != nil {
_, err = statsCopy(rhs, lhs, sess, dirUplink)
} else {
_, err = io.Copy(rhs, lhs)
}
if err != nil {
cls(dirUplink, true) // interrupt the conn if the error is not nil (not EOF)
} else {
cls(dirUplink, false) // half close uplink direction of the TCP conn if possible
}
upCh <- struct{}{}
}()
// Downlink
if h.sessionStater != nil && sess != nil {
_, err = statsCopy(lhs, rhs, sess, dirDownlink)
} else {
_, err = io.Copy(lhs, rhs)
}
if err != nil {
cls(dirDownlink, true)
} else {
cls(dirDownlink, false)
}
<-upCh // Wait for uplink done.
if h.sessionStater != nil {
h.sessionStater.RemoveSession(lhs)
}
}
func (h *tcpHandler) Handle(conn net.Conn, target *net.TCPAddr) error {
dialer, err := proxy.SOCKS5("tcp", core.ParseTCPAddr(h.proxyHost, h.proxyPort).String(), nil, nil)
if err != nil {
return err
}
// Replace with a domain name if target address IP is a fake IP.
var targetHost string
if h.fakeDns != nil && h.fakeDns.IsFakeIP(target.IP) {
targetHost = h.fakeDns.IPToHost(target.IP)
} else {
targetHost = target.IP.String()
}
dest := net.JoinHostPort(targetHost, strconv.Itoa(target.Port))
c, err := dialer.Dial(target.Network(), dest)
if err != nil {
return err
}
var process string
var sess *stats.Session
if h.sessionStater != nil {
// Get name of the process.
localHost, localPortStr, _ := net.SplitHostPort(conn.LocalAddr().String())
localPortInt, _ := strconv.Atoi(localPortStr)
process, err = lsof.GetCommandNameBySocket(target.Network(), localHost, uint16(localPortInt))
if err != nil {
process = "unknown process"
}
sess = &stats.Session{
ProcessName: process,
Network: target.Network(),
LocalAddr: conn.LocalAddr().String(),
RemoteAddr: dest,
UploadBytes: 0,
DownloadBytes: 0,
SessionStart: time.Now(),
}
h.sessionStater.AddSession(conn, sess)
}
go h.relay(conn, c, sess)
log.Access(process, "proxy", target.Network(), conn.LocalAddr().String(), dest)
return nil
}
+271
View File
@@ -0,0 +1,271 @@
package socks
import (
"errors"
"fmt"
"io"
"net"
"strconv"
"sync"
"time"
"github.com/xjasonlyu/tun2socks/common/dns"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/common/lsof"
"github.com/xjasonlyu/tun2socks/common/stats"
"github.com/xjasonlyu/tun2socks/core"
)
type udpHandler struct {
sync.Mutex
proxyHost string
proxyPort uint16
udpConns map[core.UDPConn]net.PacketConn
tcpConns map[core.UDPConn]net.Conn
remoteAddrs map[core.UDPConn]*net.UDPAddr // UDP relay server addresses
timeout time.Duration
dnsCache dns.DnsCache
fakeDns dns.FakeDns
sessionStater stats.SessionStater
}
func NewUDPHandler(proxyHost string, proxyPort uint16, timeout time.Duration, dnsCache dns.DnsCache, fakeDns dns.FakeDns, sessionStater stats.SessionStater) core.UDPConnHandler {
return &udpHandler{
proxyHost: proxyHost,
proxyPort: proxyPort,
udpConns: make(map[core.UDPConn]net.PacketConn, 8),
tcpConns: make(map[core.UDPConn]net.Conn, 8),
remoteAddrs: make(map[core.UDPConn]*net.UDPAddr, 8),
dnsCache: dnsCache,
fakeDns: fakeDns,
timeout: timeout,
sessionStater: sessionStater,
}
}
func (h *udpHandler) handleTCP(conn core.UDPConn, c net.Conn) {
buf := core.NewBytes(core.BufSize)
defer core.FreeBytes(buf)
for {
_ = c.SetDeadline(time.Time{})
_, err := c.Read(buf)
if err == io.EOF {
log.Warnf("UDP associate to %v closed by remote", c.RemoteAddr())
h.Close(conn)
return
} else if err != nil {
h.Close(conn)
return
}
}
}
func (h *udpHandler) fetchUDPInput(conn core.UDPConn, input net.PacketConn) {
buf := core.NewBytes(core.BufSize)
defer func() {
h.Close(conn)
core.FreeBytes(buf)
}()
for {
_ = input.SetDeadline(time.Now().Add(h.timeout))
n, _, err := input.ReadFrom(buf)
if err != nil {
// log.Printf("read remote failed: %v", err)
return
}
addr := SplitAddr(buf[3:])
resolvedAddr, err := net.ResolveUDPAddr("udp", addr.String())
if err != nil {
return
}
n, err = conn.WriteFrom(buf[int(3+len(addr)):n], resolvedAddr)
if n > 0 && h.sessionStater != nil {
if sess := h.sessionStater.GetSession(conn); sess != nil {
sess.AddDownloadBytes(int64(n))
}
}
if err != nil {
log.Warnf("write local failed: %v", err)
return
}
if h.dnsCache != nil {
_, port, err := net.SplitHostPort(addr.String())
if err != nil {
panic("impossible error")
}
if port == strconv.Itoa(dns.CommonDnsPort) {
h.dnsCache.Store(buf[int(3+len(addr)):n])
return // DNS response
}
}
}
}
func (h *udpHandler) Connect(conn core.UDPConn, target *net.UDPAddr) error {
if target == nil {
return h.connectInternal(conn, "")
}
// Replace with a domain name if target address IP is a fake IP.
targetHost := target.IP.String()
if h.fakeDns != nil {
if target.Port == dns.CommonDnsPort {
return nil // skip dns
}
if h.fakeDns.IsFakeIP(target.IP) {
targetHost = h.fakeDns.IPToHost(target.IP)
}
}
dest := net.JoinHostPort(targetHost, strconv.Itoa(target.Port))
return h.connectInternal(conn, dest)
}
func (h *udpHandler) connectInternal(conn core.UDPConn, dest string) error {
c, err := net.DialTimeout("tcp", core.ParseTCPAddr(h.proxyHost, h.proxyPort).String(), 4*time.Second)
if err != nil {
return err
}
_ = c.SetDeadline(time.Now().Add(4 * time.Second))
// send VER, NMETHODS, METHODS
_, _ = c.Write([]byte{5, 1, 0})
buf := make([]byte, MaxAddrLen)
// read VER METHOD
if _, err := io.ReadFull(c, buf[:2]); err != nil {
return err
}
if len(dest) != 0 {
targetAddr := ParseAddr(dest)
// write VER CMD RSV ATYP DST.ADDR DST.PORT
_, _ = c.Write(append([]byte{5, socks5UDPAssociate, 0}, targetAddr...))
} else {
_, _ = c.Write(append([]byte{5, socks5UDPAssociate, 0}, []byte{1, 0, 0, 0, 0, 0, 0}...))
}
// read VER REP RSV ATYP BND.ADDR BND.PORT
if _, err := io.ReadFull(c, buf[:3]); err != nil {
return err
}
rep := buf[1]
if rep != 0 {
return errors.New("SOCKS handshake failed")
}
remoteAddr, err := readAddr(c, buf)
if err != nil {
return err
}
resolvedRemoteAddr, err := net.ResolveUDPAddr("udp", remoteAddr.String())
if err != nil {
return errors.New("failed to resolve remote address")
}
go h.handleTCP(conn, c)
pc, err := net.ListenPacket("udp", "")
if err != nil {
return err
}
h.Lock()
h.tcpConns[conn] = c
h.udpConns[conn] = pc
h.remoteAddrs[conn] = resolvedRemoteAddr
h.Unlock()
go h.fetchUDPInput(conn, pc)
if len(dest) != 0 {
var process string
if h.sessionStater != nil {
// Get name of the process.
localHost, localPortStr, _ := net.SplitHostPort(conn.LocalAddr().String())
localPortInt, _ := strconv.Atoi(localPortStr)
process, err = lsof.GetCommandNameBySocket(conn.LocalAddr().Network(), localHost, uint16(localPortInt))
if err != nil {
process = "unknown process"
}
sess := &stats.Session{
ProcessName: process,
Network: conn.LocalAddr().Network(),
LocalAddr: conn.LocalAddr().String(),
RemoteAddr: dest,
UploadBytes: 0,
DownloadBytes: 0,
SessionStart: time.Now(),
}
h.sessionStater.AddSession(conn, sess)
}
log.Access(process, "proxy", "udp", conn.LocalAddr().String(), dest)
}
return nil
}
func (h *udpHandler) ReceiveTo(conn core.UDPConn, data []byte, addr *net.UDPAddr) error {
h.Lock()
pc, ok1 := h.udpConns[conn]
remoteAddr, ok2 := h.remoteAddrs[conn]
h.Unlock()
// use system DNS instead of force override
if ok1 && ok2 {
var targetHost string
if h.fakeDns != nil && h.fakeDns.IsFakeIP(addr.IP) {
targetHost = h.fakeDns.IPToHost(addr.IP)
} else {
targetHost = addr.IP.String()
}
dest := net.JoinHostPort(targetHost, strconv.Itoa(addr.Port))
buf := append([]byte{0, 0, 0}, ParseAddr(dest)...)
buf = append(buf, data[:]...)
n, err := pc.WriteTo(buf, remoteAddr)
if n > 0 && h.sessionStater != nil {
if sess := h.sessionStater.GetSession(conn); sess != nil {
sess.AddUploadBytes(int64(n))
}
}
if err != nil {
h.Close(conn)
return errors.New(fmt.Sprintf("write remote failed: %v", err))
}
return nil
} else {
h.Close(conn)
return errors.New(fmt.Sprintf("proxy connection %v->%v does not exists", conn.LocalAddr(), addr))
}
}
func (h *udpHandler) Close(conn core.UDPConn) {
_ = conn.Close()
h.Lock()
defer h.Unlock()
if c, ok := h.tcpConns[conn]; ok {
_ = c.Close()
delete(h.tcpConns, conn)
}
if pc, ok := h.udpConns[conn]; ok {
_ = pc.Close()
delete(h.udpConns, conn)
}
delete(h.remoteAddrs, conn)
if h.sessionStater != nil {
h.sessionStater.RemoveSession(conn)
}
}
+61
View File
@@ -0,0 +1,61 @@
package v2ray
import (
// The following are necessary as they register handlers in their init functions.
// Required features. Can't remove unless there is replacements.
_ "v2ray.com/core/app/dispatcher"
_ "v2ray.com/core/app/proxyman/inbound"
_ "v2ray.com/core/app/proxyman/outbound"
// Default commander and all its services. This is an optional feature.
// _ "v2ray.com/core/app/commander"
// _ "v2ray.com/core/app/log/command"
// _ "v2ray.com/core/app/proxyman/command"
// _ "v2ray.com/core/app/stats/command"
// Other optional features.
_ "v2ray.com/core/app/dns"
_ "v2ray.com/core/app/log"
_ "v2ray.com/core/app/policy"
_ "v2ray.com/core/app/router"
_ "v2ray.com/core/app/stats"
// Inbound and outbound proxies.
_ "v2ray.com/core/proxy/blackhole"
_ "v2ray.com/core/proxy/dokodemo"
_ "v2ray.com/core/proxy/freedom"
_ "v2ray.com/core/proxy/http"
_ "v2ray.com/core/proxy/mtproto"
_ "v2ray.com/core/proxy/shadowsocks"
_ "v2ray.com/core/proxy/socks"
_ "v2ray.com/core/proxy/vmess/inbound"
_ "v2ray.com/core/proxy/vmess/outbound"
// Transports
_ "v2ray.com/core/transport/internet/domainsocket"
_ "v2ray.com/core/transport/internet/http"
_ "v2ray.com/core/transport/internet/kcp"
_ "v2ray.com/core/transport/internet/quic"
_ "v2ray.com/core/transport/internet/tcp"
_ "v2ray.com/core/transport/internet/tls"
_ "v2ray.com/core/transport/internet/udp"
_ "v2ray.com/core/transport/internet/websocket"
// Transport headers
_ "v2ray.com/core/transport/internet/headers/http"
_ "v2ray.com/core/transport/internet/headers/noop"
_ "v2ray.com/core/transport/internet/headers/srtp"
_ "v2ray.com/core/transport/internet/headers/tls"
_ "v2ray.com/core/transport/internet/headers/utp"
_ "v2ray.com/core/transport/internet/headers/wechat"
_ "v2ray.com/core/transport/internet/headers/wireguard"
// JSON config support. Choose only one from the two below.
// The following line loads JSON from v2ctl
// _ "v2ray.com/core/main/json"
// The following line loads JSON internally
_ "v2ray.com/core/main/jsonem"
// Load config from file or http(s)
// _ "v2ray.com/core/main/confloader/external"
)
+14
View File
@@ -0,0 +1,14 @@
// +build !ios,!android
package v2ray
import (
_ "v2ray.com/core/app/commander"
_ "v2ray.com/core/app/log/command"
_ "v2ray.com/core/app/proxyman/command"
_ "v2ray.com/core/app/stats/command"
_ "v2ray.com/core/app/reverse"
_ "v2ray.com/core/transport/internet/domainsocket"
)
+58
View File
@@ -0,0 +1,58 @@
package v2ray
import (
"context"
"errors"
"fmt"
"io"
"net"
vcore "v2ray.com/core"
vnet "v2ray.com/core/common/net"
vsession "v2ray.com/core/common/session"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/core"
)
type tcpHandler struct {
ctx context.Context
v *vcore.Instance
}
func (h *tcpHandler) handleInput(conn net.Conn, input io.ReadCloser) {
defer func() {
conn.Close()
input.Close()
}()
io.Copy(conn, input)
}
func (h *tcpHandler) handleOutput(conn net.Conn, output io.WriteCloser) {
defer func() {
conn.Close()
output.Close()
}()
io.Copy(output, conn)
}
func NewTCPHandler(ctx context.Context, instance *vcore.Instance) core.TCPConnHandler {
return &tcpHandler{
ctx: ctx,
v: instance,
}
}
func (h *tcpHandler) Handle(conn net.Conn, target *net.TCPAddr) error {
dest := vnet.DestinationFromAddr(target)
sid := vsession.NewID()
ctx := vsession.ContextWithID(h.ctx, sid)
c, err := vcore.Dial(ctx, h.v, dest)
if err != nil {
return errors.New(fmt.Sprintf("dial V proxy connection failed: %v", err))
}
go h.handleInput(conn, c)
go h.handleOutput(conn, c)
log.Infof("new proxy connection for target: %s:%s", target.Network(), target.String())
return nil
}
+138
View File
@@ -0,0 +1,138 @@
package v2ray
import (
"context"
"errors"
"fmt"
"net"
"sync"
"time"
vcore "v2ray.com/core"
vsession "v2ray.com/core/common/session"
vsignal "v2ray.com/core/common/signal"
vtask "v2ray.com/core/common/task"
"github.com/xjasonlyu/tun2socks/common/log"
"github.com/xjasonlyu/tun2socks/core"
)
type udpConnEntry struct {
conn net.PacketConn
// `ReadFrom` method of PacketConn given by V2Ray
// won't return the correct remote address, we treat
// all data receive from V2Ray are coming from the
// same remote host, i.e. the `target` that passed
// to `Connect`.
target *net.UDPAddr
updater vsignal.ActivityUpdater
}
type udpHandler struct {
sync.Mutex
ctx context.Context
v *vcore.Instance
conns map[core.UDPConn]*udpConnEntry
timeout time.Duration // Maybe override by V2Ray local policies for some conns.
}
func (h *udpHandler) fetchInput(conn core.UDPConn) {
h.Lock()
c, ok := h.conns[conn]
h.Unlock()
if !ok {
return
}
buf := core.NewBytes(core.BufSize)
defer core.FreeBytes(buf)
for {
n, _, err := c.conn.ReadFrom(buf)
if err != nil && n <= 0 {
h.Close(conn)
conn.Close()
return
}
c.updater.Update()
_, err = conn.WriteFrom(buf[:n], c.target)
if err != nil {
h.Close(conn)
conn.Close()
return
}
}
}
func NewUDPHandler(ctx context.Context, instance *vcore.Instance, timeout time.Duration) core.UDPConnHandler {
return &udpHandler{
ctx: ctx,
v: instance,
conns: make(map[core.UDPConn]*udpConnEntry, 16),
timeout: timeout,
}
}
func (h *udpHandler) Connect(conn core.UDPConn, target *net.UDPAddr) error {
if target == nil {
return errors.New("nil target is not allowed")
}
sid := vsession.NewID()
ctx := vsession.ContextWithID(h.ctx, sid)
ctx, cancel := context.WithCancel(ctx)
pc, err := vcore.DialUDP(ctx, h.v)
if err != nil {
return errors.New(fmt.Sprintf("dial V proxy connection failed: %v", err))
}
timer := vsignal.CancelAfterInactivity(ctx, cancel, h.timeout)
h.Lock()
h.conns[conn] = &udpConnEntry{
conn: pc,
target: target,
updater: timer,
}
h.Unlock()
fetchTask := func() error {
h.fetchInput(conn)
return nil
}
go func() {
if err := vtask.Run(ctx, fetchTask); err != nil {
pc.Close()
}
}()
log.Infof("new proxy connection for target: %s:%s", target.Network(), target.String())
return nil
}
func (h *udpHandler) ReceiveTo(conn core.UDPConn, data []byte, addr *net.UDPAddr) error {
h.Lock()
c, ok := h.conns[conn]
h.Unlock()
if ok {
_, err := c.conn.WriteTo(data, addr)
c.updater.Update()
if err != nil {
h.Close(conn)
return errors.New(fmt.Sprintf("write remote failed: %v", err))
}
return nil
} else {
h.Close(conn)
return errors.New(fmt.Sprintf("proxy connection %v->%v does not exists", conn.LocalAddr(), c.target))
}
}
func (h *udpHandler) Close(conn core.UDPConn) {
h.Lock()
defer h.Unlock()
if c, found := h.conns[conn]; found {
c.conn.Close()
}
delete(h.conns, conn)
}