From 9ddf3e147c3a0da50475a8309978a47e0c9a7232 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Wed, 20 Jun 2018 08:24:17 +0300 Subject: [PATCH 1/5] Use bytes.Buffer pools for wrappers --- wrappers/buffer_pool.go | 27 +++++++++++++++++++++++++++ wrappers/streamcipherrwc.go | 24 ++++++++++++------------ 2 files changed, 39 insertions(+), 12 deletions(-) create mode 100644 wrappers/buffer_pool.go diff --git a/wrappers/buffer_pool.go b/wrappers/buffer_pool.go new file mode 100644 index 0000000..ead700d --- /dev/null +++ b/wrappers/buffer_pool.go @@ -0,0 +1,27 @@ +package wrappers + +import ( + "bytes" + "sync" +) + +var bufPool sync.Pool + +func getBuffer() *bytes.Buffer { + buf := bufPool.Get().(*bytes.Buffer) + buf.Reset() + + return buf +} + +func putBuffer(buf *bytes.Buffer) { + bufPool.Put(buf) +} + +func init() { + bufPool = sync.Pool{ + New: func() interface{} { + return &bytes.Buffer{} + }, + } +} diff --git a/wrappers/streamcipherrwc.go b/wrappers/streamcipherrwc.go index 77243b7..1d7d73c 100644 --- a/wrappers/streamcipherrwc.go +++ b/wrappers/streamcipherrwc.go @@ -22,20 +22,20 @@ func (c *StreamCipherReadWriteCloser) Read(p []byte) (n int, err error) { // Write writes into connection. func (c *StreamCipherReadWriteCloser) Write(p []byte) (int, error) { - encrypted := make([]byte, len(p)) + // This is to decrease an amount of allocations. Unfortunately, escape + // analysis in (at least Golang 1.10) is absolutely not perfect. For + // example, it understands that we want to have a slice locally, right? + // But since slice is effectively 2 ints + uintptr to [number]byte, the + // most heavyweight part is placed in heap. + buf := getBuffer() + defer putBuffer(buf) + buf.Grow(len(p)) + buf.Write(p) + + encrypted := buf.Bytes() c.encryptor.XORKeyStream(encrypted, p) - allWritten := 0 - for len(encrypted) > 0 { - n, err := c.conn.Write(encrypted) - allWritten += n - if err != nil { - return allWritten, err - } - encrypted = encrypted[n:] - } - - return allWritten, nil + return c.conn.Write(encrypted) } // Close closes underlying connection. From c76d8307c201e6dc15fdb634a7ec563c60638e34 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Wed, 20 Jun 2018 09:19:27 +0300 Subject: [PATCH 2/5] Use pool for frames --- Gopkg.lock | 4 ++-- client/direct.go | 1 + obfuscated2/frame.go | 26 ++++++++++++++++---------- obfuscated2/frame_pool.go | 24 ++++++++++++++++++++++++ obfuscated2/frame_test.go | 9 +++++---- obfuscated2/obfuscated2.go | 18 ++++++++++-------- obfuscated2/obfuscated2_test.go | 10 +++++----- telegram/direct.go | 4 +++- 8 files changed, 66 insertions(+), 30 deletions(-) create mode 100644 obfuscated2/frame_pool.go diff --git a/Gopkg.lock b/Gopkg.lock index e8b5d5c..88ad0e6 100644 --- a/Gopkg.lock +++ b/Gopkg.lock @@ -43,8 +43,8 @@ [[projects]] name = "github.com/stretchr/testify" packages = ["assert"] - revision = "12b6f73e6084dad08a7c6e575284b177ecafbc71" - version = "v1.2.1" + revision = "f35b8ab0b5a2cef36673838d662e249dd9c94686" + version = "v1.2.2" [[projects]] name = "go.uber.org/atomic" diff --git a/client/direct.go b/client/direct.go index 14776c4..cc06c3a 100644 --- a/client/direct.go +++ b/client/direct.go @@ -18,6 +18,7 @@ func DirectInit(conn net.Conn, conf *config.Config) (int16, io.ReadWriteCloser, if err != nil { return 0, nil, errors.Annotate(err, "Cannot extract frame") } + defer obfuscated2.ReturnFrame(frame) obfs2, dc, err := obfuscated2.ParseObfuscated2ClientFrame(conf.Secret, frame) if err != nil { diff --git a/obfuscated2/frame.go b/obfuscated2/frame.go index 9bf2bc1..eaf6521 100644 --- a/obfuscated2/frame.go +++ b/obfuscated2/frame.go @@ -68,29 +68,35 @@ func (f Frame) Valid() bool { // Invert inverts frame for extracting encryption keys. Pkease check that link: // https://blog.susanka.eu/how-telegram-obfuscates-its-mtproto-traffic/ -func (f Frame) Invert() Frame { - reversed := make(Frame, FrameLen) - copy(reversed, f) +func (f Frame) Invert() *Frame { + reversed := MakeFrame() + copy(*reversed, f) for i := 0; i < frameLenKey+frameLenIV; i++ { - reversed[frameOffsetFirst+i] = f[frameOffsetIV-1-i] + (*reversed)[frameOffsetFirst+i] = f[frameOffsetIV-1-i] } return reversed } // ExtractFrame extracts exact obfuscated2 handshake frame from given reader. -func ExtractFrame(conn io.Reader) (Frame, error) { - buf := &bytes.Buffer{} +func ExtractFrame(conn io.Reader) (*Frame, error) { + frame := MakeFrame() + buf := bytes.NewBuffer(*frame) + buf.Reset() + if _, err := io.CopyN(buf, conn, FrameLen); err != nil { + ReturnFrame(frame) return nil, errors.Annotate(err, "Cannot extract obfuscated header") } + copy(*frame, buf.Bytes()) - return Frame(buf.Bytes()), nil + return frame, nil } -func generateFrame() Frame { - data := make(Frame, FrameLen) +func generateFrame() *Frame { + frame := MakeFrame() + data := *frame for { if _, err := rand.Read(data); err != nil { @@ -112,6 +118,6 @@ func generateFrame() Frame { copy(data.Magic(), tgMagicBytes) - return data + return frame } } diff --git a/obfuscated2/frame_pool.go b/obfuscated2/frame_pool.go new file mode 100644 index 0000000..431073f --- /dev/null +++ b/obfuscated2/frame_pool.go @@ -0,0 +1,24 @@ +package obfuscated2 + +import "sync" + +var framePool sync.Pool + +// MakeFrame returns new pointer to the handshake frame. +func MakeFrame() *Frame { + return framePool.Get().(*Frame) +} + +// ReturnFrame returns pointer to the handshake frame back to the pool. +func ReturnFrame(f *Frame) { + framePool.Put(f) +} + +func init() { + framePool = sync.Pool{ + New: func() interface{} { + data := make(Frame, FrameLen) + return &data + }, + } +} diff --git a/obfuscated2/frame_test.go b/obfuscated2/frame_test.go index b772e3c..1a3f2b1 100644 --- a/obfuscated2/frame_test.go +++ b/obfuscated2/frame_test.go @@ -1,6 +1,7 @@ package obfuscated2 import ( + "bytes" "testing" "github.com/stretchr/testify/assert" @@ -47,21 +48,21 @@ func TestFrameValid(t *testing.T) { func TestFrameDoubleInvert(t *testing.T) { frame := makeFrame() - assert.Equal(t, frame, frame.Invert().Invert()) + assert.True(t, bytes.Equal(frame, *frame.Invert().Invert())) } func TestFrameInvert(t *testing.T) { frame := makeFrame() reversed := frame.Invert() - assert.Exactly(t, frame[:8], reversed[:8]) - assert.Exactly(t, frame[56:], reversed[56:]) + assert.Exactly(t, frame[:8], (*reversed)[:8]) + assert.Exactly(t, frame[56:], (*reversed)[56:]) toCompare := make([]byte, 48) for i := 0; i < 48; i++ { toCompare[i] = frame[55-i] } - assert.Equal(t, []byte(reversed[8:56]), toCompare) + assert.Equal(t, []byte((*reversed)[8:56]), toCompare) } func TestFrameGenerateValid(t *testing.T) { diff --git a/obfuscated2/obfuscated2.go b/obfuscated2/obfuscated2.go index 9c70c5d..15bde67 100644 --- a/obfuscated2/obfuscated2.go +++ b/obfuscated2/obfuscated2.go @@ -19,7 +19,7 @@ type Obfuscated2 struct { // details: http://telegra.ph/telegram-blocks-wtf-05-26 // // Beware, link above is in russian. -func ParseObfuscated2ClientFrame(secret []byte, frame Frame) (*Obfuscated2, int16, error) { +func ParseObfuscated2ClientFrame(secret []byte, frame *Frame) (*Obfuscated2, int16, error) { decHasher := sha256.New() decHasher.Write(frame.Key()) // nolint: errcheck decHasher.Write(secret) // nolint: errcheck @@ -31,8 +31,9 @@ func ParseObfuscated2ClientFrame(secret []byte, frame Frame) (*Obfuscated2, int1 encHasher.Write(secret) // nolint: errcheck encryptor := makeStreamCipher(encHasher.Sum(nil), invertedFrame.IV()) - decryptedFrame := make(Frame, FrameLen) - decryptor.XORKeyStream(decryptedFrame, frame) + decryptedFrame := MakeFrame() + defer ReturnFrame(decryptedFrame) + decryptor.XORKeyStream(*decryptedFrame, *frame) if !decryptedFrame.Valid() { return nil, 0, errors.New("Unknown protocol") } @@ -48,17 +49,18 @@ func ParseObfuscated2ClientFrame(secret []byte, frame Frame) (*Obfuscated2, int1 // MakeTelegramObfuscated2Frame creates new handshake frame to send to // Telegram. // https://blog.susanka.eu/how-telegram-obfuscates-its-mtproto-traffic/ -func MakeTelegramObfuscated2Frame() (*Obfuscated2, Frame) { +func MakeTelegramObfuscated2Frame() (*Obfuscated2, *Frame) { frame := generateFrame() encryptor := makeStreamCipher(frame.Key(), frame.IV()) decryptorFrame := frame.Invert() decryptor := makeStreamCipher(decryptorFrame.Key(), decryptorFrame.IV()) - copyFrame := make(Frame, frameOffsetIV) - copy(copyFrame, frame) - encryptor.XORKeyStream(frame, frame) - copy(frame, copyFrame) + copyFrame := MakeFrame() + defer ReturnFrame(copyFrame) + copy((*copyFrame)[:frameOffsetIV], (*frame)[:frameOffsetIV]) + encryptor.XORKeyStream(*frame, *frame) + copy((*frame)[:frameOffsetIV], (*copyFrame)[:frameOffsetIV]) obfs := &Obfuscated2{ Decryptor: decryptor, diff --git a/obfuscated2/obfuscated2_test.go b/obfuscated2/obfuscated2_test.go index 49255df..54e84c8 100644 --- a/obfuscated2/obfuscated2_test.go +++ b/obfuscated2/obfuscated2_test.go @@ -12,7 +12,7 @@ func TestObfs2TelegramFrameDecrypt(t *testing.T) { decryptor := makeStreamCipher(frame.Key(), frame.IV()) decrypted := make(Frame, FrameLen) - decryptor.XORKeyStream(decrypted, frame) + decryptor.XORKeyStream(decrypted, *frame) assert.True(t, decrypted.Valid()) } @@ -42,8 +42,8 @@ func TestObfs2Full(t *testing.T) { encryptor := makeStreamCipher(clientKey, clientFrame.IV()) encrypted := make(Frame, FrameLen) - encryptor.XORKeyStream(encrypted, clientFrame) - copy(encrypted[:56], clientFrame[:56]) + encryptor.XORKeyStream(encrypted, *clientFrame) + copy(encrypted[:56], (*clientFrame)[:56]) invertedClientFrame := clientFrame.Invert() clientHasher = sha256.New() @@ -52,13 +52,13 @@ func TestObfs2Full(t *testing.T) { invertedClientKey := clientHasher.Sum(nil) clientDecryptor := makeStreamCipher(invertedClientKey, invertedClientFrame.IV()) - clientObfs, _, err := ParseObfuscated2ClientFrame(secret, encrypted) + clientObfs, _, err := ParseObfuscated2ClientFrame(secret, &encrypted) assert.Nil(t, err) tgObfs, tgFrame := MakeTelegramObfuscated2Frame() tgDecryptor := makeStreamCipher(tgFrame.Key(), tgFrame.IV()) decrypted := make(Frame, FrameLen) - tgDecryptor.XORKeyStream(decrypted, tgFrame) + tgDecryptor.XORKeyStream(decrypted, *tgFrame) assert.True(t, decrypted.Valid()) tgInvertedFrame := tgFrame.Invert() diff --git a/telegram/direct.go b/telegram/direct.go index f2437ea..0c9e8fd 100644 --- a/telegram/direct.go +++ b/telegram/direct.go @@ -43,7 +43,9 @@ func (t *directTelegram) Dial(dcIdx int16) (io.ReadWriteCloser, error) { func (t *directTelegram) Init(conn io.ReadWriteCloser) (io.ReadWriteCloser, error) { obfs2, frame := obfuscated2.MakeTelegramObfuscated2Frame() - if n, err := conn.Write(frame); err != nil || n != len(frame) { + defer obfuscated2.ReturnFrame(frame) + + if n, err := conn.Write(*frame); err != nil || n != len(*frame) { return nil, errors.Annotate(err, "Cannot write hadnshake frame") } From 2c938d284894a257a26e0800e42286e876cb975e Mon Sep 17 00:00:00 2001 From: 9seconds Date: Wed, 20 Jun 2018 09:33:57 +0300 Subject: [PATCH 3/5] Reuse buffers on socket pumping --- proxy/copy_pool.go | 16 ++++++++++++++++ proxy/server.go | 21 +++++++++++++-------- 2 files changed, 29 insertions(+), 8 deletions(-) create mode 100644 proxy/copy_pool.go diff --git a/proxy/copy_pool.go b/proxy/copy_pool.go new file mode 100644 index 0000000..b6eb784 --- /dev/null +++ b/proxy/copy_pool.go @@ -0,0 +1,16 @@ +package proxy + +import "sync" + +const copyBufferSize = 30 * 1024 + +var copyPool sync.Pool + +func init() { + copyPool = sync.Pool{ + New: func() interface{} { + data := make([]byte, copyBufferSize) + return &data + }, + } +} diff --git a/proxy/server.go b/proxy/server.go index bd5564d..4e2d335 100644 --- a/proxy/server.go +++ b/proxy/server.go @@ -83,14 +83,10 @@ func (s *Server) accept(conn net.Conn) { wait := &sync.WaitGroup{} wait.Add(2) - go func() { - defer wait.Done() - io.Copy(clientConn, tgConn) // nolint: errcheck - }() - go func() { - defer wait.Done() - io.Copy(tgConn, clientConn) // nolint: errcheck - }() + + go s.pipe(clientConn, tgConn, wait) + go s.pipe(tgConn, clientConn, wait) + <-ctx.Done() wait.Wait() @@ -131,6 +127,15 @@ func (s *Server) getTelegramStream(ctx context.Context, cancel context.CancelFun return conn, nil } +func (s *Server) pipe(dst io.Writer, src io.Reader, wait *sync.WaitGroup) { + defer wait.Done() + + buf := copyPool.Get().(*[]byte) + defer copyPool.Put(buf) + + io.CopyBuffer(dst, src, *buf) // nolint: errcheck +} + // NewServer creates new instance of MTPROTO proxy. func NewServer(conf *config.Config, logger *zap.SugaredLogger, stat *Stats) *Server { return &Server{ From 4178ffe44b8b3dd690dc97963cda3748f8bbb061 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Wed, 20 Jun 2018 09:36:06 +0300 Subject: [PATCH 4/5] Fix run-mtg --- run-mtg.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/run-mtg.sh b/run-mtg.sh index 7cc965e..96d72d3 100755 --- a/run-mtg.sh +++ b/run-mtg.sh @@ -12,7 +12,7 @@ STAT_PORT=3129 chmod 0400 "$SECRET_PATH" ) -# docker pull "$IMAGE_NAME" +docker pull "$IMAGE_NAME" docker ps --filter "Name=$CONTAINER_NAME" -aq | xargs -r docker rm -fv docker run \ -d \ From 0a924018e586dc1d22613f8e9c90ef5b85fa716b Mon Sep 17 00:00:00 2001 From: 9seconds Date: Wed, 20 Jun 2018 09:41:24 +0300 Subject: [PATCH 5/5] Correct math/rand seed initialization --- main.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/main.go b/main.go index 6ba0cb6..cf3d140 100644 --- a/main.go +++ b/main.go @@ -5,7 +5,9 @@ package main import ( "encoding/json" "io" + "math/rand" "os" + "time" "go.uber.org/zap" "go.uber.org/zap/zapcore" @@ -80,6 +82,8 @@ var ( ) func main() { + rand.Seed(time.Now().UTC().UnixNano()) + app.Version(version) kingpin.MustParse(app.Parse(os.Args[1:]))