mirror of
https://github.com/XTLS/Xray-core.git
synced 2026-10-01 21:45:44 +00:00
Compare commits
2
Commits
log
...
mux-brutal
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b33fde2e44 | ||
|
|
65da3d9c9d |
@@ -411,6 +411,7 @@ type MultiplexingConfig struct {
|
|||||||
XudpConcurrency int32 `protobuf:"varint,3,opt,name=xudpConcurrency,proto3" json:"xudpConcurrency,omitempty"`
|
XudpConcurrency int32 `protobuf:"varint,3,opt,name=xudpConcurrency,proto3" json:"xudpConcurrency,omitempty"`
|
||||||
// "reject" (default), "allow" or "skip".
|
// "reject" (default), "allow" or "skip".
|
||||||
XudpProxyUDP443 string `protobuf:"bytes,4,opt,name=xudpProxyUDP443,proto3" json:"xudpProxyUDP443,omitempty"`
|
XudpProxyUDP443 string `protobuf:"bytes,4,opt,name=xudpProxyUDP443,proto3" json:"xudpProxyUDP443,omitempty"`
|
||||||
|
BrutalBPS uint64 `protobuf:"varint,5,opt,name=brutalBPS,proto3" json:"brutalBPS,omitempty"`
|
||||||
unknownFields protoimpl.UnknownFields
|
unknownFields protoimpl.UnknownFields
|
||||||
sizeCache protoimpl.SizeCache
|
sizeCache protoimpl.SizeCache
|
||||||
}
|
}
|
||||||
@@ -473,6 +474,13 @@ func (x *MultiplexingConfig) GetXudpProxyUDP443() string {
|
|||||||
return ""
|
return ""
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (x *MultiplexingConfig) GetBrutalBPS() uint64 {
|
||||||
|
if x != nil {
|
||||||
|
return x.BrutalBPS
|
||||||
|
}
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
|
||||||
var File_app_proxyman_config_proto protoreflect.FileDescriptor
|
var File_app_proxyman_config_proto protoreflect.FileDescriptor
|
||||||
|
|
||||||
const file_app_proxyman_config_proto_rawDesc = "" +
|
const file_app_proxyman_config_proto_rawDesc = "" +
|
||||||
@@ -503,12 +511,13 @@ const file_app_proxyman_config_proto_rawDesc = "" +
|
|||||||
"\x0eproxy_settings\x18\x03 \x01(\v2$.xray.transport.internet.ProxyConfigR\rproxySettings\x12T\n" +
|
"\x0eproxy_settings\x18\x03 \x01(\v2$.xray.transport.internet.ProxyConfigR\rproxySettings\x12T\n" +
|
||||||
"\x12multiplex_settings\x18\x04 \x01(\v2%.xray.app.proxyman.MultiplexingConfigR\x11multiplexSettings\x12\x19\n" +
|
"\x12multiplex_settings\x18\x04 \x01(\v2%.xray.app.proxyman.MultiplexingConfigR\x11multiplexSettings\x12\x19\n" +
|
||||||
"\bvia_cidr\x18\x05 \x01(\tR\aviaCidr\x12P\n" +
|
"\bvia_cidr\x18\x05 \x01(\tR\aviaCidr\x12P\n" +
|
||||||
"\x0ftarget_strategy\x18\x06 \x01(\x0e2'.xray.transport.internet.DomainStrategyR\x0etargetStrategy\"\xa4\x01\n" +
|
"\x0ftarget_strategy\x18\x06 \x01(\x0e2'.xray.transport.internet.DomainStrategyR\x0etargetStrategy\"\xc2\x01\n" +
|
||||||
"\x12MultiplexingConfig\x12\x18\n" +
|
"\x12MultiplexingConfig\x12\x18\n" +
|
||||||
"\aenabled\x18\x01 \x01(\bR\aenabled\x12 \n" +
|
"\aenabled\x18\x01 \x01(\bR\aenabled\x12 \n" +
|
||||||
"\vconcurrency\x18\x02 \x01(\x05R\vconcurrency\x12(\n" +
|
"\vconcurrency\x18\x02 \x01(\x05R\vconcurrency\x12(\n" +
|
||||||
"\x0fxudpConcurrency\x18\x03 \x01(\x05R\x0fxudpConcurrency\x12(\n" +
|
"\x0fxudpConcurrency\x18\x03 \x01(\x05R\x0fxudpConcurrency\x12(\n" +
|
||||||
"\x0fxudpProxyUDP443\x18\x04 \x01(\tR\x0fxudpProxyUDP443BU\n" +
|
"\x0fxudpProxyUDP443\x18\x04 \x01(\tR\x0fxudpProxyUDP443\x12\x1c\n" +
|
||||||
|
"\tbrutalBPS\x18\x05 \x01(\x04R\tbrutalBPSBU\n" +
|
||||||
"\x15com.xray.app.proxymanP\x01Z&github.com/xtls/xray-core/app/proxyman\xaa\x02\x11Xray.App.Proxymanb\x06proto3"
|
"\x15com.xray.app.proxymanP\x01Z&github.com/xtls/xray-core/app/proxyman\xaa\x02\x11Xray.App.Proxymanb\x06proto3"
|
||||||
|
|
||||||
var (
|
var (
|
||||||
|
|||||||
@@ -68,4 +68,5 @@ message MultiplexingConfig {
|
|||||||
int32 xudpConcurrency = 3;
|
int32 xudpConcurrency = 3;
|
||||||
// "reject" (default), "allow" or "skip".
|
// "reject" (default), "allow" or "skip".
|
||||||
string xudpProxyUDP443 = 4;
|
string xudpProxyUDP443 = 4;
|
||||||
|
uint64 brutalBPS = 5;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -136,7 +136,8 @@ func NewHandler(ctx context.Context, config *core.OutboundHandlerConfig) (outbou
|
|||||||
Dialer: h,
|
Dialer: h,
|
||||||
Strategy: mux.ClientStrategy{
|
Strategy: mux.ClientStrategy{
|
||||||
MaxConcurrency: uint32(config.Concurrency),
|
MaxConcurrency: uint32(config.Concurrency),
|
||||||
MaxConnection: 128,
|
MaxReuseTimes: 32768,
|
||||||
|
BrutalBPS: config.BrutalBPS,
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
@@ -157,7 +158,8 @@ func NewHandler(ctx context.Context, config *core.OutboundHandlerConfig) (outbou
|
|||||||
Dialer: h,
|
Dialer: h,
|
||||||
Strategy: mux.ClientStrategy{
|
Strategy: mux.ClientStrategy{
|
||||||
MaxConcurrency: uint32(config.XudpConcurrency),
|
MaxConcurrency: uint32(config.XudpConcurrency),
|
||||||
MaxConnection: 128,
|
MaxReuseTimes: 32768,
|
||||||
|
BrutalBPS: config.BrutalBPS,
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -75,6 +75,12 @@ func (mb MultiBuffer) Copy(b []byte) int {
|
|||||||
return total
|
return total
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (mb MultiBuffer) Bytes() []byte {
|
||||||
|
b := make([]byte, mb.Len())
|
||||||
|
mb.Copy(b)
|
||||||
|
return b
|
||||||
|
}
|
||||||
|
|
||||||
// ReadFrom reads all content from reader until EOF.
|
// ReadFrom reads all content from reader until EOF.
|
||||||
func ReadFrom(reader io.Reader) (MultiBuffer, error) {
|
func ReadFrom(reader io.Reader) (MultiBuffer, error) {
|
||||||
mb := make(MultiBuffer, 0, 16)
|
mb := make(MultiBuffer, 0, 16)
|
||||||
|
|||||||
+29
-2
@@ -2,6 +2,7 @@ package mux
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"encoding/binary"
|
||||||
goerrors "errors"
|
goerrors "errors"
|
||||||
"io"
|
"io"
|
||||||
"sync"
|
"sync"
|
||||||
@@ -9,6 +10,7 @@ import (
|
|||||||
|
|
||||||
"github.com/xtls/xray-core/common"
|
"github.com/xtls/xray-core/common"
|
||||||
"github.com/xtls/xray-core/common/buf"
|
"github.com/xtls/xray-core/common/buf"
|
||||||
|
"github.com/xtls/xray-core/common/dice"
|
||||||
"github.com/xtls/xray-core/common/errors"
|
"github.com/xtls/xray-core/common/errors"
|
||||||
"github.com/xtls/xray-core/common/net"
|
"github.com/xtls/xray-core/common/net"
|
||||||
"github.com/xtls/xray-core/common/protocol"
|
"github.com/xtls/xray-core/common/protocol"
|
||||||
@@ -170,7 +172,8 @@ func (f *DialingWorkerFactory) Create() (*ClientWorker, error) {
|
|||||||
|
|
||||||
type ClientStrategy struct {
|
type ClientStrategy struct {
|
||||||
MaxConcurrency uint32
|
MaxConcurrency uint32
|
||||||
MaxConnection uint32
|
MaxReuseTimes uint32
|
||||||
|
BrutalBPS uint64
|
||||||
}
|
}
|
||||||
|
|
||||||
type ClientWorker struct {
|
type ClientWorker struct {
|
||||||
@@ -198,6 +201,9 @@ func NewClientWorker(stream transport.Link, s ClientStrategy) (*ClientWorker, er
|
|||||||
|
|
||||||
go c.fetchOutput()
|
go c.fetchOutput()
|
||||||
go c.monitor()
|
go c.monitor()
|
||||||
|
if s.BrutalBPS > 0 {
|
||||||
|
go c.sendSetBrutal(s.BrutalBPS)
|
||||||
|
}
|
||||||
|
|
||||||
return c, nil
|
return c, nil
|
||||||
}
|
}
|
||||||
@@ -288,7 +294,7 @@ func fetchInput(ctx context.Context, s *Session, output buf.Writer) {
|
|||||||
|
|
||||||
func (m *ClientWorker) IsClosing() bool {
|
func (m *ClientWorker) IsClosing() bool {
|
||||||
sm := m.sessionManager
|
sm := m.sessionManager
|
||||||
if m.strategy.MaxConnection > 0 && sm.Count() >= int(m.strategy.MaxConnection) {
|
if m.strategy.MaxReuseTimes > 0 && sm.Count() >= int(m.strategy.MaxReuseTimes) {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
@@ -417,3 +423,24 @@ func (m *ClientWorker) fetchOutput() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (m *ClientWorker) sendSetBrutal(sendBPS uint64) {
|
||||||
|
meta := FrameMetadata{
|
||||||
|
SessionID: 91,
|
||||||
|
SessionStatus: SessionStatusSetBrutal,
|
||||||
|
}
|
||||||
|
meta.Option.Set(OptionData)
|
||||||
|
frame := buf.New()
|
||||||
|
common.Must(meta.WriteTo(frame))
|
||||||
|
lengthByte := frame.Extend(2)
|
||||||
|
speedByte := frame.Extend(int32(200 + dice.Roll(200)))
|
||||||
|
binary.BigEndian.PutUint64(speedByte, sendBPS)
|
||||||
|
binary.BigEndian.PutUint16(lengthByte, uint16(len(speedByte)))
|
||||||
|
errors.LogError(context.Background(), "Start sending SetBrutal frame with speed: ", sendBPS)
|
||||||
|
err := m.link.Writer.WriteMultiBuffer(buf.MultiBuffer{frame})
|
||||||
|
if err != nil {
|
||||||
|
frame.Release()
|
||||||
|
errors.LogError(context.Background(), "failed to send SetBrutal frame: ", err)
|
||||||
|
}
|
||||||
|
errors.LogInfo(context.Background(), "SetBrutal frame sent successfully")
|
||||||
|
}
|
||||||
|
|||||||
@@ -58,7 +58,7 @@ func TestClientWorkerClose(t *testing.T) {
|
|||||||
Writer: w1,
|
Writer: w1,
|
||||||
}, mux.ClientStrategy{
|
}, mux.ClientStrategy{
|
||||||
MaxConcurrency: 4,
|
MaxConcurrency: 4,
|
||||||
MaxConnection: 4,
|
MaxReuseTimes: 4,
|
||||||
})
|
})
|
||||||
common.Must(err)
|
common.Must(err)
|
||||||
|
|
||||||
@@ -68,7 +68,7 @@ func TestClientWorkerClose(t *testing.T) {
|
|||||||
Writer: w2,
|
Writer: w2,
|
||||||
}, mux.ClientStrategy{
|
}, mux.ClientStrategy{
|
||||||
MaxConcurrency: 4,
|
MaxConcurrency: 4,
|
||||||
MaxConnection: 4,
|
MaxReuseTimes: 4,
|
||||||
})
|
})
|
||||||
common.Must(err)
|
common.Must(err)
|
||||||
|
|
||||||
|
|||||||
@@ -21,6 +21,8 @@ const (
|
|||||||
SessionStatusKeep SessionStatus = 0x02
|
SessionStatusKeep SessionStatus = 0x02
|
||||||
SessionStatusEnd SessionStatus = 0x03
|
SessionStatusEnd SessionStatus = 0x03
|
||||||
SessionStatusKeepAlive SessionStatus = 0x04
|
SessionStatusKeepAlive SessionStatus = 0x04
|
||||||
|
|
||||||
|
SessionStatusSetBrutal SessionStatus = 0x91
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package mux
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"encoding/binary"
|
||||||
"io"
|
"io"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -16,6 +17,7 @@ import (
|
|||||||
"github.com/xtls/xray-core/core"
|
"github.com/xtls/xray-core/core"
|
||||||
"github.com/xtls/xray-core/features/routing"
|
"github.com/xtls/xray-core/features/routing"
|
||||||
"github.com/xtls/xray-core/transport"
|
"github.com/xtls/xray-core/transport"
|
||||||
|
"github.com/xtls/xray-core/transport/internet/brutal"
|
||||||
"github.com/xtls/xray-core/transport/pipe"
|
"github.com/xtls/xray-core/transport/pipe"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -333,6 +335,40 @@ func (w *ServerWorker) handleStatusEnd(meta *FrameMetadata, reader *buf.Buffered
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// note: do not return error, just log it
|
||||||
|
func (w *ServerWorker) HandleSetBrutal(ctx context.Context, meta *FrameMetadata, reader *buf.BufferedReader) error {
|
||||||
|
if meta.Option.Has(OptionData) == false {
|
||||||
|
errors.LogError(ctx, "SetBrutal frame missing data")
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
chunkReader := NewStreamReader(reader)
|
||||||
|
data, err := chunkReader.ReadMultiBuffer()
|
||||||
|
if err != nil {
|
||||||
|
errors.LogError(ctx, "unexpected error when reading brutal data: ", err)
|
||||||
|
}
|
||||||
|
speed := binary.BigEndian.Uint64(data.Bytes())
|
||||||
|
|
||||||
|
inbound := session.InboundFromContext(ctx)
|
||||||
|
if inbound == nil || inbound.Conn == nil {
|
||||||
|
errors.LogError(ctx, "no inbound connection found for brutal set")
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
conn := inbound.Conn
|
||||||
|
tcpConn, ok := conn.(*net.TCPConn)
|
||||||
|
if !ok {
|
||||||
|
errors.LogError(ctx, "brutal can only be set on TCP connections")
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
err = brutal.SetBrutal(tcpConn, speed)
|
||||||
|
if err != nil {
|
||||||
|
errors.LogError(ctx, "failed to set brutal: ", err)
|
||||||
|
return nil
|
||||||
|
} else {
|
||||||
|
errors.LogInfo(ctx, "successfully set brutal speed: ", speed)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func (w *ServerWorker) handleFrame(ctx context.Context, reader *buf.BufferedReader) error {
|
func (w *ServerWorker) handleFrame(ctx context.Context, reader *buf.BufferedReader) error {
|
||||||
var meta FrameMetadata
|
var meta FrameMetadata
|
||||||
err := meta.Unmarshal(reader, session.IsReverseMuxFromContext(ctx))
|
err := meta.Unmarshal(reader, session.IsReverseMuxFromContext(ctx))
|
||||||
@@ -349,6 +385,8 @@ func (w *ServerWorker) handleFrame(ctx context.Context, reader *buf.BufferedRead
|
|||||||
err = w.handleStatusNew(session.ContextWithIsReverseMux(ctx, false), &meta, reader)
|
err = w.handleStatusNew(session.ContextWithIsReverseMux(ctx, false), &meta, reader)
|
||||||
case SessionStatusKeep:
|
case SessionStatusKeep:
|
||||||
err = w.handleStatusKeep(&meta, reader)
|
err = w.handleStatusKeep(&meta, reader)
|
||||||
|
case SessionStatusSetBrutal:
|
||||||
|
err = w.HandleSetBrutal(ctx, &meta, reader)
|
||||||
default:
|
default:
|
||||||
status := meta.SessionStatus
|
status := meta.SessionStatus
|
||||||
return errors.New("unknown status: ", status).AtError()
|
return errors.New("unknown status: ", status).AtError()
|
||||||
|
|||||||
@@ -56,11 +56,12 @@ func (m *SessionManager) Allocate(Strategy *ClientStrategy) *Session {
|
|||||||
defer m.Unlock()
|
defer m.Unlock()
|
||||||
|
|
||||||
MaxConcurrency := int(Strategy.MaxConcurrency)
|
MaxConcurrency := int(Strategy.MaxConcurrency)
|
||||||
MaxConnection := uint16(Strategy.MaxConnection)
|
MaxReuseTimes := uint16(Strategy.MaxReuseTimes)
|
||||||
|
|
||||||
if m.closed || (MaxConcurrency > 0 && len(m.sessions) >= MaxConcurrency) || (MaxConnection > 0 && m.count >= MaxConnection) {
|
if m.closed || (MaxConcurrency > 0 && len(m.sessions) >= MaxConcurrency) || (MaxReuseTimes > 0 && m.count >= MaxReuseTimes) {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
errors.LogInfo(context.Background(), "Allocated mux.cool session id ", m.count, "/", MaxReuseTimes)
|
||||||
|
|
||||||
m.count++
|
m.count++
|
||||||
s := &Session{
|
s := &Session{
|
||||||
|
|||||||
@@ -102,6 +102,7 @@ func (c *SniffingConfig) Build() (*proxyman.SniffingConfig, error) {
|
|||||||
type MuxConfig struct {
|
type MuxConfig struct {
|
||||||
Enabled bool `json:"enabled"`
|
Enabled bool `json:"enabled"`
|
||||||
Concurrency int16 `json:"concurrency"`
|
Concurrency int16 `json:"concurrency"`
|
||||||
|
BrutalBPS uint64 `json:"brutalBPS"`
|
||||||
XudpConcurrency int16 `json:"xudpConcurrency"`
|
XudpConcurrency int16 `json:"xudpConcurrency"`
|
||||||
XudpProxyUDP443 string `json:"xudpProxyUDP443"`
|
XudpProxyUDP443 string `json:"xudpProxyUDP443"`
|
||||||
}
|
}
|
||||||
@@ -118,6 +119,7 @@ func (m *MuxConfig) Build() (*proxyman.MultiplexingConfig, error) {
|
|||||||
return &proxyman.MultiplexingConfig{
|
return &proxyman.MultiplexingConfig{
|
||||||
Enabled: m.Enabled,
|
Enabled: m.Enabled,
|
||||||
Concurrency: int32(m.Concurrency),
|
Concurrency: int32(m.Concurrency),
|
||||||
|
BrutalBPS: m.BrutalBPS,
|
||||||
XudpConcurrency: int32(m.XudpConcurrency),
|
XudpConcurrency: int32(m.XudpConcurrency),
|
||||||
XudpProxyUDP443: m.XudpProxyUDP443,
|
XudpProxyUDP443: m.XudpProxyUDP443,
|
||||||
}, nil
|
}, nil
|
||||||
|
|||||||
@@ -0,0 +1,53 @@
|
|||||||
|
//go:build linux
|
||||||
|
|
||||||
|
package brutal
|
||||||
|
|
||||||
|
import (
|
||||||
|
"os"
|
||||||
|
"unsafe"
|
||||||
|
|
||||||
|
"github.com/xtls/xray-core/common/net"
|
||||||
|
"golang.org/x/sys/unix"
|
||||||
|
)
|
||||||
|
|
||||||
|
//go:linkname setsockopt syscall.setsockopt
|
||||||
|
func setsockopt(s int, level int, name int, val unsafe.Pointer, vallen uintptr) (err error)
|
||||||
|
|
||||||
|
const (
|
||||||
|
TCP_BRUTAL_PARAMS = 23301
|
||||||
|
)
|
||||||
|
|
||||||
|
type TCPBrutalParams struct {
|
||||||
|
Rate uint64
|
||||||
|
CwndGain uint32
|
||||||
|
}
|
||||||
|
|
||||||
|
func setBrutalFD(fd uintptr, sendBPS uint64) error {
|
||||||
|
err := unix.SetsockoptString(int(fd), unix.IPPROTO_TCP, unix.TCP_CONGESTION, "brutal")
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
params := TCPBrutalParams{
|
||||||
|
Rate: sendBPS,
|
||||||
|
CwndGain: 20, // hysteria2 default
|
||||||
|
}
|
||||||
|
err = setsockopt(int(fd), unix.IPPROTO_TCP, TCP_BRUTAL_PARAMS, unsafe.Pointer(¶ms), unsafe.Sizeof(params))
|
||||||
|
if err != nil {
|
||||||
|
return os.NewSyscallError("setsockopt IPPROTO_TCP TCP_BRUTAL_PARAMS", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func SetBrutal(conn *net.TCPConn, sendBPS uint64) error {
|
||||||
|
syscallConn, err := conn.SyscallConn()
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
syscallConn.Control(func(fd uintptr) {
|
||||||
|
err = setBrutalFD(fd, sendBPS)
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,12 @@
|
|||||||
|
//go:build !linux
|
||||||
|
|
||||||
|
package brutal
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/xtls/xray-core/common/errors"
|
||||||
|
"github.com/xtls/xray-core/common/net"
|
||||||
|
)
|
||||||
|
|
||||||
|
func SetBrutal(conn *net.TCPConn, sendBPS uint64) error {
|
||||||
|
return errors.New("brutal not available on this platform")
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user