2020-11-25 19:01:53 +08:00
package inbound
import (
2023-01-06 05:37:16 +00:00
"bytes"
2020-11-25 19:01:53 +08:00
"context"
2023-01-26 22:43:58 -05:00
gotls "crypto/tls"
2025-08-28 04:55:36 +00:00
"encoding/base64"
2020-11-25 19:01:53 +08:00
"io"
2023-01-06 05:37:16 +00:00
"reflect"
2020-11-25 19:01:53 +08:00
"strconv"
2021-01-13 23:13:51 +08:00
"strings"
2020-11-25 19:01:53 +08:00
"time"
2023-01-06 05:37:16 +00:00
"unsafe"
2020-11-25 19:01:53 +08:00
2025-12-31 06:00:45 -05:00
"github.com/xtls/xray-core/app/dispatcher"
2025-09-09 14:19:12 +00:00
"github.com/xtls/xray-core/app/reverse"
2020-12-04 09:36:16 +08:00
"github.com/xtls/xray-core/common"
"github.com/xtls/xray-core/common/buf"
"github.com/xtls/xray-core/common/errors"
"github.com/xtls/xray-core/common/log"
2025-09-09 14:19:12 +00:00
"github.com/xtls/xray-core/common/mux"
2020-12-04 09:36:16 +08:00
"github.com/xtls/xray-core/common/net"
"github.com/xtls/xray-core/common/protocol"
"github.com/xtls/xray-core/common/retry"
2025-09-09 14:19:12 +00:00
"github.com/xtls/xray-core/common/serial"
2020-12-04 09:36:16 +08:00
"github.com/xtls/xray-core/common/session"
"github.com/xtls/xray-core/common/signal"
"github.com/xtls/xray-core/common/task"
2023-02-15 16:07:12 +00:00
"github.com/xtls/xray-core/core"
2026-03-21 19:16:24 +08:00
"github.com/xtls/xray-core/features"
2020-12-04 09:36:16 +08:00
"github.com/xtls/xray-core/features/dns"
2026-03-21 19:16:24 +08:00
"github.com/xtls/xray-core/features/extension"
2020-12-04 09:36:16 +08:00
feature_inbound "github.com/xtls/xray-core/features/inbound"
2025-09-09 14:19:12 +00:00
"github.com/xtls/xray-core/features/outbound"
2020-12-04 09:36:16 +08:00
"github.com/xtls/xray-core/features/policy"
"github.com/xtls/xray-core/features/routing"
2025-12-31 06:00:45 -05:00
"github.com/xtls/xray-core/features/stats"
2023-09-02 11:37:50 -04:00
"github.com/xtls/xray-core/proxy"
2020-12-04 09:36:16 +08:00
"github.com/xtls/xray-core/proxy/vless"
"github.com/xtls/xray-core/proxy/vless/encoding"
2025-08-28 04:55:36 +00:00
"github.com/xtls/xray-core/proxy/vless/encryption"
2025-09-01 11:15:32 -04:00
"github.com/xtls/xray-core/transport"
2023-02-15 16:07:12 +00:00
"github.com/xtls/xray-core/transport/internet/reality"
2021-12-14 19:28:47 -05:00
"github.com/xtls/xray-core/transport/internet/stat"
2020-12-04 09:36:16 +08:00
"github.com/xtls/xray-core/transport/internet/tls"
2020-11-25 19:01:53 +08:00
)
func init () {
common . Must ( common . RegisterConfig (( * Config )( nil ), func ( ctx context . Context , config interface {}) ( interface {}, error ) {
var dc dns . Client
if err := core . RequireFeatures ( ctx , func ( d dns . Client ) error {
dc = d
return nil
}); err != nil {
return nil , err
}
2024-09-13 17:51:26 +03:00
c := config .( * Config )
validator := new ( vless . MemoryValidator )
2026-05-07 19:10:48 +08:00
for _ , user := range c . Users {
2024-09-13 17:51:26 +03:00
u , err := user . ToMemoryUser ()
if err != nil {
return nil , errors . New ( "failed to get VLESS user" ). Base ( err ). AtError ()
}
if err := validator . Add ( u ); err != nil {
return nil , errors . New ( "failed to initiate user" ). Base ( err ). AtError ()
}
}
return New ( ctx , c , dc , validator )
2020-11-25 19:01:53 +08:00
}))
}
// Handler is an inbound connection handler that handles messages in VLess protocol.
type Handler struct {
2025-09-09 14:19:12 +00:00
inboundHandlerManager feature_inbound . Manager
policyManager policy . Manager
2025-12-31 06:00:45 -05:00
stats stats . Manager
2025-09-09 14:19:12 +00:00
validator vless . Validator
decryption * encryption . ServerInstance
outboundHandlerManager outbound . Manager
2026-03-21 19:16:24 +08:00
observer features . Feature
2025-12-31 06:00:45 -05:00
defaultDispatcher routing . Dispatcher
2025-09-09 14:19:12 +00:00
ctx context . Context
fallbacks map [ string ] map [ string ] map [ string ] * Fallback // or nil
2020-11-25 19:01:53 +08:00
// regexps map[string]*regexp.Regexp // or nil
}
// New creates a new VLess inbound handler.
2024-09-13 17:51:26 +03:00
func New ( ctx context . Context , config * Config , dc dns . Client , validator vless . Validator ) ( * Handler , error ) {
2020-11-25 19:01:53 +08:00
v := core . MustFromContext ( ctx )
handler := & Handler {
2025-09-09 14:19:12 +00:00
inboundHandlerManager : v . GetFeature ( feature_inbound . ManagerType ()).( feature_inbound . Manager ),
policyManager : v . GetFeature ( policy . ManagerType ()).( policy . Manager ),
2025-12-31 06:00:45 -05:00
stats : v . GetFeature ( stats . ManagerType ()).( stats . Manager ),
2025-09-09 14:19:12 +00:00
validator : validator ,
outboundHandlerManager : v . GetFeature ( outbound . ManagerType ()).( outbound . Manager ),
2026-03-21 19:16:24 +08:00
observer : v . GetFeature ( extension . ObservatoryType ()),
2025-12-31 06:00:45 -05:00
defaultDispatcher : v . GetFeature ( routing . DispatcherType ()).( routing . Dispatcher ),
2025-09-09 14:19:12 +00:00
ctx : ctx ,
2020-11-25 19:01:53 +08:00
}
2025-08-28 04:55:36 +00:00
if config . Decryption != "" && config . Decryption != "none" {
s := strings . Split ( config . Decryption , "." )
var nfsSKeysBytes [][] byte
for _ , r := range s {
b , _ := base64 . RawURLEncoding . DecodeString ( r )
nfsSKeysBytes = append ( nfsSKeysBytes , b )
}
handler . decryption = & encryption . ServerInstance {}
2025-09-02 23:37:14 +00:00
if err := handler . decryption . Init ( nfsSKeysBytes , config . XorMode , config . SecondsFrom , config . SecondsTo , config . Padding ); err != nil {
2025-08-28 04:55:36 +00:00
return nil , errors . New ( "failed to use decryption" ). Base ( err ). AtError ()
}
}
2020-11-25 19:01:53 +08:00
if config . Fallbacks != nil {
2021-01-13 23:13:51 +08:00
handler . fallbacks = make ( map [ string ] map [ string ] map [ string ] * Fallback )
2020-11-25 19:01:53 +08:00
// handler.regexps = make(map[string]*regexp.Regexp)
for _ , fb := range config . Fallbacks {
2021-01-13 23:13:51 +08:00
if handler . fallbacks [ fb . Name ] == nil {
handler . fallbacks [ fb . Name ] = make ( map [ string ] map [ string ] * Fallback )
2020-11-25 19:01:53 +08:00
}
2021-01-13 23:13:51 +08:00
if handler . fallbacks [ fb . Name ][ fb . Alpn ] == nil {
handler . fallbacks [ fb . Name ][ fb . Alpn ] = make ( map [ string ] * Fallback )
}
handler . fallbacks [ fb . Name ][ fb . Alpn ][ fb . Path ] = fb
2020-11-25 19:01:53 +08:00
/*
if fb.Path != "" {
if r, err := regexp.Compile(fb.Path); err != nil {
2024-06-29 14:32:57 -04:00
return nil, errors.New("invalid path regexp").Base(err).AtError()
2020-11-25 19:01:53 +08:00
} else {
handler.regexps[fb.Path] = r
}
}
*/
}
2021-01-15 11:36:31 +00:00
if handler . fallbacks [ "" ] != nil {
for name , apfb := range handler . fallbacks {
if name != "" {
for alpn := range handler . fallbacks [ "" ] {
if apfb [ alpn ] == nil {
apfb [ alpn ] = make ( map [ string ] * Fallback )
}
}
}
}
}
2021-01-14 21:55:52 +00:00
for _ , apfb := range handler . fallbacks {
if apfb [ "" ] != nil {
for alpn , pfb := range apfb {
if alpn != "" { // && alpn != "h2" {
for path , fb := range apfb [ "" ] {
if pfb [ path ] == nil {
pfb [ path ] = fb
}
}
}
}
}
}
2020-11-25 19:01:53 +08:00
if handler . fallbacks [ "" ] != nil {
2021-01-14 21:55:52 +00:00
for name , apfb := range handler . fallbacks {
if name != "" {
for alpn , pfb := range handler . fallbacks [ "" ] {
for path , fb := range pfb {
if apfb [ alpn ][ path ] == nil {
apfb [ alpn ][ path ] = fb
}
2020-11-25 19:01:53 +08:00
}
}
}
}
}
}
return handler , nil
}
2023-01-28 00:39:36 -05:00
func isMuxAndNotXUDP ( request * protocol . RequestHeader , first * buf . Buffer ) bool {
if request . Command != protocol . RequestCommandMux {
return false
}
if first . Len () < 7 {
return true
}
firstBytes := first . Bytes ()
return !( firstBytes [ 2 ] == 0 && // ID high
firstBytes [ 3 ] == 0 && // ID low
firstBytes [ 6 ] == 2 ) // Network type: UDP
}
2025-09-09 14:19:12 +00:00
func ( h * Handler ) GetReverse ( a * vless . MemoryAccount ) ( * Reverse , error ) {
u := h . validator . Get ( a . ID . UUID ())
if u == nil {
return nil , errors . New ( "reverse: user " + a . ID . String () + " doesn't exist anymore" )
}
a = u . Account .( * vless . MemoryAccount )
if a . Reverse == nil || a . Reverse . Tag == "" {
return nil , errors . New ( "reverse: user " + a . ID . String () + " is not allowed to create reverse proxy" )
}
r := h . outboundHandlerManager . GetHandler ( a . Reverse . Tag )
if r == nil {
picker , _ := reverse . NewStaticMuxPicker ()
r = & Reverse { tag : a . Reverse . Tag , picker : picker , client : & mux . ClientManager { Picker : picker }}
for len ( h . outboundHandlerManager . ListHandlers ( h . ctx )) == 0 {
time . Sleep ( time . Second ) // prevents this outbound from becoming the default outbound
}
if err := h . outboundHandlerManager . AddHandler ( h . ctx , r ); err != nil {
return nil , err
}
}
if r , ok := r .( * Reverse ); ok {
return r , nil
}
return nil , errors . New ( "reverse: outbound " + a . Reverse . Tag + " is not type Reverse" )
}
func ( h * Handler ) RemoveReverse ( u * protocol . MemoryUser ) {
if u != nil {
a := u . Account .( * vless . MemoryAccount )
if a . Reverse != nil && a . Reverse . Tag != "" {
h . outboundHandlerManager . RemoveHandler ( h . ctx , a . Reverse . Tag )
}
}
}
2020-11-25 19:01:53 +08:00
// Close implements common.Closable.Close().
func ( h * Handler ) Close () error {
2025-09-04 14:03:55 +00:00
if h . decryption != nil {
h . decryption . Close ()
}
2025-09-09 14:19:12 +00:00
for _ , u := range h . validator . GetAll () {
h . RemoveReverse ( u )
}
2020-11-25 19:01:53 +08:00
return errors . Combine ( common . Close ( h . validator ))
}
// AddUser implements proxy.UserManager.AddUser().
func ( h * Handler ) AddUser ( ctx context . Context , u * protocol . MemoryUser ) error {
return h . validator . Add ( u )
}
// RemoveUser implements proxy.UserManager.RemoveUser().
func ( h * Handler ) RemoveUser ( ctx context . Context , e string ) error {
2025-09-09 14:19:12 +00:00
h . RemoveReverse ( h . validator . GetByEmail ( e ))
2020-11-25 19:01:53 +08:00
return h . validator . Del ( e )
}
2024-11-03 00:25:23 -04:00
// GetUser implements proxy.UserManager.GetUser().
func ( h * Handler ) GetUser ( ctx context . Context , email string ) * protocol . MemoryUser {
return h . validator . GetByEmail ( email )
}
// GetUsers implements proxy.UserManager.GetUsers().
func ( h * Handler ) GetUsers ( ctx context . Context ) [] * protocol . MemoryUser {
return h . validator . GetAll ()
}
// GetUsersCount implements proxy.UserManager.GetUsersCount().
func ( h * Handler ) GetUsersCount ( context . Context ) int64 {
return h . validator . GetCount ()
}
2020-11-25 19:01:53 +08:00
// Network implements proxy.Inbound.Network().
func ( * Handler ) Network () [] net . Network {
return [] net . Network { net . Network_TCP , net . Network_UNIX }
}
// Process implements proxy.Inbound.Process().
2025-12-31 06:00:45 -05:00
func ( h * Handler ) Process ( ctx context . Context , network net . Network , connection stat . Connection , dispatch routing . Dispatcher ) error {
2025-12-23 17:44:54 +08:00
iConn := stat . TryUnwrapStatsConn ( connection )
2020-11-25 19:01:53 +08:00
2025-08-28 04:55:36 +00:00
if h . decryption != nil {
var err error
2025-08-31 04:09:28 +00:00
if connection , err = h . decryption . Handshake ( connection , nil ); err != nil {
2025-08-28 04:55:36 +00:00
return errors . New ( "ML-KEM-768 handshake failed" ). Base ( err ). AtInfo ()
}
}
2020-11-25 19:01:53 +08:00
sessionPolicy := h . policyManager . ForLevel ( 0 )
if err := connection . SetReadDeadline ( time . Now (). Add ( sessionPolicy . Timeouts . Handshake )); err != nil {
2024-06-29 14:32:57 -04:00
return errors . New ( "unable to set read deadline" ). Base ( err ). AtWarning ()
2020-11-25 19:01:53 +08:00
}
2023-01-17 11:18:58 +08:00
first := buf . FromBytes ( make ([] byte , buf . Size ))
first . Clear ()
2025-07-23 14:53:37 +02:00
firstLen , errR := first . ReadFrom ( connection )
if errR != nil {
return errR
}
2024-06-29 14:32:57 -04:00
errors . LogInfo ( ctx , "firstLen = " , firstLen )
2020-11-25 19:01:53 +08:00
reader := & buf . BufferedReader {
Reader : buf . NewReader ( connection ),
Buffer : buf . MultiBuffer { first },
}
2025-08-17 18:13:56 +00:00
var userSentID [] byte // not MemoryAccount.ID
2020-11-25 19:01:53 +08:00
var request * protocol . RequestHeader
var requestAddons * encoding . Addons
var err error
2021-01-14 21:55:52 +00:00
napfb := h . fallbacks
isfb := napfb != nil
2020-11-25 19:01:53 +08:00
if isfb && firstLen < 18 {
2024-06-29 14:32:57 -04:00
err = errors . New ( "fallback directly" )
2020-11-25 19:01:53 +08:00
} else {
2025-08-17 18:13:56 +00:00
userSentID , request , requestAddons , isfb , err = encoding . DecodeRequestHeader ( isfb , first , reader , h . validator )
2020-11-25 19:01:53 +08:00
}
if err != nil {
if isfb {
if err := connection . SetReadDeadline ( time . Time {}); err != nil {
2024-06-29 14:32:57 -04:00
errors . LogWarningInner ( ctx , err , "unable to set back read deadline" )
2020-11-25 19:01:53 +08:00
}
2024-06-29 14:32:57 -04:00
errors . LogInfoInner ( ctx , err , "fallback starts" )
2020-11-25 19:01:53 +08:00
2021-01-13 23:13:51 +08:00
name := ""
2020-11-25 19:01:53 +08:00
alpn := ""
2021-01-14 21:55:52 +00:00
if tlsConn , ok := iConn .( * tls . Conn ); ok {
cs := tlsConn . ConnectionState ()
name = cs . ServerName
alpn = cs . NegotiatedProtocol
2024-07-12 00:20:06 +02:00
errors . LogInfo ( ctx , "realName = " + name )
errors . LogInfo ( ctx , "realAlpn = " + alpn )
2023-02-15 16:07:12 +00:00
} else if realityConn , ok := iConn .( * reality . Conn ); ok {
cs := realityConn . ConnectionState ()
name = cs . ServerName
alpn = cs . NegotiatedProtocol
2024-07-12 00:20:06 +02:00
errors . LogInfo ( ctx , "realName = " + name )
errors . LogInfo ( ctx , "realAlpn = " + alpn )
2021-01-14 21:55:52 +00:00
}
2021-01-22 07:37:55 +08:00
name = strings . ToLower ( name )
alpn = strings . ToLower ( alpn )
2021-01-14 21:55:52 +00:00
if len ( napfb ) > 1 || napfb [ "" ] == nil {
2021-01-15 09:43:39 +00:00
if name != "" && napfb [ name ] == nil {
match := ""
for n := range napfb {
if n != "" && strings . Contains ( name , n ) && len ( n ) > len ( match ) {
match = n
}
2021-01-13 23:13:51 +08:00
}
2021-01-15 09:43:39 +00:00
name = match
2021-01-13 23:13:51 +08:00
}
2021-01-14 21:55:52 +00:00
}
if napfb [ name ] == nil {
name = ""
}
apfb := napfb [ name ]
if apfb == nil {
2024-06-29 14:32:57 -04:00
return errors . New ( `failed to find the default "name" config` ). AtWarning ()
2021-01-14 21:55:52 +00:00
}
2021-01-13 23:13:51 +08:00
2021-01-14 21:55:52 +00:00
if apfb [ alpn ] == nil {
alpn = ""
2020-11-25 19:01:53 +08:00
}
2021-01-14 21:55:52 +00:00
pfb := apfb [ alpn ]
2020-11-25 19:01:53 +08:00
if pfb == nil {
2024-06-29 14:32:57 -04:00
return errors . New ( `failed to find the default "alpn" config` ). AtWarning ()
2020-11-25 19:01:53 +08:00
}
path := ""
if len ( pfb ) > 1 || pfb [ "" ] == nil {
/*
if lines := bytes.Split(firstBytes, []byte{'\r', '\n'}); len(lines) > 1 {
if s := bytes.Split(lines[0], []byte{' '}); len(s) == 3 {
if len(s[0]) < 8 && len(s[1]) > 0 && len(s[2]) == 8 {
2024-06-29 14:32:57 -04:00
errors.New("realPath = " + string(s[1])).AtInfo().WriteToLog(sid)
2020-11-25 19:01:53 +08:00
for _, fb := range pfb {
if fb.Path != "" && h.regexps[fb.Path].Match(s[1]) {
path = fb.Path
break
}
}
}
}
}
*/
if firstLen >= 18 && first . Byte ( 4 ) != '*' { // not h2c
firstBytes := first . Bytes ()
for i := 4 ; i <= 8 ; i ++ { // 5 -> 9
if firstBytes [ i ] == '/' && firstBytes [ i - 1 ] == ' ' {
search := len ( firstBytes )
if search > 64 {
search = 64 // up to about 60
}
for j := i + 1 ; j < search ; j ++ {
k := firstBytes [ j ]
if k == '\r' || k == '\n' { // avoid logging \r or \n
break
}
2021-03-12 11:50:59 +00:00
if k == '?' || k == ' ' {
2020-11-25 19:01:53 +08:00
path = string ( firstBytes [ i : j ])
2024-07-12 00:20:06 +02:00
errors . LogInfo ( ctx , "realPath = " + path )
2020-11-25 19:01:53 +08:00
if pfb [ path ] == nil {
path = ""
}
break
}
}
break
}
}
}
}
fb := pfb [ path ]
if fb == nil {
2024-06-29 14:32:57 -04:00
return errors . New ( `failed to find the default "path" config` ). AtWarning ()
2020-11-25 19:01:53 +08:00
}
ctx , cancel := context . WithCancel ( ctx )
timer := signal . CancelAfterInactivity ( ctx , cancel , sessionPolicy . Timeouts . ConnectionIdle )
ctx = policy . ContextWithBufferPolicy ( ctx , sessionPolicy . Buffer )
var conn net . Conn
if err := retry . ExponentialBackoff ( 5 , 100 ). On ( func () error {
var dialer net . Dialer
conn , err = dialer . DialContext ( ctx , fb . Type , fb . Dest )
if err != nil {
return err
}
return nil
}); err != nil {
2024-06-29 14:32:57 -04:00
return errors . New ( "failed to dial to " + fb . Dest ). Base ( err ). AtWarning ()
2020-11-25 19:01:53 +08:00
}
defer conn . Close ()
serverReader := buf . NewReader ( conn )
serverWriter := buf . NewWriter ( conn )
postRequest := func () error {
defer timer . SetTimeout ( sessionPolicy . Timeouts . DownlinkOnly )
if fb . Xver != 0 {
2021-01-23 21:06:15 +00:00
ipType := 4
remoteAddr , remotePort , err := net . SplitHostPort ( connection . RemoteAddr (). String ())
if err != nil {
ipType = 0
}
localAddr , localPort , err := net . SplitHostPort ( connection . LocalAddr (). String ())
if err != nil {
ipType = 0
}
if ipType == 4 {
2021-01-22 11:26:57 +08:00
for i := 0 ; i < len ( remoteAddr ); i ++ {
if remoteAddr [ i ] == ':' {
ipType = 6
break
}
2020-11-25 19:01:53 +08:00
}
}
pro := buf . New ()
defer pro . Release ()
switch fb . Xver {
case 1 :
2021-01-22 11:26:57 +08:00
if ipType == 0 {
pro . Write ([] byte ( "PROXY UNKNOWN\r\n" ))
break
}
if ipType == 4 {
2020-11-25 19:01:53 +08:00
pro . Write ([] byte ( "PROXY TCP4 " + remoteAddr + " " + localAddr + " " + remotePort + " " + localPort + "\r\n" ))
} else {
pro . Write ([] byte ( "PROXY TCP6 " + remoteAddr + " " + localAddr + " " + remotePort + " " + localPort + "\r\n" ))
}
case 2 :
2021-01-22 11:26:57 +08:00
pro . Write ([] byte ( "\x0D\x0A\x0D\x0A\x00\x0D\x0A\x51\x55\x49\x54\x0A" )) // signature
if ipType == 0 {
pro . Write ([] byte ( "\x20\x00\x00\x00" )) // v2 + LOCAL + UNSPEC + UNSPEC + 0 bytes
break
}
if ipType == 4 {
pro . Write ([] byte ( "\x21\x11\x00\x0C" )) // v2 + PROXY + AF_INET + STREAM + 12 bytes
2020-11-25 19:01:53 +08:00
pro . Write ( net . ParseIP ( remoteAddr ). To4 ())
pro . Write ( net . ParseIP ( localAddr ). To4 ())
} else {
2021-01-22 11:26:57 +08:00
pro . Write ([] byte ( "\x21\x21\x00\x24" )) // v2 + PROXY + AF_INET6 + STREAM + 36 bytes
2020-11-25 19:01:53 +08:00
pro . Write ( net . ParseIP ( remoteAddr ). To16 ())
pro . Write ( net . ParseIP ( localAddr ). To16 ())
}
p1 , _ := strconv . ParseUint ( remotePort , 10 , 16 )
p2 , _ := strconv . ParseUint ( localPort , 10 , 16 )
pro . Write ([] byte { byte ( p1 >> 8 ), byte ( p1 ), byte ( p2 >> 8 ), byte ( p2 )})
}
if err := serverWriter . WriteMultiBuffer ( buf . MultiBuffer { pro }); err != nil {
2024-06-29 14:32:57 -04:00
return errors . New ( "failed to set PROXY protocol v" , fb . Xver ). Base ( err ). AtWarning ()
2020-11-25 19:01:53 +08:00
}
}
if err := buf . Copy ( reader , serverWriter , buf . UpdateActivity ( timer )); err != nil {
2024-06-29 14:32:57 -04:00
return errors . New ( "failed to fallback request payload" ). Base ( err ). AtInfo ()
2020-11-25 19:01:53 +08:00
}
return nil
}
writer := buf . NewWriter ( connection )
getResponse := func () error {
defer timer . SetTimeout ( sessionPolicy . Timeouts . UplinkOnly )
if err := buf . Copy ( serverReader , writer , buf . UpdateActivity ( timer )); err != nil {
2024-06-29 14:32:57 -04:00
return errors . New ( "failed to deliver response payload" ). Base ( err ). AtInfo ()
2020-11-25 19:01:53 +08:00
}
return nil
}
if err := task . Run ( ctx , task . OnSuccess ( postRequest , task . Close ( serverWriter )), task . OnSuccess ( getResponse , task . Close ( writer ))); err != nil {
common . Interrupt ( serverReader )
common . Interrupt ( serverWriter )
2024-06-29 14:32:57 -04:00
return errors . New ( "fallback ends" ). Base ( err ). AtInfo ()
2020-11-25 19:01:53 +08:00
}
return nil
}
if errors . Cause ( err ) != io . EOF {
log . Record ( & log . AccessMessage {
From : connection . RemoteAddr (),
To : "" ,
Status : log . AccessRejected ,
Reason : err ,
})
2024-06-29 14:32:57 -04:00
err = errors . New ( "invalid request from " , connection . RemoteAddr ()). Base ( err ). AtInfo ()
2020-11-25 19:01:53 +08:00
}
return err
}
if err := connection . SetReadDeadline ( time . Time {}); err != nil {
2024-06-29 14:32:57 -04:00
errors . LogWarningInner ( ctx , err , "unable to set back read deadline" )
2020-11-25 19:01:53 +08:00
}
2024-06-29 14:32:57 -04:00
errors . LogInfo ( ctx , "received request for " , request . Destination ())
2020-11-25 19:01:53 +08:00
inbound := session . InboundFromContext ( ctx )
if inbound == nil {
panic ( "no inbound metadata" )
}
2023-04-06 10:21:35 +00:00
inbound . Name = "vless"
2020-11-25 19:01:53 +08:00
inbound . User = request . User
2025-08-18 08:50:43 +00:00
inbound . VlessRoute = net . PortFromBytes ( userSentID [ 6 : 8 ])
2020-11-25 19:01:53 +08:00
account := request . User . Account .( * vless . MemoryAccount )
2025-11-23 04:23:48 +00:00
if account . Reverse != nil && request . Command != protocol . RequestCommandRvs {
return errors . New ( "for safety reasons, user " + account . ID . String () + " is not allowed to use forward proxy" )
}
2020-11-25 19:01:53 +08:00
responseAddons := & encoding . Addons {
// Flow: requestAddons.Flow,
}
2023-01-06 05:37:16 +00:00
var input * bytes . Reader
var rawInput * bytes . Buffer
2020-11-25 19:01:53 +08:00
switch requestAddons . Flow {
2023-03-04 05:39:26 -05:00
case vless . XRV :
2023-03-04 15:39:27 +00:00
if account . Flow == requestAddons . Flow {
2024-05-13 21:52:24 -04:00
inbound . CanSpliceCopy = 2
2020-11-25 19:01:53 +08:00
switch request . Command {
case protocol . RequestCommandUDP :
2024-06-29 14:32:57 -04:00
return errors . New ( requestAddons . Flow + " doesn't support UDP" ). AtWarning ()
2025-09-09 14:19:12 +00:00
case protocol . RequestCommandMux , protocol . RequestCommandRvs :
inbound . CanSpliceCopy = 3
2023-04-16 21:15:36 +00:00
fallthrough // we will break Mux connections that contain TCP requests
2020-11-25 19:01:53 +08:00
case protocol . RequestCommandTCP :
2023-03-04 05:39:26 -05:00
var t reflect . Type
var p uintptr
2025-08-31 04:09:28 +00:00
if commonConn , ok := connection .( * encryption . CommonConn ); ok {
2025-09-02 18:15:08 +00:00
if _ , ok := commonConn . Conn .( * encryption . XorConn ); ok || ! proxy . IsRAWTransportWithoutSecurity ( iConn ) {
inbound . CanSpliceCopy = 3 // full-random xorConn / non-RAW transport / another securityConn should not be penetrated
2025-08-31 04:09:28 +00:00
}
t = reflect . TypeOf ( commonConn ). Elem ()
p = uintptr ( unsafe . Pointer ( commonConn ))
} else if tlsConn , ok := iConn .( * tls . Conn ); ok {
2023-03-04 05:39:26 -05:00
if tlsConn . ConnectionState (). Version != gotls . VersionTLS13 {
2024-06-29 14:32:57 -04:00
return errors . New ( `failed to use ` + requestAddons . Flow + `, found outer tls version ` , tlsConn . ConnectionState (). Version ). AtWarning ()
2020-11-25 19:01:53 +08:00
}
2023-03-04 05:39:26 -05:00
t = reflect . TypeOf ( tlsConn . Conn ). Elem ()
p = uintptr ( unsafe . Pointer ( tlsConn . Conn ))
} else if realityConn , ok := iConn .( * reality . Conn ); ok {
t = reflect . TypeOf ( realityConn . Conn ). Elem ()
p = uintptr ( unsafe . Pointer ( realityConn . Conn ))
2020-11-25 19:01:53 +08:00
} else {
2024-06-29 14:32:57 -04:00
return errors . New ( "XTLS only supports TLS and REALITY directly for now." ). AtWarning ()
2023-03-04 05:39:26 -05:00
}
i , _ := t . FieldByName ( "input" )
r , _ := t . FieldByName ( "rawInput" )
input = ( * bytes . Reader )( unsafe . Pointer ( p + i . Offset ))
rawInput = ( * bytes . Buffer )( unsafe . Pointer ( p + r . Offset ))
2020-11-25 19:01:53 +08:00
}
} else {
2024-11-25 17:16:29 +01:00
return errors . New ( "account " + account . ID . String () + " is not able to use the flow " + requestAddons . Flow ). AtWarning ()
2020-11-25 19:01:53 +08:00
}
2023-03-04 15:39:27 +00:00
case "" :
2024-05-13 21:52:24 -04:00
inbound . CanSpliceCopy = 3
2023-04-16 21:15:36 +00:00
if account . Flow == vless . XRV && ( request . Command == protocol . RequestCommandTCP || isMuxAndNotXUDP ( request , first )) {
2024-11-25 17:16:29 +01:00
return errors . New ( "account " + account . ID . String () + " is rejected since the client flow is empty. Note that the pure TLS proxy has certain TLS in TLS characters." ). AtWarning ()
2022-12-04 18:24:46 -05:00
}
2020-11-25 19:01:53 +08:00
default :
2024-06-29 14:32:57 -04:00
return errors . New ( "unknown request flow " + requestAddons . Flow ). AtWarning ()
2020-11-25 19:01:53 +08:00
}
if request . Command != protocol . RequestCommandMux {
ctx = log . ContextWithAccessMessage ( ctx , & log . AccessMessage {
From : connection . RemoteAddr (),
To : request . Destination (),
Status : log . AccessAccepted ,
Reason : "" ,
Email : request . User . Email ,
})
2023-04-12 23:20:38 +08:00
} else if account . Flow == vless . XRV {
ctx = session . ContextWithAllowedNetwork ( ctx , net . Network_UDP )
2020-11-25 19:01:53 +08:00
}
2025-08-17 18:13:56 +00:00
trafficState := proxy . NewTrafficState ( userSentID )
2025-09-01 11:15:32 -04:00
clientReader := encoding . DecodeBodyAddons ( reader , request , requestAddons )
if requestAddons . Flow == vless . XRV {
clientReader = proxy . NewVisionReader ( clientReader , trafficState , true , ctx , connection , input , rawInput , nil )
2020-11-25 19:01:53 +08:00
}
2025-09-01 11:15:32 -04:00
bufferWriter := buf . NewBufferedWriter ( buf . NewWriter ( connection ))
if err := encoding . EncodeResponseHeader ( bufferWriter , request , responseAddons ); err != nil {
return errors . New ( "failed to encode response header" ). Base ( err ). AtWarning ()
2020-11-25 19:01:53 +08:00
}
2025-09-01 11:15:32 -04:00
clientWriter := encoding . EncodeBodyAddons ( bufferWriter , request , requestAddons , trafficState , false , ctx , connection , nil )
bufferWriter . SetFlushNext ()
2025-09-09 14:19:12 +00:00
if request . Command == protocol . RequestCommandRvs {
r , err := h . GetReverse ( account )
if err != nil {
return err
}
2026-03-21 19:16:24 +08:00
return r . NewMux ( ctx , dispatcher . WrapLink ( ctx , h . policyManager , h . stats , & transport . Link { Reader : clientReader , Writer : clientWriter }), h . observer )
2025-09-09 14:19:12 +00:00
}
2025-12-31 06:00:45 -05:00
if err := dispatch . DispatchLink ( ctx , request . Destination (), & transport . Link {
2025-09-03 23:25:17 +00:00
Reader : clientReader ,
2025-09-01 11:15:32 -04:00
Writer : clientWriter },
); err != nil {
return errors . New ( "failed to dispatch request" ). Base ( err )
2020-11-25 19:01:53 +08:00
}
return nil
}
2025-09-09 14:19:12 +00:00
type Reverse struct {
tag string
picker * reverse . StaticMuxPicker
client * mux . ClientManager
}
func ( r * Reverse ) Tag () string {
return r . tag
}
2026-03-21 19:16:24 +08:00
func ( r * Reverse ) NewMux ( ctx context . Context , link * transport . Link , observer features . Feature ) error {
2025-09-09 14:19:12 +00:00
muxClient , err := mux . NewClientWorker ( * link , mux . ClientStrategy {})
if err != nil {
return errors . New ( "failed to create mux client worker" ). Base ( err ). AtWarning ()
}
worker , err := reverse . NewPortalWorker ( muxClient )
if err != nil {
return errors . New ( "failed to create portal worker" ). Base ( err ). AtWarning ()
}
r . picker . AddWorker ( worker )
2026-03-21 19:16:24 +08:00
if burstObs , ok := observer .( extension . BurstObservatory ); ok {
go burstObs . Check ([] string { r . Tag ()})
}
2025-09-09 14:19:12 +00:00
select {
case <- ctx . Done ():
case <- muxClient . WaitClosed ():
}
return nil
}
func ( r * Reverse ) Dispatch ( ctx context . Context , link * transport . Link ) {
outbounds := session . OutboundsFromContext ( ctx )
ob := outbounds [ len ( outbounds ) - 1 ]
if ob != nil {
if ob . Target . Network == net . Network_UDP && ob . OriginalTarget . Address != nil && ob . OriginalTarget . Address != ob . Target . Address {
link . Reader = & buf . EndpointOverrideReader { Reader : link . Reader , Dest : ob . Target . Address , OriginalDest : ob . OriginalTarget . Address }
link . Writer = & buf . EndpointOverrideWriter { Writer : link . Writer , Dest : ob . Target . Address , OriginalDest : ob . OriginalTarget . Address }
}
2025-10-15 07:05:52 +00:00
r . client . Dispatch ( session . ContextWithIsReverseMux ( ctx , true ), link )
2025-09-09 14:19:12 +00:00
}
}
func ( r * Reverse ) Start () error {
return nil
}
func ( r * Reverse ) Close () error {
return nil
}
func ( r * Reverse ) SenderSettings () * serial . TypedMessage {
return nil
}
func ( r * Reverse ) ProxySettings () * serial . TypedMessage {
return nil
}