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
This commit is contained in:
Niko Marmeladkov 2026-06-16 16:07:55 +03:00
parent 4c593412fb
commit da8366eb8a
Signed by untrusted user who does not match committer: Niko
GPG key ID: E3B955F9442D44E3
11 changed files with 1896 additions and 10 deletions

155
NEW_FEATURES.md Normal file
View file

@ -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`.

View file

@ -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) {

View file

@ -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))
}

177
extras/network/client.go Normal file
View file

@ -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
}

37
extras/network/config.go Normal file
View file

@ -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
}

73
extras/network/frame.go Normal file
View file

@ -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})
}

75
extras/network/pool.go Normal file
View file

@ -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
}

265
extras/network/server.go Normal file
View file

@ -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
}

98
extras/network/tun.go Normal file
View file

@ -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)
}

440
extras/obfs/filemask.go Normal file
View file

@ -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()
}

View file

@ -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)
}
}