Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 14 additions & 4 deletions app/proxyman/config.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions app/proxyman/config.proto
Original file line number Diff line number Diff line change
Expand Up @@ -68,4 +68,6 @@ message MultiplexingConfig {
int32 xudpConcurrency = 3;
// "reject" (default), "allow" or "skip".
string xudpProxyUDP443 = 4;
// MaxReuseTimes for an connection
int32 maxReuseTimes = 5;
}
10 changes: 8 additions & 2 deletions app/proxyman/outbound/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,12 @@ func NewHandler(ctx context.Context, config *core.OutboundHandlerConfig) (outbou

if h.senderSettings != nil && h.senderSettings.MultiplexSettings != nil {
if config := h.senderSettings.MultiplexSettings; config.Enabled {
// MaxReuseTimes use 60000 as default, and it also means the upper limit of MaxReuseTimes
// In mux cool spec, connection ID is 2 bytes, so physical limit is 65535, bu we reserve some IDs for future use
MaxReuseTimes := uint32(60000)
if config.MaxReuseTimes != 0 && config.MaxReuseTimes < 60000 {
MaxReuseTimes = uint32(config.MaxReuseTimes)
}
if config.Concurrency < 0 {
h.mux = &mux.ClientManager{Enabled: false}
}
Expand All @@ -136,7 +142,7 @@ func NewHandler(ctx context.Context, config *core.OutboundHandlerConfig) (outbou
Dialer: h,
Strategy: mux.ClientStrategy{
MaxConcurrency: uint32(config.Concurrency),
MaxConnection: 128,
MaxReuseTimes: MaxReuseTimes,
},
},
},
Expand All @@ -157,7 +163,7 @@ func NewHandler(ctx context.Context, config *core.OutboundHandlerConfig) (outbou
Dialer: h,
Strategy: mux.ClientStrategy{
MaxConcurrency: uint32(config.XudpConcurrency),
MaxConnection: 128,
MaxReuseTimes: 128,
},
},
},
Expand Down
8 changes: 8 additions & 0 deletions common/buf/copy.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,10 @@ type readError struct {
error
}

func NewReadError(err error) error {
return readError{err}
}

func (e readError) Error() string {
return e.error.Error()
}
Expand All @@ -74,6 +78,10 @@ type writeError struct {
error
}

func NewWriteError(err error) error {
return writeError{err}
}

func (e writeError) Error() string {
return e.error.Error()
}
Expand Down
70 changes: 70 additions & 0 deletions common/mux/bench_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
package mux_test

import (
"context"
"testing"

"github.com/xtls/xray-core/common"
"github.com/xtls/xray-core/common/buf"
"github.com/xtls/xray-core/common/mux"
"github.com/xtls/xray-core/common/net"
"github.com/xtls/xray-core/common/session"
"github.com/xtls/xray-core/transport"
"github.com/xtls/xray-core/transport/pipe"
)

func BenchmarkMuxThroughput(b *testing.B) {
serverCtx := session.ContextWithOutbounds(context.Background(), []*session.Outbound{{}})
muxServerUplink, muxServerDownlink := newLinkPair()
dispatcher := TestDispatcher{
OnDispatch: func(ctx context.Context, dest net.Destination) (*transport.Link, error) {
inputReader, inputWriter := pipe.New(pipe.WithSizeLimit(512 * 1024))
outputReader, outputWriter := pipe.New(pipe.WithSizeLimit(512 * 1024))
go func() {
defer outputWriter.Close()
for {
mb, err := inputReader.ReadMultiBuffer()
if err != nil {
break
}
buf.ReleaseMulti(mb)
}
}()
return &transport.Link{
Reader: outputReader,
Writer: inputWriter,
}, nil
},
}
_, err := mux.NewServerWorker(serverCtx, &dispatcher, muxServerUplink)
common.Must(err)
client, err := mux.NewClientWorker(*muxServerDownlink, mux.ClientStrategy{})
common.Must(err)
clientCtx := session.ContextWithOutbounds(context.Background(), []*session.Outbound{{
Target: net.TCPDestination(net.DomainAddress("www.example.com"), 80),
}})
muxClientUplink, muxClientDownlink := newLinkPair()
go func() {
for {
mb, err := muxClientDownlink.Reader.ReadMultiBuffer()
if err != nil {
break
}
buf.ReleaseMulti(mb)
}
}()
ok := client.Dispatch(clientCtx, muxClientUplink)
if !ok {
b.Fatal("failed to dispatch")
}
data := buf.FromBytes(make([]byte, 8192))
b.SetBytes(int64(8192))
b.ResetTimer()

for i := 0; i < b.N; i++ {
err := muxClientUplink.Writer.WriteMultiBuffer(buf.MultiBuffer{data})
if err != nil {
b.Fatal(err)
}
}
}
27 changes: 21 additions & 6 deletions common/mux/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -170,7 +170,7 @@ func (f *DialingWorkerFactory) Create() (*ClientWorker, error) {

type ClientStrategy struct {
MaxConcurrency uint32
MaxConnection uint32
MaxReuseTimes uint32
}

type ClientWorker struct {
Expand All @@ -179,6 +179,7 @@ type ClientWorker struct {
done *done.Instance
timer *time.Ticker
strategy ClientStrategy
timeCretaed time.Time
}

var (
Expand All @@ -194,6 +195,7 @@ func NewClientWorker(stream transport.Link, s ClientStrategy) (*ClientWorker, er
done: done.New(),
timer: time.NewTicker(time.Second * 16),
strategy: s,
timeCretaed: time.Now(),
}

go c.fetchOutput()
Expand Down Expand Up @@ -288,7 +290,7 @@ func fetchInput(ctx context.Context, s *Session, output buf.Writer) {

func (m *ClientWorker) IsClosing() bool {
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 false
Expand Down Expand Up @@ -318,6 +320,7 @@ func (m *ClientWorker) Dispatch(ctx context.Context, link *transport.Link) bool
if s == nil {
return false
}
errors.LogInfo(ctx, "Allocated mux.cool sub connection ID: ", s.ID, "/", m.strategy.MaxReuseTimes, " living: ", m.ActiveConnections(), "/", m.strategy.MaxConcurrency, " age: ", time.Since(m.timeCretaed).Truncate(time.Second))
s.input = link.Reader
s.output = link.Writer
go fetchInput(ctx, s, m.link.Writer)
Expand All @@ -332,14 +335,14 @@ func (m *ClientWorker) Dispatch(ctx context.Context, link *transport.Link) bool

func (m *ClientWorker) handleStatueKeepAlive(meta *FrameMetadata, reader *buf.BufferedReader) error {
if meta.Option.Has(OptionData) {
return buf.Copy(NewStreamReader(reader), buf.Discard)
return CopyChunk(reader, buf.Discard)
}
return nil
}

func (m *ClientWorker) handleStatusNew(meta *FrameMetadata, reader *buf.BufferedReader) error {
if meta.Option.Has(OptionData) {
return buf.Copy(NewStreamReader(reader), buf.Discard)
return CopyChunk(reader, buf.Discard)
}
return nil
}
Expand All @@ -355,7 +358,19 @@ func (m *ClientWorker) handleStatusKeep(meta *FrameMetadata, reader *buf.Buffere
closingWriter := NewResponseWriter(meta.SessionID, m.link.Writer, protocol.TransferTypeStream)
closingWriter.Close()

return buf.Copy(NewStreamReader(reader), buf.Discard)
return CopyChunk(reader, buf.Discard)
}

if s.transferType == protocol.TransferTypeStream {
err := CopyChunk(reader, s.output)
if err != nil && buf.IsWriteError(err) {
errors.LogInfoInner(context.Background(), err, "failed to write to downstream. closing session ", s.ID)
s.Close(false)
// down stream can have a write err but don't return the err to terminate the whole mux connection
// because it's still available for other sessions
return nil
}
return err
}

rr := s.NewReader(reader, &meta.Target)
Expand All @@ -374,7 +389,7 @@ func (m *ClientWorker) handleStatusEnd(meta *FrameMetadata, reader *buf.Buffered
s.Close(false)
}
if meta.Option.Has(OptionData) {
return buf.Copy(NewStreamReader(reader), buf.Discard)
return CopyChunk(reader, buf.Discard)
}
return nil
}
Expand Down
4 changes: 2 additions & 2 deletions common/mux/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ func TestClientWorkerClose(t *testing.T) {
Writer: w1,
}, mux.ClientStrategy{
MaxConcurrency: 4,
MaxConnection: 4,
MaxReuseTimes: 4,
})
common.Must(err)

Expand All @@ -68,7 +68,7 @@ func TestClientWorkerClose(t *testing.T) {
Writer: w2,
}, mux.ClientStrategy{
MaxConcurrency: 4,
MaxConnection: 4,
MaxReuseTimes: 4,
})
common.Must(err)

Expand Down
29 changes: 29 additions & 0 deletions common/mux/reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,3 +57,32 @@ func (r *PacketReader) ReadMultiBuffer() (buf.MultiBuffer, error) {
func NewStreamReader(reader *buf.BufferedReader) buf.Reader {
return crypto.NewChunkStreamReaderWithChunkCount(crypto.PlainChunkSizeParser{}, reader, 1)
}

func CopyChunk(reader *buf.BufferedReader, writer buf.Writer) error {
size, err := serial.ReadUint16(reader)
if err != nil {
return err
}
var writeErr error
for size > 0 {
mb, readErr := reader.ReadAtMost(int32(size))
if !mb.IsEmpty() {
size -= uint16(mb.Len())
if writeErr == nil {
if err := writer.WriteMultiBuffer(mb); err != nil {
writeErr = err
}
} else {
buf.ReleaseMulti(mb)
}
continue
}
if readErr != nil {
return buf.NewReadError(readErr)
}
}
if writeErr != nil {
return buf.NewWriteError(writeErr)
}
return nil
}
29 changes: 25 additions & 4 deletions common/mux/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -157,7 +157,7 @@ func (w *ServerWorker) Close() error {

func (w *ServerWorker) handleStatusKeepAlive(meta *FrameMetadata, reader *buf.BufferedReader) error {
if meta.Option.Has(OptionData) {
return buf.Copy(NewStreamReader(reader), buf.Discard)
return CopyChunk(reader, buf.Discard)
}
return nil
}
Expand Down Expand Up @@ -264,7 +264,7 @@ func (w *ServerWorker) handleStatusNew(ctx context.Context, meta *FrameMetadata,
link, err := w.dispatcher.Dispatch(ctx, meta.Target)
if err != nil {
if meta.Option.Has(OptionData) {
buf.Copy(NewStreamReader(reader), buf.Discard)
CopyChunk(reader, buf.Discard)
}
return errors.New("failed to dispatch request.").Base(err)
}
Expand All @@ -287,6 +287,15 @@ func (w *ServerWorker) handleStatusNew(ctx context.Context, meta *FrameMetadata,
return nil
}

if s.transferType == protocol.TransferTypeStream {
err = CopyChunk(reader, s.output)
if err != nil && buf.IsWriteError(err) {
s.Close(false)
return err
}
return err
}

rr := s.NewReader(reader, &meta.Target)
err = buf.Copy(rr, s.output)

Expand All @@ -308,7 +317,19 @@ func (w *ServerWorker) handleStatusKeep(meta *FrameMetadata, reader *buf.Buffere
closingWriter := NewResponseWriter(meta.SessionID, w.link.Writer, protocol.TransferTypeStream)
closingWriter.Close()

return buf.Copy(NewStreamReader(reader), buf.Discard)
return CopyChunk(reader, buf.Discard)
}

if s.transferType == protocol.TransferTypeStream {
err := CopyChunk(reader, s.output)
if err != nil && buf.IsWriteError(err) {
errors.LogInfoInner(context.Background(), err, "failed to write to downstream writer. closing session ", s.ID)
s.Close(false)
// down stream can have a write err but don't return the err to terminate the whole mux connection
// because it's still available for other sessions
return nil
}
return err
}

rr := s.NewReader(reader, &meta.Target)
Expand All @@ -328,7 +349,7 @@ func (w *ServerWorker) handleStatusEnd(meta *FrameMetadata, reader *buf.Buffered
s.Close(false)
}
if meta.Option.Has(OptionData) {
return buf.Copy(NewStreamReader(reader), buf.Discard)
return CopyChunk(reader, buf.Discard)
}
return nil
}
Expand Down
2 changes: 1 addition & 1 deletion common/mux/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ import (
)

func newLinkPair() (*transport.Link, *transport.Link) {
opt := pipe.WithoutSizeLimit()
opt := pipe.WithSizeLimit(512 * 1024)
uplinkReader, uplinkWriter := pipe.New(opt)
downlinkReader, downlinkWriter := pipe.New(opt)

Expand Down
2 changes: 1 addition & 1 deletion common/mux/session.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ func (m *SessionManager) Allocate(Strategy *ClientStrategy) *Session {
defer m.Unlock()

MaxConcurrency := int(Strategy.MaxConcurrency)
MaxConnection := uint16(Strategy.MaxConnection)
MaxConnection := uint16(Strategy.MaxReuseTimes)

if m.closed || (MaxConcurrency > 0 && len(m.sessions) >= MaxConcurrency) || (MaxConnection > 0 && m.count >= MaxConnection) {
return nil
Expand Down
Loading
Loading