mirror of
https://github.com/XTLS/Xray-core.git
synced 2026-09-16 22:40:27 +00:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a0bf239fc1 | ||
|
|
dee64ef240 | ||
|
|
0a1b5bfb51 |
@@ -411,8 +411,10 @@ 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"`
|
||||||
unknownFields protoimpl.UnknownFields
|
// MaxReuseTimes for an connection
|
||||||
sizeCache protoimpl.SizeCache
|
MaxReuseTimes int32 `protobuf:"varint,5,opt,name=maxReuseTimes,proto3" json:"maxReuseTimes,omitempty"`
|
||||||
|
unknownFields protoimpl.UnknownFields
|
||||||
|
sizeCache protoimpl.SizeCache
|
||||||
}
|
}
|
||||||
|
|
||||||
func (x *MultiplexingConfig) Reset() {
|
func (x *MultiplexingConfig) Reset() {
|
||||||
@@ -473,6 +475,13 @@ func (x *MultiplexingConfig) GetXudpProxyUDP443() string {
|
|||||||
return ""
|
return ""
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (x *MultiplexingConfig) GetMaxReuseTimes() int32 {
|
||||||
|
if x != nil {
|
||||||
|
return x.MaxReuseTimes
|
||||||
|
}
|
||||||
|
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 +512,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\"\xca\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$\n" +
|
||||||
|
"\rmaxReuseTimes\x18\x05 \x01(\x05R\rmaxReuseTimesBU\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,6 @@ message MultiplexingConfig {
|
|||||||
int32 xudpConcurrency = 3;
|
int32 xudpConcurrency = 3;
|
||||||
// "reject" (default), "allow" or "skip".
|
// "reject" (default), "allow" or "skip".
|
||||||
string xudpProxyUDP443 = 4;
|
string xudpProxyUDP443 = 4;
|
||||||
|
// MaxReuseTimes for an connection
|
||||||
|
int32 maxReuseTimes = 5;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -121,6 +121,12 @@ func NewHandler(ctx context.Context, config *core.OutboundHandlerConfig) (outbou
|
|||||||
|
|
||||||
if h.senderSettings != nil && h.senderSettings.MultiplexSettings != nil {
|
if h.senderSettings != nil && h.senderSettings.MultiplexSettings != nil {
|
||||||
if config := h.senderSettings.MultiplexSettings; config.Enabled {
|
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 {
|
if config.Concurrency < 0 {
|
||||||
h.mux = &mux.ClientManager{Enabled: false}
|
h.mux = &mux.ClientManager{Enabled: false}
|
||||||
}
|
}
|
||||||
@@ -136,7 +142,7 @@ 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: MaxReuseTimes,
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
@@ -157,7 +163,7 @@ 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: 128,
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -56,6 +56,10 @@ type readError struct {
|
|||||||
error
|
error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func NewReadError(err error) error {
|
||||||
|
return readError{err}
|
||||||
|
}
|
||||||
|
|
||||||
func (e readError) Error() string {
|
func (e readError) Error() string {
|
||||||
return e.error.Error()
|
return e.error.Error()
|
||||||
}
|
}
|
||||||
@@ -74,6 +78,10 @@ type writeError struct {
|
|||||||
error
|
error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func NewWriteError(err error) error {
|
||||||
|
return writeError{err}
|
||||||
|
}
|
||||||
|
|
||||||
func (e writeError) Error() string {
|
func (e writeError) Error() string {
|
||||||
return e.error.Error()
|
return e.error.Error()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
+21
-6
@@ -170,7 +170,7 @@ func (f *DialingWorkerFactory) Create() (*ClientWorker, error) {
|
|||||||
|
|
||||||
type ClientStrategy struct {
|
type ClientStrategy struct {
|
||||||
MaxConcurrency uint32
|
MaxConcurrency uint32
|
||||||
MaxConnection uint32
|
MaxReuseTimes uint32
|
||||||
}
|
}
|
||||||
|
|
||||||
type ClientWorker struct {
|
type ClientWorker struct {
|
||||||
@@ -179,6 +179,7 @@ type ClientWorker struct {
|
|||||||
done *done.Instance
|
done *done.Instance
|
||||||
timer *time.Ticker
|
timer *time.Ticker
|
||||||
strategy ClientStrategy
|
strategy ClientStrategy
|
||||||
|
timeCretaed time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@@ -194,6 +195,7 @@ func NewClientWorker(stream transport.Link, s ClientStrategy) (*ClientWorker, er
|
|||||||
done: done.New(),
|
done: done.New(),
|
||||||
timer: time.NewTicker(time.Second * 16),
|
timer: time.NewTicker(time.Second * 16),
|
||||||
strategy: s,
|
strategy: s,
|
||||||
|
timeCretaed: time.Now(),
|
||||||
}
|
}
|
||||||
|
|
||||||
go c.fetchOutput()
|
go c.fetchOutput()
|
||||||
@@ -288,7 +290,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
|
||||||
@@ -318,6 +320,7 @@ func (m *ClientWorker) Dispatch(ctx context.Context, link *transport.Link) bool
|
|||||||
if s == nil {
|
if s == nil {
|
||||||
return false
|
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.input = link.Reader
|
||||||
s.output = link.Writer
|
s.output = link.Writer
|
||||||
go fetchInput(ctx, s, m.link.Writer)
|
go fetchInput(ctx, s, m.link.Writer)
|
||||||
@@ -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 {
|
func (m *ClientWorker) handleStatueKeepAlive(meta *FrameMetadata, reader *buf.BufferedReader) error {
|
||||||
if meta.Option.Has(OptionData) {
|
if meta.Option.Has(OptionData) {
|
||||||
return buf.Copy(NewStreamReader(reader), buf.Discard)
|
return CopyChunk(reader, buf.Discard)
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *ClientWorker) handleStatusNew(meta *FrameMetadata, reader *buf.BufferedReader) error {
|
func (m *ClientWorker) handleStatusNew(meta *FrameMetadata, reader *buf.BufferedReader) error {
|
||||||
if meta.Option.Has(OptionData) {
|
if meta.Option.Has(OptionData) {
|
||||||
return buf.Copy(NewStreamReader(reader), buf.Discard)
|
return CopyChunk(reader, buf.Discard)
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -355,7 +358,19 @@ func (m *ClientWorker) handleStatusKeep(meta *FrameMetadata, reader *buf.Buffere
|
|||||||
closingWriter := NewResponseWriter(meta.SessionID, m.link.Writer, protocol.TransferTypeStream)
|
closingWriter := NewResponseWriter(meta.SessionID, m.link.Writer, protocol.TransferTypeStream)
|
||||||
closingWriter.Close()
|
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)
|
rr := s.NewReader(reader, &meta.Target)
|
||||||
@@ -374,7 +389,7 @@ func (m *ClientWorker) handleStatusEnd(meta *FrameMetadata, reader *buf.Buffered
|
|||||||
s.Close(false)
|
s.Close(false)
|
||||||
}
|
}
|
||||||
if meta.Option.Has(OptionData) {
|
if meta.Option.Has(OptionData) {
|
||||||
return buf.Copy(NewStreamReader(reader), buf.Discard)
|
return CopyChunk(reader, buf.Discard)
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|
||||||
|
|||||||
@@ -57,3 +57,32 @@ func (r *PacketReader) ReadMultiBuffer() (buf.MultiBuffer, error) {
|
|||||||
func NewStreamReader(reader *buf.BufferedReader) buf.Reader {
|
func NewStreamReader(reader *buf.BufferedReader) buf.Reader {
|
||||||
return crypto.NewChunkStreamReaderWithChunkCount(crypto.PlainChunkSizeParser{}, reader, 1)
|
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
|
||||||
|
}
|
||||||
|
|||||||
+25
-4
@@ -157,7 +157,7 @@ func (w *ServerWorker) Close() error {
|
|||||||
|
|
||||||
func (w *ServerWorker) handleStatusKeepAlive(meta *FrameMetadata, reader *buf.BufferedReader) error {
|
func (w *ServerWorker) handleStatusKeepAlive(meta *FrameMetadata, reader *buf.BufferedReader) error {
|
||||||
if meta.Option.Has(OptionData) {
|
if meta.Option.Has(OptionData) {
|
||||||
return buf.Copy(NewStreamReader(reader), buf.Discard)
|
return CopyChunk(reader, buf.Discard)
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -264,7 +264,7 @@ func (w *ServerWorker) handleStatusNew(ctx context.Context, meta *FrameMetadata,
|
|||||||
link, err := w.dispatcher.Dispatch(ctx, meta.Target)
|
link, err := w.dispatcher.Dispatch(ctx, meta.Target)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if meta.Option.Has(OptionData) {
|
if meta.Option.Has(OptionData) {
|
||||||
buf.Copy(NewStreamReader(reader), buf.Discard)
|
CopyChunk(reader, buf.Discard)
|
||||||
}
|
}
|
||||||
return errors.New("failed to dispatch request.").Base(err)
|
return errors.New("failed to dispatch request.").Base(err)
|
||||||
}
|
}
|
||||||
@@ -287,6 +287,15 @@ func (w *ServerWorker) handleStatusNew(ctx context.Context, meta *FrameMetadata,
|
|||||||
return nil
|
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)
|
rr := s.NewReader(reader, &meta.Target)
|
||||||
err = buf.Copy(rr, s.output)
|
err = buf.Copy(rr, s.output)
|
||||||
|
|
||||||
@@ -308,7 +317,19 @@ func (w *ServerWorker) handleStatusKeep(meta *FrameMetadata, reader *buf.Buffere
|
|||||||
closingWriter := NewResponseWriter(meta.SessionID, w.link.Writer, protocol.TransferTypeStream)
|
closingWriter := NewResponseWriter(meta.SessionID, w.link.Writer, protocol.TransferTypeStream)
|
||||||
closingWriter.Close()
|
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)
|
rr := s.NewReader(reader, &meta.Target)
|
||||||
@@ -328,7 +349,7 @@ func (w *ServerWorker) handleStatusEnd(meta *FrameMetadata, reader *buf.Buffered
|
|||||||
s.Close(false)
|
s.Close(false)
|
||||||
}
|
}
|
||||||
if meta.Option.Has(OptionData) {
|
if meta.Option.Has(OptionData) {
|
||||||
return buf.Copy(NewStreamReader(reader), buf.Discard)
|
return CopyChunk(reader, buf.Discard)
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -15,7 +15,7 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
func newLinkPair() (*transport.Link, *transport.Link) {
|
func newLinkPair() (*transport.Link, *transport.Link) {
|
||||||
opt := pipe.WithoutSizeLimit()
|
opt := pipe.WithSizeLimit(512 * 1024)
|
||||||
uplinkReader, uplinkWriter := pipe.New(opt)
|
uplinkReader, uplinkWriter := pipe.New(opt)
|
||||||
downlinkReader, downlinkWriter := pipe.New(opt)
|
downlinkReader, downlinkWriter := pipe.New(opt)
|
||||||
|
|
||||||
|
|||||||
@@ -56,7 +56,7 @@ 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)
|
MaxConnection := 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) || (MaxConnection > 0 && m.count >= MaxConnection) {
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -101,6 +101,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"`
|
||||||
|
MaxReuseTimes int32 `json:"maxReuseTimes"`
|
||||||
XudpConcurrency int16 `json:"xudpConcurrency"`
|
XudpConcurrency int16 `json:"xudpConcurrency"`
|
||||||
XudpProxyUDP443 string `json:"xudpProxyUDP443"`
|
XudpProxyUDP443 string `json:"xudpProxyUDP443"`
|
||||||
}
|
}
|
||||||
@@ -117,6 +118,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),
|
||||||
|
MaxReuseTimes: m.MaxReuseTimes,
|
||||||
XudpConcurrency: int32(m.XudpConcurrency),
|
XudpConcurrency: int32(m.XudpConcurrency),
|
||||||
XudpProxyUDP443: m.XudpProxyUDP443,
|
XudpProxyUDP443: m.XudpProxyUDP443,
|
||||||
}, nil
|
}, nil
|
||||||
|
|||||||
Reference in New Issue
Block a user