diff --git a/proxy/direct.go b/proxy/direct.go new file mode 100644 index 0000000..19e7f69 --- /dev/null +++ b/proxy/direct.go @@ -0,0 +1,43 @@ +package proxy + +import ( + "io" + "sync" + + "go.uber.org/zap" + + "github.com/9seconds/mtg/conntypes" + "github.com/9seconds/mtg/obfuscated2" + "github.com/9seconds/mtg/protocol" +) + +const directPipeBufferSize = 1024 * 1024 + +func directConnection(request *protocol.TelegramRequest) error { + telegramConnRaw, err := obfuscated2.TelegramProtocol(request) + if err != nil { + return err + } + telegramConn := telegramConnRaw.(conntypes.StreamReadWriteCloser) + defer telegramConn.Close() + + wg := &sync.WaitGroup{} + wg.Add(2) + + go directPipe(telegramConn, request.ClientConn, wg, request.Logger) + go directPipe(request.ClientConn, telegramConn, wg, request.Logger) + + <-request.Ctx.Done() + wg.Wait() + + return request.Ctx.Err() +} + +func directPipe(dst io.Writer, src io.Reader, wg *sync.WaitGroup, logger *zap.SugaredLogger) { + defer wg.Done() + + buf := make([]byte, directPipeBufferSize) + if _, err := io.CopyBuffer(dst, src, buf); err != nil { + logger.Debugw("Cannot pump sockets", "error", err) + } +} diff --git a/proxy/middle.go b/proxy/middle.go new file mode 100644 index 0000000..d806885 --- /dev/null +++ b/proxy/middle.go @@ -0,0 +1,60 @@ +package proxy + +import ( + "sync" + + "go.uber.org/zap" + + "github.com/9seconds/mtg/conntypes" + "github.com/9seconds/mtg/protocol" + "github.com/9seconds/mtg/wrappers/packetack" +) + +func middleConnection(request *protocol.TelegramRequest) error { + telegramConn := packetack.NewProxy(request) + defer telegramConn.Close() + + var clientConn conntypes.PacketAckFullReadWriteCloser + switch request.ClientProtocol.ConnectionType() { + case conntypes.ConnectionTypeAbridged: + clientConn = packetack.NewClientAbridged(request.ClientConn) + case conntypes.ConnectionTypeIntermediate: + clientConn = packetack.NewClientIntermediate(request.ClientConn) + case conntypes.ConnectionTypeSecure: + clientConn = packetack.NewClientIntermediateSecure(request.ClientConn) + default: + panic("unknown connection type") + } + + wg := &sync.WaitGroup{} + wg.Add(2) + + go middlePipe(telegramConn, clientConn, wg, request.Logger) + go middlePipe(clientConn, telegramConn, wg, request.Logger) + + <-request.Ctx.Done() + wg.Wait() + + return request.Ctx.Err() +} + +func middlePipe(dst conntypes.PacketAckWriter, + src conntypes.PacketAckReader, + wg *sync.WaitGroup, + logger *zap.SugaredLogger) { + defer wg.Done() + + for { + acks := conntypes.ConnectionAcks{} + packet, err := src.Read(&acks) + if err != nil { + logger.Debugw("Cannot read packet", "error", err) + return + } + + if err = dst.Write(packet, &acks); err != nil { + logger.Debugw("Cannot send packet", "error", err) + return + } + } +} diff --git a/proxy/proxy.go b/proxy/proxy.go index b7d04a3..46bc62f 100644 --- a/proxy/proxy.go +++ b/proxy/proxy.go @@ -2,23 +2,18 @@ package proxy import ( "context" - "io" "net" - "sync" "go.uber.org/zap" "github.com/9seconds/mtg/config" "github.com/9seconds/mtg/conntypes" - "github.com/9seconds/mtg/obfuscated2" "github.com/9seconds/mtg/protocol" "github.com/9seconds/mtg/stats" "github.com/9seconds/mtg/utils" "github.com/9seconds/mtg/wrappers/stream" ) -const directPipeBufferSize = 1024 * 1024 - type Proxy struct { Logger *zap.SugaredLogger Context context.Context @@ -89,46 +84,10 @@ func (p *Proxy) accept(conn net.Conn) { } if len(config.C.AdTag) > 0 { - err = p.acceptMiddleProxyConnection(req) + err = middleConnection(req) } else { - err = p.acceptDirectConnection(req) + err = directConnection(req) } logger.Infow("Client disconnected", "error", err, "addr", conn.RemoteAddr()) } - -func (p *Proxy) acceptDirectConnection(request *protocol.TelegramRequest) error { - telegramConnRaw, err := obfuscated2.TelegramProtocol(request) - if err != nil { - return err - } - telegramConn := telegramConnRaw.(conntypes.StreamReadWriteCloser) - defer telegramConn.Close() - - wg := &sync.WaitGroup{} - wg.Add(2) - - go p.directPipe(telegramConn, request.ClientConn, wg, request.Logger) - go p.directPipe(request.ClientConn, telegramConn, wg, request.Logger) - - <-request.Ctx.Done() - wg.Wait() - - return request.Ctx.Err() -} - -func (p *Proxy) directPipe(dst io.Writer, - src io.Reader, - wg *sync.WaitGroup, - logger *zap.SugaredLogger) { - defer wg.Done() - - buf := make([]byte, directPipeBufferSize) - if _, err := io.CopyBuffer(dst, src, buf); err != nil { - logger.Debugw("Cannot pump sockets", "error", err) - } -} - -func (p *Proxy) acceptMiddleProxyConnection(request *protocol.TelegramRequest) error { - return nil -}