From da8366eb8a1625fa8527c5bb67bb4623e543aa3f Mon Sep 17 00:00:00 2001 From: Niko Marmeladkov Date: Tue, 16 Jun 2026 16:07:55 +0300 Subject: [PATCH] feat: L3 VPN tunnel, FileMask noise, Hysteria outbound - Network: Layer 3 IP tunnel over QUIC with TUN interfaces and IP pool - FileMask: new noise/obfuscation layer masking traffic as encrypted file downloads - Hysteria outbound: chain Hysteria servers via pluggable outbound --- NEW_FEATURES.md | 155 +++++++++++ app/cmd/client.go | 35 +++ app/cmd/server.go | 240 ++++++++++++++++- extras/network/client.go | 177 +++++++++++++ extras/network/config.go | 37 +++ extras/network/frame.go | 73 ++++++ extras/network/pool.go | 75 ++++++ extras/network/server.go | 265 +++++++++++++++++++ extras/network/tun.go | 98 +++++++ extras/obfs/filemask.go | 440 ++++++++++++++++++++++++++++++++ extras/outbounds/ob_hysteria.go | 311 ++++++++++++++++++++++ 11 files changed, 1896 insertions(+), 10 deletions(-) create mode 100644 NEW_FEATURES.md create mode 100644 extras/network/client.go create mode 100644 extras/network/config.go create mode 100644 extras/network/frame.go create mode 100644 extras/network/pool.go create mode 100644 extras/network/server.go create mode 100644 extras/network/tun.go create mode 100644 extras/obfs/filemask.go create mode 100644 extras/outbounds/ob_hysteria.go diff --git a/NEW_FEATURES.md b/NEW_FEATURES.md new file mode 100644 index 0000000..51c41a6 --- /dev/null +++ b/NEW_FEATURES.md @@ -0,0 +1,155 @@ +# Новые возможности (относительно upstream) + +## 1. L3 VPN туннель (Network) + +Полноценный IP-туннель поверх QUIC. На сервере создаётся TUN-интерфейс, клиенты получают IP из пула, весь L3-трафик маршрутизируется через зашифрованное QUIC-соединение. + +### Конфиг сервера (config.json) + +```json +{ + "network": { + "enabled": true, + "listen": ":5199", + "token": "mysecret", + "tun": { + "name": "hynt0", + "mtu": 1400 + }, + "pool": "10.10.0.0/24" + } +} +``` + +Параметры: +- `enabled` — включить туннель +- `listen` — адрес для QUIC-слушателя (по умолч. `:5199`) +- `token` — общий секрет для аутентификации клиентов +- `tun.name` — имя TUN-интерфейса на сервере +- `tun.mtu` — MTU TUN-интерфейса +- `pool` — CIDR подсеть для выдачи IP клиентам + +### Конфиг клиента (client.json) + +```json +{ + "network": { + "server": "your.server.com:5199", + "token": "mysecret", + "tun": { + "name": "hynt0_c", + "mtu": 1400 + } + } +} +``` + +Параметры: +- `server` — адрес сервера (обязательно) +- `token` — общий секрет (должен совпадать с серверным) +- `tun.name` — имя TUN-интерфейса на клиенте +- `tun.mtu` — MTU TUN-интерфейса + +--- + +## 2. FileMask шум (Noise) + +Новый слой обфускации, который маскирует трафик под скачивание зашифрованных файлов. Работает поверх существующей обфускации (Salamander/Gecko). Сервер шлёт клиенту случайные куски реальных файлов с диска, имитируя активную загрузку. + +### Конфиг сервера + +```json +{ + "noise": { + "type": "filemask", + "filemask": { + "dir": "/path/to/files", + "maxRate": "512 kbps", + "minPacketSize": 512, + "maxPacketSize": 1400, + "idleThreshold": "3s" + } + } +} +``` + +Параметры: +- `type` — `"filemask"` для включения; `"none"` или пустое значение — отключено +- `filemask.dir` — **обязательно**, директория с файлами, содержимым которых будет маскироваться трафик +- `filemask.maxRate` — максимальная скорость шума (по умолч. `"512 kbps"`) +- `filemask.minPacketSize` — минимальный размер пакета шума +- `filemask.maxPacketSize` — максимальный размер пакета шума +- `filemask.idleThreshold` — таймаут бездействия клиента, после которого начинает генерироваться шум + +Файлы из указанной директории шифруются AES-GCM на лету и отправляются клиенту. + +--- + +## 3. Hysteria Outbound + +Возможность использовать другой Hysteria сервер как upstream/выходной узел (outbound). Позволяет строить цепочки: клиент -> сервер А -> сервер Б (через hysteria outbound). + +### Конфиг сервера + +```json +{ + "outbounds": [ + { + "name": "upstream-hy2", + "type": "hysteria", + "hysteria": { + "server": "upstream.example.com:443", + "auth": "mypassword", + "tls": { + "sni": "upstream.example.com", + "insecure": false, + "pinSHA256": "abc123...", + "ca": "/path/to/ca.pem" + }, + "quic": { + "initStreamReceiveWindow": 8388608, + "maxStreamReceiveWindow": 8388608, + "initConnReceiveWindow": 20971520, + "maxConnReceiveWindow": 20971520, + "maxIdleTimeout": "30s", + "disablePathMTUDiscovery": false + }, + "bandwidth": { + "up": "100 mbps", + "down": "200 mbps" + }, + "congestion": { + "type": "bbr", + "bbrProfile": "auto" + }, + "obfs": { + "type": "salamander", + "salamander": { + "password": "obfspass" + } + }, + "transport": { + "type": "udp", + "udp": { + "hopInterval": "5s" + } + }, + "fastOpen": true + } + } + ] +} +``` + +Поддерживаемые параметры внутри `hysteria`: +- `server` — адрес целевого Hysteria сервера (обязательно) +- `auth` — пароль аутентификации (обязательно) +- `tls` — TLS настройки (SNI, insecure, pinSHA256, CA) +- `quic` — настройки QUIC (окна приёма, таймауты) +- `bandwidth` — лимиты пропускной способности +- `congestion` — алгоритм контроля перегрузки (bbr/cubic/brutal) +- `obfs` — обфускация для соединения с upstream (salamander/gecko) +- `transport` — транспорт (udp, с опциональным port hopping) +- `fastOpen` — включить fast open + +Используется в разделе `outbounds` наравне с `direct`, `socks5`, `http`. diff --git a/app/cmd/client.go b/app/cmd/client.go index b23f326..9b84cc1 100644 --- a/app/cmd/client.go +++ b/app/cmd/client.go @@ -37,6 +37,7 @@ import ( "github.com/apernet/hysteria/app/v2/internal/utils" "github.com/apernet/hysteria/core/v2/client" "github.com/apernet/hysteria/extras/v2/correctnet" + "github.com/apernet/hysteria/extras/v2/network" "github.com/apernet/hysteria/extras/v2/obfs" "github.com/apernet/hysteria/extras/v2/realm" "github.com/apernet/hysteria/extras/v2/transport/udphop" @@ -89,6 +90,7 @@ type clientConfig struct { UDPTProxy *udpTProxyConfig `mapstructure:"udpTProxy"` TCPRedirect *tcpRedirectConfig `mapstructure:"tcpRedirect"` TUN *tunConfig `mapstructure:"tun"` + Network *clientConfigNetwork `mapstructure:"network"` } type clientConfigRealm struct { @@ -231,6 +233,15 @@ type tunConfig struct { } `mapstructure:"route"` } +type clientConfigNetwork struct { + Server string `mapstructure:"server"` + Token string `mapstructure:"token"` + TUN struct { + Name string `mapstructure:"name"` + MTU int `mapstructure:"mtu"` + } `mapstructure:"tun"` +} + func (c *clientConfig) fillServerAddr(hyConfig *client.Config) error { if c.Server == "" { return configError{Field: "server", Err: errors.New("server address is empty")} @@ -862,6 +873,11 @@ func runClient(v *viper.Viper) { return clientTUN(*config.TUN, c) }) } + if config.Network != nil { + runner.Add("Network", func() error { + return clientNetwork(*config.Network) + }) + } signalChan := make(chan os.Signal, 1) signal.Notify(signalChan, os.Interrupt, syscall.SIGTERM) @@ -1178,6 +1194,25 @@ func clientTUN(config tunConfig, c client.Client) error { return server.Serve() } +func clientNetwork(config clientConfigNetwork) error { + nc, err := network.NewNetworkClient(network.UpstreamConfig{ + Server: config.Server, + Token: config.Token, + }, config.TUN.Name, config.TUN.MTU) + if err != nil { + return configError{Field: "network", Err: err} + } + logger.Info("network client connecting", zap.String("server", config.Server)) + if err := nc.Connect(); err != nil { + return configError{Field: "network", Err: err} + } + logger.Info("network client connected", + zap.String("server", config.Server), + zap.String("ip", nc.AssignedIP().String())) + <-nc.Done() + return nil +} + // parseServerAddrString parses server address string. // Server address can be in either "host:port" or "host" format (in which case we assume port 443). func parseServerAddrString(addrStr string) (host, port, hostPort string) { diff --git a/app/cmd/server.go b/app/cmd/server.go index 3cde769..df19aba 100644 --- a/app/cmd/server.go +++ b/app/cmd/server.go @@ -37,10 +37,12 @@ import ( "github.com/apernet/hysteria/app/v2/internal/firewall" "github.com/apernet/hysteria/app/v2/internal/utils" + "github.com/apernet/hysteria/core/v2/client" "github.com/apernet/hysteria/core/v2/server" "github.com/apernet/hysteria/extras/v2/auth" "github.com/apernet/hysteria/extras/v2/correctnet" "github.com/apernet/hysteria/extras/v2/masq" + "github.com/apernet/hysteria/extras/v2/network" "github.com/apernet/hysteria/extras/v2/obfs" "github.com/apernet/hysteria/extras/v2/outbounds" "github.com/apernet/hysteria/extras/v2/realm" @@ -67,6 +69,7 @@ type serverConfig struct { Listen string `mapstructure:"listen"` Realm serverConfigRealm `mapstructure:"realm"` Obfs serverConfigObfs `mapstructure:"obfs"` + Noise serverConfigNoise `mapstructure:"noise"` TLS *serverConfigTLS `mapstructure:"tls"` ACME *serverConfigACME `mapstructure:"acme"` QUIC serverConfigQUIC `mapstructure:"quic"` @@ -83,6 +86,7 @@ type serverConfig struct { Outbounds []serverConfigOutboundEntry `mapstructure:"outbounds"` TrafficStats serverConfigTrafficStats `mapstructure:"trafficStats"` Masquerade serverConfigMasquerade `mapstructure:"masquerade"` + Network *serverConfigNetwork `mapstructure:"network"` } type serverConfigRealm struct { @@ -111,6 +115,19 @@ type serverConfigObfs struct { Gecko serverConfigObfsGecko `mapstructure:"gecko"` } +type serverConfigNoiseFileMask struct { + Dir string `mapstructure:"dir"` + MaxRate string `mapstructure:"maxRate"` + MinPacketSize int `mapstructure:"minPacketSize"` + MaxPacketSize int `mapstructure:"maxPacketSize"` + IdleThreshold time.Duration `mapstructure:"idleThreshold"` +} + +type serverConfigNoise struct { + Type string `mapstructure:"type"` + FileMask serverConfigNoiseFileMask `mapstructure:"filemask"` +} + type serverConfigTLS struct { Cert string `mapstructure:"cert"` Key string `mapstructure:"key"` @@ -253,12 +270,43 @@ type serverConfigOutboundHTTP struct { Insecure bool `mapstructure:"insecure"` } +type serverConfigOutboundTLS struct { + SNI string `mapstructure:"sni"` + Insecure bool `mapstructure:"insecure"` + PinSHA256 string `mapstructure:"pinSHA256"` + CA string `mapstructure:"ca"` +} + +type serverConfigOutboundTransportUDP struct { + HopInterval time.Duration `mapstructure:"hopInterval"` + MinHopInterval time.Duration `mapstructure:"minHopInterval"` + MaxHopInterval time.Duration `mapstructure:"maxHopInterval"` +} + +type serverConfigOutboundTransport struct { + Type string `mapstructure:"type"` + UDP serverConfigOutboundTransportUDP `mapstructure:"udp"` +} + +type serverConfigOutboundHysteria struct { + Server string `mapstructure:"server"` + Auth string `mapstructure:"auth"` + TLS serverConfigOutboundTLS `mapstructure:"tls"` + QUIC serverConfigQUIC `mapstructure:"quic"` + Bandwidth serverConfigBandwidth `mapstructure:"bandwidth"` + Congestion serverConfigCongestion `mapstructure:"congestion"` + Obfs serverConfigObfs `mapstructure:"obfs"` + Transport serverConfigOutboundTransport `mapstructure:"transport"` + FastOpen bool `mapstructure:"fastOpen"` +} + type serverConfigOutboundEntry struct { - Name string `mapstructure:"name"` - Type string `mapstructure:"type"` - Direct serverConfigOutboundDirect `mapstructure:"direct"` - SOCKS5 serverConfigOutboundSOCKS5 `mapstructure:"socks5"` - HTTP serverConfigOutboundHTTP `mapstructure:"http"` + Name string `mapstructure:"name"` + Type string `mapstructure:"type"` + Direct serverConfigOutboundDirect `mapstructure:"direct"` + SOCKS5 serverConfigOutboundSOCKS5 `mapstructure:"socks5"` + HTTP serverConfigOutboundHTTP `mapstructure:"http"` + Hysteria serverConfigOutboundHysteria `mapstructure:"hysteria"` } type serverConfigTrafficStats struct { @@ -283,6 +331,19 @@ type serverConfigMasqueradeString struct { StatusCode int `mapstructure:"statusCode"` } +type serverConfigNetworkTUN struct { + Name string `mapstructure:"name"` + MTU int `mapstructure:"mtu"` +} + +type serverConfigNetwork struct { + Enabled bool `mapstructure:"enabled"` + Listen string `mapstructure:"listen"` + Token string `mapstructure:"token"` + TUN serverConfigNetworkTUN `mapstructure:"tun"` + Pool string `mapstructure:"pool"` +} + type serverConfigMasquerade struct { Type string `mapstructure:"type"` File serverConfigMasqueradeFile `mapstructure:"file"` @@ -390,17 +451,18 @@ func (c *serverConfig) fillRealmConn(hyConfig *server.Config, addr *realm.Addr) } func (c *serverConfig) wrapObfs(conn net.PacketConn) (net.PacketConn, error) { + packetConn := conn + switch strings.ToLower(c.Obfs.Type) { case "", "plain": - return conn, nil case "salamander": - wrapped, err := obfs.WrapPacketConnSalamander(conn, []byte(c.Obfs.Salamander.Password)) + wrapped, err := obfs.WrapPacketConnSalamander(packetConn, []byte(c.Obfs.Salamander.Password)) if err != nil { return nil, configError{Field: "obfs.salamander.password", Err: err} } - return wrapped, nil + packetConn = wrapped case "gecko": - wrapped, err := obfs.WrapPacketConnGecko(conn, obfs.GeckoOptions{ + wrapped, err := obfs.WrapPacketConnGecko(packetConn, obfs.GeckoOptions{ Password: []byte(c.Obfs.Gecko.Password), MinPacketSize: c.Obfs.Gecko.MinPacketSize, MaxPacketSize: c.Obfs.Gecko.MaxPacketSize, @@ -408,10 +470,38 @@ func (c *serverConfig) wrapObfs(conn net.PacketConn) (net.PacketConn, error) { if err != nil { return nil, configError{Field: "obfs.gecko", Err: err} } - return wrapped, nil + packetConn = wrapped default: return nil, configError{Field: "obfs.type", Err: errors.New("unsupported obfuscation type")} } + + switch strings.ToLower(c.Noise.Type) { + case "", "none": + case "filemask": + maxRate := c.Noise.FileMask.MaxRate + if maxRate == "" { + maxRate = "512 kbps" + } + maxRateBPS, err := utils.ConvBandwidth(maxRate) + if err != nil { + return nil, configError{Field: "noise.filemask.maxRate", Err: err} + } + wrapped, err := obfs.WrapPacketConnFileMask(packetConn, obfs.FileMaskOptions{ + Dir: c.Noise.FileMask.Dir, + MaxRate: int(maxRateBPS), + MinPacketSize: c.Noise.FileMask.MinPacketSize, + MaxPacketSize: c.Noise.FileMask.MaxPacketSize, + IdleThreshold: c.Noise.FileMask.IdleThreshold, + }) + if err != nil { + return nil, configError{Field: "noise.filemask", Err: err} + } + packetConn = wrapped + default: + return nil, configError{Field: "noise.type", Err: errors.New("unsupported noise type")} + } + + return packetConn, nil } func resolveServerListenAddr(listenAddr string) (*net.UDPAddr, eUtils.PortUnion, error) { @@ -1206,6 +1296,109 @@ func serverConfigOutboundHTTPToOutbound(c serverConfigOutboundHTTP) (outbounds.P return outbounds.NewHTTPOutbound(c.URL, c.Insecure) } +func serverConfigOutboundHysteriaToOutbound(c serverConfigOutboundHysteria) (outbounds.PluggableOutbound, error) { + if c.Server == "" { + return nil, configError{Field: "outbounds.hysteria.server", Err: errors.New("empty server address")} + } + if c.Auth == "" { + return nil, configError{Field: "outbounds.hysteria.auth", Err: errors.New("empty auth")} + } + // TLS config + tlsCfg := outbounds.HysteriaOutboundTLSConfig{ + ServerName: c.TLS.SNI, + InsecureSkipVerify: c.TLS.Insecure, + PinSHA256: c.TLS.PinSHA256, + } + if c.TLS.CA != "" { + ca, err := os.ReadFile(c.TLS.CA) + if err != nil { + return nil, configError{Field: "outbounds.hysteria.tls.ca", Err: err} + } + cPool := x509.NewCertPool() + if !cPool.AppendCertsFromPEM(ca) { + return nil, configError{Field: "outbounds.hysteria.tls.ca", Err: errors.New("failed to parse CA certificate")} + } + tlsCfg.RootCAs = cPool + } + // Normalize congestion type + normalizedType, err := normalizeCongestionType(c.Congestion.Type) + if err != nil { + return nil, configError{Field: "outbounds.hysteria.congestion.type", Err: err} + } + var bbrProfile string + if normalizedType == congestionTypeBBR { + bbrProfile, err = normalizeBBRProfile(c.Congestion.BBRProfile) + if err != nil { + return nil, configError{Field: "outbounds.hysteria.congestion.bbrProfile", Err: err} + } + } + // Obfuscation config + obfsPassword := "" + switch strings.ToLower(c.Obfs.Type) { + case "salamander": + obfsPassword = c.Obfs.Salamander.Password + case "gecko": + obfsPassword = c.Obfs.Gecko.Password + } + cfg := &outbounds.HysteriaOutboundConfig{ + ServerAddr: c.Server, + Auth: c.Auth, + TLSConfig: tlsCfg, + QUICConfig: client.QUICConfig{}, + BandwidthConfig: client.BandwidthConfig{}, + CongestionType: normalizedType, + BBRProfile: bbrProfile, + FastOpen: c.FastOpen, + Obfs: outbounds.HysteriaOutboundObfsConfig{ + Type: c.Obfs.Type, + Password: obfsPassword, + MinPacketSize: c.Obfs.Gecko.MinPacketSize, + MaxPacketSize: c.Obfs.Gecko.MaxPacketSize, + }, + Transport: outbounds.HysteriaOutboundTransportConfig{ + Type: c.Transport.Type, + HopInterval: c.Transport.UDP.HopInterval, + MinHopInterval: c.Transport.UDP.MinHopInterval, + MaxHopInterval: c.Transport.UDP.MaxHopInterval, + }, + } + // QUIC config + if c.QUIC.InitStreamReceiveWindow != 0 { + cfg.QUICConfig.InitialStreamReceiveWindow = c.QUIC.InitStreamReceiveWindow + } + if c.QUIC.MaxStreamReceiveWindow != 0 { + cfg.QUICConfig.MaxStreamReceiveWindow = c.QUIC.MaxStreamReceiveWindow + } + if c.QUIC.InitConnectionReceiveWindow != 0 { + cfg.QUICConfig.InitialConnectionReceiveWindow = c.QUIC.InitConnectionReceiveWindow + } + if c.QUIC.MaxConnectionReceiveWindow != 0 { + cfg.QUICConfig.MaxConnectionReceiveWindow = c.QUIC.MaxConnectionReceiveWindow + } + if c.QUIC.MaxIdleTimeout != 0 { + cfg.QUICConfig.MaxIdleTimeout = c.QUIC.MaxIdleTimeout + } + if c.QUIC.DisablePathMTUDiscovery { + cfg.QUICConfig.DisablePathMTUDiscovery = true + } + // Bandwidth config + if c.Bandwidth.Up != "" { + up, err := utils.ConvBandwidth(c.Bandwidth.Up) + if err != nil { + return nil, configError{Field: "outbounds.hysteria.bandwidth.up", Err: err} + } + cfg.BandwidthConfig.MaxTx = up + } + if c.Bandwidth.Down != "" { + down, err := utils.ConvBandwidth(c.Bandwidth.Down) + if err != nil { + return nil, configError{Field: "outbounds.hysteria.bandwidth.down", Err: err} + } + cfg.BandwidthConfig.MaxRx = down + } + return outbounds.NewHysteriaOutbound(cfg) +} + func (c *serverConfig) fillRequestHook(hyConfig *server.Config) error { if c.Sniff.Enable { s := &sniff.Sniffer{ @@ -1257,6 +1450,8 @@ func (c *serverConfig) fillOutboundConfig(hyConfig *server.Config) error { ob, err = serverConfigOutboundSOCKS5ToOutbound(entry.SOCKS5) case "http": ob, err = serverConfigOutboundHTTPToOutbound(entry.HTTP) + case "hysteria": + ob, err = serverConfigOutboundHysteriaToOutbound(entry.Hysteria) default: err = configError{Field: "outbounds.type", Err: errors.New("unsupported outbound type")} } @@ -1604,6 +1799,28 @@ func runServer(v *viper.Viper) { go runCheckUpdateServer() } + var ns *network.NetworkServer + if config.Network != nil && config.Network.Enabled { + ns, err = network.NewNetworkServer(network.Config{ + Enabled: true, + Listen: config.Network.Listen, + Token: config.Network.Token, + TUN: network.TUNConfig{ + Name: config.Network.TUN.Name, + MTU: config.Network.TUN.MTU, + }, + Pool: config.Network.Pool, + }) + if err != nil { + logger.Fatal("failed to initialize network server", zap.Error(err)) + } + go func() { + if err := ns.Serve(); err != nil { + logger.Fatal("network server error", zap.Error(err)) + } + }() + } + signalChan := make(chan os.Signal, 1) signal.Notify(signalChan, os.Interrupt, syscall.SIGTERM) defer signal.Stop(signalChan) @@ -1616,6 +1833,9 @@ func runServer(v *viper.Viper) { select { case <-signalChan: logger.Info("received signal, shutting down gracefully") + if ns != nil { + ns.Close() + } if err := s.Close(); err != nil { logger.Error("failed to shut down server cleanly", zap.Error(err)) } diff --git a/extras/network/client.go b/extras/network/client.go new file mode 100644 index 0000000..dd8d2ef --- /dev/null +++ b/extras/network/client.go @@ -0,0 +1,177 @@ +package network + +import ( + "context" + "fmt" + "log" + "net" + "sync" + "time" + + "github.com/apernet/quic-go" +) + +type NetworkClient struct { + config UpstreamConfig + tun *TUNInterface + assignedIP net.IP + gateway net.IP + + stream *quic.Stream + conn *quic.Conn + ctx context.Context + cancel context.CancelFunc + wg sync.WaitGroup +} + +func NewNetworkClient(config UpstreamConfig, tunName string, mtu int) (*NetworkClient, error) { + if tunName == "" { + tunName = defaultTUNName + "_c" + } + if mtu <= 0 { + mtu = defaultTUNMTU + } + ctx, cancel := context.WithCancel(context.Background()) + return &NetworkClient{ + config: config, + ctx: ctx, + cancel: cancel, + }, nil +} + +func (c *NetworkClient) Connect() error { + tlsCfg := c.config.TLS + if tlsCfg == nil { + var err error + tlsCfg, err = GenerateNetworkTLSConfig() + if err != nil { + return fmt.Errorf("generate TLS: %w", err) + } + tlsCfg.InsecureSkipVerify = true + } + conn, err := quic.DialAddr(c.ctx, c.config.Server, tlsCfg, nil) + if err != nil { + return fmt.Errorf("dial: %w", err) + } + c.conn = conn + + stream, err := conn.OpenStream() + if err != nil { + conn.CloseWithError(1, "stream error") + return fmt.Errorf("open stream: %w", err) + } + c.stream = stream + + if err := SendAuth(stream, c.config.Token); err != nil { + return fmt.Errorf("send auth: %w", err) + } + + f, err := ReceiveFrame(stream) + if err != nil { + return fmt.Errorf("recv auth: %w", err) + } + switch f.Type { + case FrameAuthOK: + if len(f.Payload) < 4 { + return fmt.Errorf("invalid IP in auth response") + } + c.assignedIP = net.IP(f.Payload[:4]) + c.gateway = make(net.IP, 4) + copy(c.gateway, c.assignedIP) + c.gateway[3] = 1 + log.Printf("network: assigned IP %s, gateway %s", c.assignedIP, c.gateway) + case FrameAuthErr: + return fmt.Errorf("auth rejected: %s", string(f.Payload)) + default: + return fmt.Errorf("unexpected auth response: %d", f.Type) + } + + c.tun, err = OpenTUN(defaultTUNName+"_c", c.assignedIP.String()+"/32", defaultTUNMTU) + if err != nil { + return fmt.Errorf("open TUN: %w", err) + } + c.tun.Route(c.assignedIP.String() + "/32") + + go func() { + ticker := time.NewTicker(defaultKeepalive) + defer ticker.Stop() + for { + select { + case <-c.ctx.Done(): + return + case <-ticker.C: + SendKeepalive(stream) + } + } + }() + + c.wg.Add(2) + go c.tunToStream() + go c.streamToTun() + + return nil +} + +func (c *NetworkClient) tunToStream() { + defer c.wg.Done() + for { + pkt, err := c.tun.Read() + if err != nil { + if c.ctx.Err() != nil { + return + } + log.Printf("TUN read: %v", err) + return + } + if err := SendData(c.stream, pkt); err != nil { + log.Printf("send: %v", err) + return + } + } +} + +func (c *NetworkClient) streamToTun() { + defer c.wg.Done() + for { + f, err := ReceiveFrame(c.stream) + if err != nil { + if c.ctx.Err() != nil { + return + } + log.Printf("recv: %v", err) + return + } + switch f.Type { + case FrameData: + if err := c.tun.Write(f.Payload); err != nil { + log.Printf("TUN write: %v", err) + return + } + case FrameKeepalive: + default: + } + } +} + +func (c *NetworkClient) AssignedIP() net.IP { + return c.assignedIP +} + +func (c *NetworkClient) Done() <-chan struct{} { + return c.ctx.Done() +} + +func (c *NetworkClient) Close() error { + c.cancel() + if c.stream != nil { + c.stream.Close() + } + if c.conn != nil { + c.conn.CloseWithError(0, "bye") + } + if c.tun != nil { + c.tun.Close() + } + c.wg.Wait() + return nil +} diff --git a/extras/network/config.go b/extras/network/config.go new file mode 100644 index 0000000..339021d --- /dev/null +++ b/extras/network/config.go @@ -0,0 +1,37 @@ +package network + +import ( + "crypto/tls" + "time" +) + +const ( + defaultNetworkListen = ":5199" + defaultTUNName = "hynt0" + defaultTUNMTU = 1400 + defaultPoolCIDR = "10.10.0.0/24" + defaultKeepalive = 30 * time.Second + readBufferSize = 65535 +) + +type Config struct { + Enabled bool + Listen string + Token string + TUN TUNConfig + Pool string + Keepalive time.Duration + TLS *tls.Config + Upstream *UpstreamConfig +} + +type TUNConfig struct { + Name string + MTU int +} + +type UpstreamConfig struct { + Server string + Token string + TLS *tls.Config +} diff --git a/extras/network/frame.go b/extras/network/frame.go new file mode 100644 index 0000000..7f8879e --- /dev/null +++ b/extras/network/frame.go @@ -0,0 +1,73 @@ +package network + +import ( + "encoding/binary" + "fmt" + "io" + "net" + + "github.com/apernet/quic-go" +) + +type FrameType byte + +const ( + FrameAuth FrameType = 0x01 + FrameAuthOK FrameType = 0x02 + FrameAuthErr FrameType = 0x03 + FrameData FrameType = 0x04 + FrameKeepalive FrameType = 0x05 +) + +type Frame struct { + Type FrameType + Payload []byte +} + +func SendFrame(stream *quic.Stream, f *Frame) error { + buf := make([]byte, 4+1+len(f.Payload)) + binary.BigEndian.PutUint32(buf[:4], uint32(1+len(f.Payload))) + buf[4] = byte(f.Type) + copy(buf[5:], f.Payload) + _, err := stream.Write(buf) + return err +} + +func ReceiveFrame(stream *quic.Stream) (*Frame, error) { + var lenBuf [4]byte + if _, err := io.ReadFull(stream, lenBuf[:]); err != nil { + return nil, err + } + frameLen := binary.BigEndian.Uint32(lenBuf[:]) + if frameLen < 1 || frameLen > 65535 { + return nil, fmt.Errorf("invalid frame length: %d", frameLen) + } + frameData := make([]byte, frameLen) + if _, err := io.ReadFull(stream, frameData); err != nil { + return nil, err + } + return &Frame{ + Type: FrameType(frameData[0]), + Payload: frameData[1:], + }, nil +} + +func SendAuth(stream *quic.Stream, token string) error { + return SendFrame(stream, &Frame{Type: FrameAuth, Payload: []byte(token)}) +} + +func SendAuthOK(stream *quic.Stream, assignedIP net.IP) error { + return SendFrame(stream, &Frame{Type: FrameAuthOK, Payload: assignedIP.To4()}) +} + +func SendAuthErr(stream *quic.Stream, reason string) error { + return SendFrame(stream, &Frame{Type: FrameAuthErr, Payload: []byte(reason)}) +} + +func SendData(stream *quic.Stream, pkt []byte) error { + return SendFrame(stream, &Frame{Type: FrameData, Payload: pkt}) +} + +func SendKeepalive(stream *quic.Stream) error { + return SendFrame(stream, &Frame{Type: FrameKeepalive}) +} diff --git a/extras/network/pool.go b/extras/network/pool.go new file mode 100644 index 0000000..58d6073 --- /dev/null +++ b/extras/network/pool.go @@ -0,0 +1,75 @@ +package network + +import ( + "fmt" + "net" + "sync" +) + +type IPPool struct { + network *net.IPNet + gateway net.IP + leased map[string]bool + mu sync.Mutex + start net.IP + size int +} + +func NewIPPool(netIP net.IP, ipnet *net.IPNet) (*IPPool, error) { + ones, bits := ipnet.Mask.Size() + if bits != 32 { + return nil, fmt.Errorf("only IPv4 supported") + } + size := 1 << (bits - ones) + if size < 3 { + return nil, fmt.Errorf("pool too small: %s", ipnet) + } + start := make(net.IP, 4) + copy(start, netIP.Mask(ipnet.Mask)) + start[3]++ + gateway := make(net.IP, 4) + copy(gateway, start) + pool := &IPPool{ + network: ipnet, + gateway: gateway, + leased: make(map[string]bool), + start: start, + size: size - 2, + } + pool.leased[gateway.String()] = true + return pool, nil +} + +func (p *IPPool) Allocate() net.IP { + p.mu.Lock() + defer p.mu.Unlock() + ip := make(net.IP, 4) + copy(ip, p.start) + for i := 0; i < p.size; i++ { + if !p.leased[ip.String()] { + p.leased[ip.String()] = true + return ip + } + for j := 3; j >= 0; j-- { + ip[j]++ + if ip[j] != 0 { + break + } + } + } + return nil +} + +func (p *IPPool) Release(ip net.IP) { + p.mu.Lock() + defer p.mu.Unlock() + delete(p.leased, ip.String()) +} + +func (p *IPPool) Gateway() net.IP { + return p.gateway +} + +func (p *IPPool) Network() *net.IPNet { + return p.network +} diff --git a/extras/network/server.go b/extras/network/server.go new file mode 100644 index 0000000..ac9e838 --- /dev/null +++ b/extras/network/server.go @@ -0,0 +1,265 @@ +package network + +import ( + "context" + "crypto/rand" + "crypto/rsa" + "crypto/tls" + "crypto/x509" + "crypto/x509/pkix" + "fmt" + "log" + "math/big" + "net" + "os/exec" + "sync" + "time" + + "github.com/apernet/quic-go" +) + +type NetworkServer struct { + config Config + tun *TUNInterface + pool *IPPool + + ln *quic.Listener + clients map[string]*networkClient + clientsMu sync.Mutex + ctx context.Context + cancel context.CancelFunc + wg sync.WaitGroup +} + +type networkClient struct { + ip net.IP + stream *quic.Stream +} + +func NewNetworkServer(config Config) (*NetworkServer, error) { + ip, ipnet, err := net.ParseCIDR(config.Pool) + if err != nil { + return nil, fmt.Errorf("parse pool: %w", err) + } + pool, err := NewIPPool(ip, ipnet) + if err != nil { + return nil, fmt.Errorf("create pool: %w", err) + } + tunName := config.TUN.Name + if tunName == "" { + tunName = defaultTUNName + } + mtu := config.TUN.MTU + if mtu <= 0 { + mtu = defaultTUNMTU + } + tunDev, err := OpenTUN(tunName, pool.Gateway().String()+"/24", mtu) + if err != nil { + return nil, fmt.Errorf("open TUN: %w", err) + } + ctx, cancel := context.WithCancel(context.Background()) + return &NetworkServer{ + config: config, + tun: tunDev, + pool: pool, + clients: make(map[string]*networkClient), + ctx: ctx, + cancel: cancel, + }, nil +} + +func (s *NetworkServer) Serve() error { + listenAddr := s.config.Listen + if listenAddr == "" { + listenAddr = defaultNetworkListen + } + tlsCfg := s.config.TLS + if tlsCfg == nil { + var err error + tlsCfg, err = GenerateNetworkTLSConfig() + if err != nil { + return fmt.Errorf("generate TLS: %w", err) + } + } + ln, err := quic.ListenAddr(listenAddr, tlsCfg, nil) + if err != nil { + return fmt.Errorf("listen: %w", err) + } + s.ln = ln + log.Printf("network server listening on %s", listenAddr) + log.Printf("TUN %s: %s", s.tun.Name, s.pool.Gateway()) + log.Printf("IP pool: %s", s.config.Pool) + + s.wg.Add(1) + go s.tunToClients() + + for { + conn, err := ln.Accept(s.ctx) + if err != nil { + if s.ctx.Err() != nil { + return nil + } + return fmt.Errorf("accept: %w", err) + } + go s.handleConn(conn) + } +} + +func (s *NetworkServer) handleConn(conn *quic.Conn) { + stream, err := conn.AcceptStream(s.ctx) + if err != nil { + log.Printf("accept stream: %v", err) + return + } + assignedIP, err := s.handleAuth(stream) + if err != nil { + log.Printf("auth error from %s: %v", conn.RemoteAddr(), err) + stream.Close() + return + } + s.tun.Route(assignedIP.String() + "/32") + log.Printf("client %s authenticated, IP: %s", conn.RemoteAddr(), assignedIP) + + nc := &networkClient{ + ip: assignedIP, + stream: stream, + } + s.clientsMu.Lock() + s.clients[assignedIP.String()] = nc + s.clientsMu.Unlock() + + stopCh := make(chan struct{}) + go func() { + ticker := time.NewTicker(defaultKeepalive) + defer ticker.Stop() + for { + select { + case <-stopCh: + return + case <-ticker.C: + SendKeepalive(stream) + } + } + }() + + s.clientToTun(nc) + close(stopCh) + + s.clientsMu.Lock() + delete(s.clients, assignedIP.String()) + s.clientsMu.Unlock() + s.pool.Release(assignedIP) + exec.Command("ip", "route", "del", assignedIP.String()+"/32", "dev", s.tun.Name).Run() + log.Printf("client %s (%s) disconnected", conn.RemoteAddr(), assignedIP) +} + +func (s *NetworkServer) handleAuth(stream *quic.Stream) (net.IP, error) { + f, err := ReceiveFrame(stream) + if err != nil { + return nil, fmt.Errorf("recv auth: %w", err) + } + if f.Type != FrameAuth { + return nil, fmt.Errorf("expected auth, got %d", f.Type) + } + if string(f.Payload) != s.config.Token { + SendAuthErr(stream, "invalid token") + return nil, fmt.Errorf("invalid token") + } + assignedIP := s.pool.Allocate() + if assignedIP == nil { + SendAuthErr(stream, "no IP available") + return nil, fmt.Errorf("no IP available") + } + if err := SendAuthOK(stream, assignedIP); err != nil { + s.pool.Release(assignedIP) + return nil, fmt.Errorf("send auth ok: %w", err) + } + return assignedIP, nil +} + +func (s *NetworkServer) clientToTun(nc *networkClient) { + for { + f, err := ReceiveFrame(nc.stream) + if err != nil { + return + } + switch f.Type { + case FrameData: + if err := s.tun.Write(f.Payload); err != nil { + log.Printf("write TUN: %v", err) + return + } + case FrameKeepalive: + default: + } + } +} + +func (s *NetworkServer) tunToClients() { + defer s.wg.Done() + for { + pkt, err := s.tun.Read() + if err != nil { + if s.ctx.Err() != nil { + return + } + log.Printf("read TUN: %v", err) + return + } + if len(pkt) < 20 { + continue + } + dstIP := net.IP(pkt[16:20]).String() + s.clientsMu.Lock() + nc, ok := s.clients[dstIP] + s.clientsMu.Unlock() + if !ok { + continue + } + if err := SendData(nc.stream, pkt); err != nil { + log.Printf("send to %s: %v", dstIP, err) + return + } + } +} + +func (s *NetworkServer) Close() error { + s.cancel() + if s.ln != nil { + s.ln.Close() + } + s.tun.Close() + s.wg.Wait() + return nil +} + +func GenerateNetworkTLSConfig() (*tls.Config, error) { + key, err := rsa.GenerateKey(rand.Reader, 2048) + if err != nil { + return nil, fmt.Errorf("generate key: %w", err) + } + template := x509.Certificate{ + SerialNumber: big.NewInt(1), + Subject: pkix.Name{ + Organization: []string{"Hysteria Network"}, + }, + NotBefore: time.Now(), + NotAfter: time.Now().Add(10 * 365 * 24 * time.Hour), + KeyUsage: x509.KeyUsageKeyEncipherment | x509.KeyUsageDigitalSignature, + ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth}, + BasicConstraintsValid: true, + } + template.DNSNames = append(template.DNSNames, "localhost") + certDER, err := x509.CreateCertificate(rand.Reader, &template, &template, &key.PublicKey, key) + if err != nil { + return nil, fmt.Errorf("create cert: %w", err) + } + cert := tls.Certificate{ + Certificate: [][]byte{certDER}, + PrivateKey: key, + } + return &tls.Config{ + Certificates: []tls.Certificate{cert}, + NextProtos: []string{"hysteria-network"}, + }, nil +} diff --git a/extras/network/tun.go b/extras/network/tun.go new file mode 100644 index 0000000..fc03d08 --- /dev/null +++ b/extras/network/tun.go @@ -0,0 +1,98 @@ +package network + +import ( + "fmt" + "os" + "os/exec" + "syscall" + "unsafe" + + "golang.org/x/sys/unix" +) + +const ( + TUNSETIFF = 0x400454ca + TUNSETOFFLOAD = 0x400454d0 + IFF_TUN = 0x0001 + IFF_NO_PI = 0x1000 +) + +type ifReq struct { + Name [16]byte + Flags uint16 + pad [22]byte +} + +type TUNInterface struct { + Name string + Address string + MTU int + fd int +} + +func OpenTUN(name string, address string, mtu int) (*TUNInterface, error) { + exec.Command("modprobe", "tun").Run() + if _, err := os.Stat("/dev/net/tun"); err != nil { + return nil, fmt.Errorf("TUN not available: %w", err) + } + fd, err := unix.Open("/dev/net/tun", unix.O_RDWR, 0) + if err != nil { + return nil, fmt.Errorf("open /dev/net/tun: %w", err) + } + var req ifReq + req.Flags = IFF_TUN | IFF_NO_PI + copy(req.Name[:], name) + _, _, errno := syscall.Syscall(syscall.SYS_IOCTL, uintptr(fd), uintptr(TUNSETIFF), uintptr(unsafe.Pointer(&req))) + if errno != 0 { + unix.Close(fd) + return nil, fmt.Errorf("TUNSETIFF: %w", errno) + } + offload := 0 + syscall.Syscall(syscall.SYS_IOCTL, uintptr(fd), uintptr(TUNSETOFFLOAD), uintptr(unsafe.Pointer(&offload))) + ifName := string(req.Name[:]) + ifName = ifName[:len(name)] + t := &TUNInterface{ + Name: ifName, + MTU: mtu, + fd: fd, + } + if err := t.configure(address); err != nil { + t.Close() + return nil, fmt.Errorf("configure TUN: %w", err) + } + return t, nil +} + +func (t *TUNInterface) configure(address string) error { + if err := exec.Command("ip", "link", "set", "dev", t.Name, "mtu", fmt.Sprintf("%d", t.MTU), "up").Run(); err != nil { + return fmt.Errorf("set link up: %w", err) + } + if address != "" { + if err := exec.Command("ip", "addr", "add", address, "dev", t.Name).Run(); err != nil { + return fmt.Errorf("set addr: %w", err) + } + } + return nil +} + +func (t *TUNInterface) Route(cidr string) error { + return exec.Command("ip", "route", "replace", cidr, "dev", t.Name).Run() +} + +func (t *TUNInterface) Read() ([]byte, error) { + buf := make([]byte, t.MTU*2) + n, err := unix.Read(t.fd, buf) + if err != nil { + return nil, err + } + return buf[:n], nil +} + +func (t *TUNInterface) Write(pkt []byte) error { + _, err := unix.Write(t.fd, pkt) + return err +} + +func (t *TUNInterface) Close() error { + return unix.Close(t.fd) +} diff --git a/extras/obfs/filemask.go b/extras/obfs/filemask.go new file mode 100644 index 0000000..6772c70 --- /dev/null +++ b/extras/obfs/filemask.go @@ -0,0 +1,440 @@ +package obfs + +import ( + "context" + "crypto/aes" + "crypto/cipher" + "crypto/rand" + "fmt" + "io/fs" + "net" + "os" + "path/filepath" + "sync" + "sync/atomic" + "syscall" + "time" +) + +const ( + fmDefaultMinPacket = 512 + fmDefaultMaxPacket = 1400 + fmDefaultIdleThresh = 3 * time.Second + fmDefaultMaxRate = 512 * 1024 + fmDefaultChunkSize = 1400 + fmHeaderPackets = 4 + fmMinPause = 30 * time.Millisecond + fmMaxPause = 150 * time.Millisecond + fmLoopInterval = 100 * time.Millisecond + fmMinDownloadBytes = 512 * 1024 + fmMaxDownloadBytes = 10 * 1024 * 1024 + fmCleanupAge = 2 * time.Minute + fmQuicCIDLen = 8 + +) + +var _ net.PacketConn = (*FileMaskPacketConn)(nil) + +type FileMaskOptions struct { + Dir string `mapstructure:"dir"` + MaxRate int `mapstructure:"maxRate"` + MinPacketSize int `mapstructure:"minPacketSize"` + MaxPacketSize int `mapstructure:"maxPacketSize"` + IdleThreshold time.Duration `mapstructure:"idleThreshold"` +} + +func (o *FileMaskOptions) fill() { + if o.MinPacketSize <= 0 { + o.MinPacketSize = fmDefaultMinPacket + } + if o.MaxPacketSize <= 0 { + o.MaxPacketSize = fmDefaultMaxPacket + } + if o.MinPacketSize > o.MaxPacketSize { + o.MaxPacketSize = o.MinPacketSize + } + if o.IdleThreshold <= 0 { + o.IdleThreshold = fmDefaultIdleThresh + } + if o.MaxRate <= 0 { + o.MaxRate = fmDefaultMaxRate + } +} + +type downloadPhase int + +const ( + dpIdle downloadPhase = iota + dpHeader + dpStream + dpPause +) + +type clientState struct { + addr net.Addr + activeUnix atomic.Int64 + + phase downloadPhase + connID [fmQuicCIDLen]byte + bytesSent int64 + bytesGoal int64 + pktCount int + + pauseDeadline time.Time + headerLeft int +} + +func (cs *clientState) setActive(t time.Time) { + cs.activeUnix.Store(t.UnixNano()) +} + +func (cs *clientState) idleSince(t time.Time) time.Duration { + return t.Sub(time.Unix(0, cs.activeUnix.Load())) +} + +type FileMaskPacketConn struct { + inner net.PacketConn + options FileMaskOptions + + chunks [][]byte + chunksMu sync.RWMutex + + clients map[string]*clientState + clientsMu sync.Mutex + + bytesThisSec atomic.Int64 + secStart atomic.Int64 + + ctx context.Context + cancel context.CancelFunc + wg sync.WaitGroup +} + +func WrapPacketConnFileMask(conn net.PacketConn, opts FileMaskOptions) (net.PacketConn, error) { + if !opts.Enabled() { + return conn, nil + } + opts.fill() + + chunks, err := loadEncryptFiles(opts.Dir) + if err != nil { + return nil, fmt.Errorf("filemask: %w", err) + } + if len(chunks) == 0 { + return nil, fmt.Errorf("filemask: no files found in %s", opts.Dir) + } + + ctx, cancel := context.WithCancel(context.Background()) + f := &FileMaskPacketConn{ + inner: conn, + options: opts, + chunks: chunks, + clients: make(map[string]*clientState), + ctx: ctx, + cancel: cancel, + } + + f.secStart.Store(time.Now().Unix()) + + f.wg.Add(1) + go f.noiseLoop() + + if u, ok := conn.(udpLikePacketConn); ok { + return &fileMaskPacketConnUDP{ + FileMaskPacketConn: f, + UDPConn: u, + }, nil + } + return f, nil +} + +func (o *FileMaskOptions) Enabled() bool { + return o.Dir != "" +} + +func loadEncryptFiles(dir string) ([][]byte, error) { + key := make([]byte, 32) + if _, err := rand.Read(key); err != nil { + return nil, fmt.Errorf("failed to generate AES key: %w", err) + } + + block, err := aes.NewCipher(key) + if err != nil { + return nil, fmt.Errorf("failed to create AES cipher: %w", err) + } + aesGCM, err := cipher.NewGCM(block) + if err != nil { + return nil, fmt.Errorf("failed to create AES-GCM: %w", err) + } + + var chunks [][]byte + err = filepath.WalkDir(dir, func(path string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + if d.IsDir() { + return nil + } + if d.Type().IsRegular() { + data, readErr := os.ReadFile(path) + if readErr != nil { + return readErr + } + if len(data) == 0 { + return nil + } + encrypted := aesGCM.Seal(nil, makeNonce(aesGCM), data, nil) + for offset := 0; offset < len(encrypted); offset += fmDefaultChunkSize { + end := offset + fmDefaultChunkSize + if end > len(encrypted) { + end = len(encrypted) + } + chunk := make([]byte, end-offset) + copy(chunk, encrypted[offset:end]) + chunks = append(chunks, chunk) + } + } + return nil + }) + if err != nil { + return nil, fmt.Errorf("walk dir: %w", err) + } + + return chunks, nil +} + +func makeNonce(aead cipher.AEAD) []byte { + nonce := make([]byte, aead.NonceSize()) + rand.Read(nonce) + return nonce +} + +func (f *FileMaskPacketConn) noisePacket(chunk []byte, connID []byte) []byte { + firstByte := byte(0x40 | randIntn(64)) + pkt := make([]byte, 1+len(connID)+len(chunk)) + pkt[0] = firstByte + copy(pkt[1:], connID) + copy(pkt[1+len(connID):], chunk) + return pkt +} + +func (f *FileMaskPacketConn) WriteTo(p []byte, addr net.Addr) (int, error) { + f.trackClient(addr) + + _, err := f.inner.WriteTo(p, addr) + if err != nil { + return 0, err + } + return len(p), nil +} + +func (f *FileMaskPacketConn) ReadFrom(p []byte) (int, net.Addr, error) { + return f.inner.ReadFrom(p) +} + +func (f *FileMaskPacketConn) Close() error { + f.cancel() + f.wg.Wait() + return f.inner.Close() +} + +func (f *FileMaskPacketConn) LocalAddr() net.Addr { + return f.inner.LocalAddr() +} + +func (f *FileMaskPacketConn) SetDeadline(t time.Time) error { + return f.inner.SetDeadline(t) +} + +func (f *FileMaskPacketConn) SetReadDeadline(t time.Time) error { + return f.inner.SetReadDeadline(t) +} + +func (f *FileMaskPacketConn) SetWriteDeadline(t time.Time) error { + return f.inner.SetWriteDeadline(t) +} + +func (f *FileMaskPacketConn) trackClient(addr net.Addr) { + if addr == nil { + return + } + addrStr := addr.String() + f.clientsMu.Lock() + cs, exists := f.clients[addrStr] + if !exists { + cs = &clientState{ + addr: addr, + connID: newConnID(), + } + f.clients[addrStr] = cs + } + cs.setActive(time.Now()) + f.clientsMu.Unlock() +} + +func newConnID() [fmQuicCIDLen]byte { + var id [fmQuicCIDLen]byte + rand.Read(id[:]) + return id +} + +func (f *FileMaskPacketConn) noiseLoop() { + ticker := time.NewTicker(fmLoopInterval) + defer ticker.Stop() + defer f.wg.Done() + + for { + select { + case <-f.ctx.Done(): + return + case <-ticker.C: + f.injectNoise() + } + } +} + +func (f *FileMaskPacketConn) injectNoise() { + f.clientsMu.Lock() + entries := make([]*clientState, 0, len(f.clients)) + for _, cs := range f.clients { + entries = append(entries, cs) + } + f.clientsMu.Unlock() + + now := time.Now() + for _, cs := range entries { + f.simulateDownload(cs, now) + } + + f.cleanupClients(now) +} + +func (f *FileMaskPacketConn) simulateDownload(cs *clientState, now time.Time) { + switch cs.phase { + case dpIdle: + if cs.idleSince(now) < f.options.IdleThreshold { + return + } + cs.phase = dpHeader + cs.headerLeft = fmHeaderPackets + randIntn(3) + cs.bytesSent = 0 + cs.bytesGoal = int64(fmMinDownloadBytes + randIntn(fmMaxDownloadBytes-fmMinDownloadBytes)) + cs.pktCount = 0 + cs.connID = newConnID() + + case dpHeader: + if cs.headerLeft > 0 { + chunk := f.randomChunk() + targetSize := 300 + randIntn(300) + if len(chunk) > targetSize { + chunk = chunk[:targetSize] + } + pkt := f.noisePacket(chunk, cs.connID[:]) + if f.allowNoise(len(pkt)) { + f.inner.WriteTo(pkt, cs.addr) + cs.bytesSent += int64(len(chunk)) + cs.pktCount++ + } + cs.headerLeft-- + return + } + cs.phase = dpStream + + case dpStream: + if cs.pktCount > 0 && cs.pktCount%(fmSimPauseEvery()+randIntn(10)) == 0 { + cs.phase = dpPause + pauseDur := fmMinPause + time.Duration(randIntn(int(fmMaxPause-fmMinPause))) + cs.pauseDeadline = now.Add(pauseDur) + return + } + + if cs.bytesSent >= cs.bytesGoal { + cs.phase = dpIdle + return + } + + chunk := f.randomChunk() + var chunkLen int + if randIntn(10) < 8 { + chunkLen = min(len(chunk), f.options.MaxPacketSize) + } else { + chunkLen = min(64+randIntn(256), len(chunk)) + } + if chunkLen > len(chunk) { + chunkLen = len(chunk) + } + chunk = chunk[:chunkLen] + pkt := f.noisePacket(chunk, cs.connID[:]) + if f.allowNoise(len(pkt)) { + f.inner.WriteTo(pkt, cs.addr) + cs.bytesSent += int64(len(chunk)) + cs.pktCount++ + } + + case dpPause: + if now.Before(cs.pauseDeadline) { + return + } + cs.phase = dpStream + } +} + +func fmSimPauseEvery() int { + return 15 + randIntn(15) +} + +func (f *FileMaskPacketConn) randomChunk() []byte { + f.chunksMu.RLock() + defer f.chunksMu.RUnlock() + if len(f.chunks) == 0 { + dummy := make([]byte, 64) + rand.Read(dummy) + return dummy + } + idx := randIntn(len(f.chunks)) + return f.chunks[idx] +} + +func (f *FileMaskPacketConn) allowNoise(size int) bool { + now := time.Now().Unix() + if f.secStart.Load() != now { + f.bytesThisSec.Store(0) + f.secStart.Store(now) + } + + cur := f.bytesThisSec.Load() + if cur+int64(size) > int64(f.options.MaxRate) { + return false + } + f.bytesThisSec.Add(int64(size)) + return true +} + +func (f *FileMaskPacketConn) cleanupClients(now time.Time) { + f.clientsMu.Lock() + defer f.clientsMu.Unlock() + for addr, cs := range f.clients { + if cs.idleSince(now) > fmCleanupAge { + delete(f.clients, addr) + } + } +} + +type fileMaskPacketConnUDP struct { + *FileMaskPacketConn + UDPConn udpLikePacketConn +} + +func (c *fileMaskPacketConnUDP) SetReadBuffer(bytes int) error { + return c.UDPConn.SetReadBuffer(bytes) +} + +func (c *fileMaskPacketConnUDP) SetWriteBuffer(bytes int) error { + return c.UDPConn.SetWriteBuffer(bytes) +} + +func (c *fileMaskPacketConnUDP) SyscallConn() (syscall.RawConn, error) { + return c.UDPConn.SyscallConn() +} + + diff --git a/extras/outbounds/ob_hysteria.go b/extras/outbounds/ob_hysteria.go new file mode 100644 index 0000000..e3fdef8 --- /dev/null +++ b/extras/outbounds/ob_hysteria.go @@ -0,0 +1,311 @@ +package outbounds + +import ( + "crypto/sha256" + "crypto/tls" + "crypto/x509" + "encoding/hex" + "errors" + "net" + "strconv" + "strings" + "sync" + "time" + + "github.com/apernet/hysteria/core/v2/client" + "github.com/apernet/hysteria/extras/v2/obfs" + "github.com/apernet/hysteria/extras/v2/transport/udphop" +) + +type HysteriaOutboundObfsConfig struct { + Type string + Password string + MinPacketSize int + MaxPacketSize int +} + +type HysteriaOutboundTransportConfig struct { + Type string + HopInterval time.Duration + MinHopInterval time.Duration + MaxHopInterval time.Duration +} + +type HysteriaOutboundConfig struct { + ServerAddr string + Auth string + TLSConfig HysteriaOutboundTLSConfig + QUICConfig client.QUICConfig + BandwidthConfig client.BandwidthConfig + CongestionType string + BBRProfile string + FastOpen bool + + Obfs HysteriaOutboundObfsConfig + Transport HysteriaOutboundTransportConfig +} + +type HysteriaOutboundTLSConfig struct { + ServerName string + InsecureSkipVerify bool + PinSHA256 string + RootCAs *x509.CertPool +} + +// hysteriaOutbound is a PluggableOutbound that connects to the target using +// another Hysteria server. It acts as a Hysteria client internally, forwarding +// traffic through the next hop server. +type hysteriaOutbound struct { + configFunc func() (*client.Config, error) + client client.Client + + mu sync.Mutex + closed bool +} + +func NewHysteriaOutbound(cfg *HysteriaOutboundConfig) (PluggableOutbound, error) { + if cfg.ServerAddr == "" { + return nil, errors.New("hysteria outbound: server address is empty") + } + if cfg.Auth == "" { + return nil, errors.New("hysteria outbound: auth is empty") + } + + ho := &hysteriaOutbound{ + configFunc: func() (*client.Config, error) { + serverAddr, connFactory, err := buildServerAddrAndFactory(cfg) + if err != nil { + return nil, err + } + + tlsCfg := &tls.Config{ + ServerName: cfg.TLSConfig.ServerName, + InsecureSkipVerify: cfg.TLSConfig.InsecureSkipVerify, + RootCAs: cfg.TLSConfig.RootCAs, + } + + if cfg.TLSConfig.PinSHA256 != "" { + nHash := normalizeCertHash(cfg.TLSConfig.PinSHA256) + tlsCfg.VerifyPeerCertificate = func(rawCerts [][]byte, _ [][]*x509.Certificate) error { + cert := rawCerts[0] + hash := sha256.Sum256(cert) + hashHex := hex.EncodeToString(hash[:]) + if hashHex == nHash { + return nil + } + return errors.New("no certificate matches the pinned hash") + } + } + + hyClientCfg := &client.Config{ + ServerAddr: serverAddr, + ConnFactory: connFactory, + Auth: cfg.Auth, + QUICConfig: cfg.QUICConfig, + BandwidthConfig: cfg.BandwidthConfig, + CongestionConfig: client.CongestionConfig{Type: cfg.CongestionType, BBRProfile: cfg.BBRProfile}, + FastOpen: cfg.FastOpen, + } + hyClientCfg.TLSConfig.ServerName = tlsCfg.ServerName + hyClientCfg.TLSConfig.InsecureSkipVerify = tlsCfg.InsecureSkipVerify + hyClientCfg.TLSConfig.VerifyPeerCertificate = tlsCfg.VerifyPeerCertificate + hyClientCfg.TLSConfig.RootCAs = tlsCfg.RootCAs + return hyClientCfg, nil + }, + } + + rc, err := client.NewReconnectableClient(ho.configFunc, nil, true) + if err != nil { + return nil, err + } + ho.client = rc + return ho, nil +} + +func (h *hysteriaOutbound) TCP(reqAddr *AddrEx) (net.Conn, error) { + h.mu.Lock() + if h.closed { + h.mu.Unlock() + return nil, errors.New("hysteria outbound: closed") + } + h.mu.Unlock() + return h.client.TCP(reqAddr.String()) +} + +func (h *hysteriaOutbound) UDP(reqAddr *AddrEx) (UDPConn, error) { + h.mu.Lock() + if h.closed { + h.mu.Unlock() + return nil, errors.New("hysteria outbound: closed") + } + h.mu.Unlock() + hyUDP, err := h.client.UDP() + if err != nil { + return nil, err + } + return &hyUDPConnAdapter{HyUDPConn: hyUDP}, nil +} + +func (h *hysteriaOutbound) CheckUDP(reqAddr *AddrEx) error { + return nil +} + +func (h *hysteriaOutbound) Close() error { + h.mu.Lock() + defer h.mu.Unlock() + if h.closed { + return nil + } + h.closed = true + return h.client.Close() +} + +type hyUDPConnAdapter struct { + HyUDPConn client.HyUDPConn +} + +func (a *hyUDPConnAdapter) ReadFrom(b []byte) (int, *AddrEx, error) { + data, addrStr, err := a.HyUDPConn.Receive() + if err != nil { + return 0, nil, err + } + n := copy(b, data) + host, portStr, err := net.SplitHostPort(addrStr) + if err != nil { + return n, &AddrEx{Host: addrStr}, nil + } + port, _ := strconv.ParseUint(portStr, 10, 16) + return n, &AddrEx{Host: host, Port: uint16(port)}, nil +} + +func (a *hyUDPConnAdapter) WriteTo(b []byte, addr *AddrEx) (int, error) { + return len(b), a.HyUDPConn.Send(b, addr.String()) +} + +func (a *hyUDPConnAdapter) Close() error { + return a.HyUDPConn.Close() +} + +func normalizeCertHash(hash string) string { + if len(hash) == 64 { + return hash + } + if len(hash) == 59 && hash[:7] == "sha256/" { + return hash[7:] + } + return hash +} + +// buildServerAddrAndFactory resolves the server address and creates a ConnFactory +// that applies obfuscation and port hopping as configured. +func buildServerAddrAndFactory(cfg *HysteriaOutboundConfig) (net.Addr, client.ConnFactory, error) { + host, portStr, err := net.SplitHostPort(cfg.ServerAddr) + if err != nil { + return nil, nil, err + } + + // Resolve host to IP + ip, err := net.ResolveIPAddr("ip", host) + if err != nil { + return nil, nil, err + } + + // Determine if port hopping is used + isHop := strings.ContainsAny(portStr, "-,") + + if isHop { + return buildHopAddrAndFactory(cfg, ip, portStr) + } + return buildPlainAddrAndFactory(cfg, ip, portStr) +} + +func buildPlainAddrAndFactory(cfg *HysteriaOutboundConfig, ip *net.IPAddr, portStr string) (net.Addr, client.ConnFactory, error) { + port, err := strconv.ParseUint(portStr, 10, 16) + if err != nil { + return nil, nil, err + } + addr := &net.UDPAddr{IP: ip.IP, Port: int(port)} + factory := &hysteriaConnFactory{ + obfsType: cfg.Obfs.Type, + obfsPassword: cfg.Obfs.Password, + obfsMinPacketSize: cfg.Obfs.MinPacketSize, + obfsMaxPacketSize: cfg.Obfs.MaxPacketSize, + } + return addr, factory, nil +} + +func buildHopAddrAndFactory(cfg *HysteriaOutboundConfig, ip *net.IPAddr, portStr string) (net.Addr, client.ConnFactory, error) { + hopAddr, err := udphop.ResolveUDPHopAddr(net.JoinHostPort(ip.IP.String(), portStr)) + if err != nil { + return nil, nil, err + } + factory := &hysteriaConnFactory{ + hopAddr: hopAddr, + hopInterval: hopIntervalFromConfig(cfg.Transport), + obfsType: cfg.Obfs.Type, + obfsPassword: cfg.Obfs.Password, + obfsMinPacketSize: cfg.Obfs.MinPacketSize, + obfsMaxPacketSize: cfg.Obfs.MaxPacketSize, + } + return hopAddr, factory, nil +} + +func hopIntervalFromConfig(cfg HysteriaOutboundTransportConfig) udphop.HopIntervalConfig { + if cfg.HopInterval != 0 { + return udphop.HopIntervalConfig{Min: cfg.HopInterval, Max: cfg.HopInterval} + } + return udphop.HopIntervalConfig{Min: cfg.MinHopInterval, Max: cfg.MaxHopInterval} +} + +// hysteriaConnFactory implements client.ConnFactory, creating a PacketConn +// optionally wrapped with port hopping and/or obfuscation. +type hysteriaConnFactory struct { + hopAddr *udphop.UDPHopAddr + hopInterval udphop.HopIntervalConfig + obfsType string + obfsPassword string + obfsMinPacketSize int + obfsMaxPacketSize int +} + +func (f *hysteriaConnFactory) New(addr net.Addr) (net.PacketConn, error) { + var conn net.PacketConn + var err error + + if f.hopAddr != nil { + conn, err = udphop.NewUDPHopPacketConn(f.hopAddr, f.hopInterval, nil) + if err != nil { + return nil, err + } + } else { + conn, err = net.ListenUDP("udp", nil) + if err != nil { + return nil, err + } + } + + if f.obfsType != "" { + conn, err = wrapObfs(conn, f.obfsType, f.obfsPassword, f.obfsMinPacketSize, f.obfsMaxPacketSize) + if err != nil { + _ = conn.Close() + return nil, err + } + } + + return conn, nil +} + +func wrapObfs(conn net.PacketConn, obfsType, password string, minPacketSize, maxPacketSize int) (net.PacketConn, error) { + switch strings.ToLower(obfsType) { + case "salamander": + return obfs.WrapPacketConnSalamander(conn, []byte(password)) + case "gecko": + return obfs.WrapPacketConnGecko(conn, obfs.GeckoOptions{ + Password: []byte(password), + MinPacketSize: minPacketSize, + MaxPacketSize: maxPacketSize, + }) + default: + return nil, errors.New("unsupported obfuscation type: " + obfsType) + } +}