Merge remote-tracking branch 'origin/master' into stable

This commit is contained in:
9seconds
2020-03-25 10:24:31 +03:00
25 changed files with 264 additions and 68 deletions
+5 -2
View File
@@ -194,11 +194,14 @@ supported environment variables:
| `MTG_STATSD_PREFIX` | `--statsd-prefix` | `mtg` | Which bucket prefix we should use. For example, if you set `mtg`, then metric `traffic.ingress` would be send as `mtg.traffic.ingress`. | | `MTG_STATSD_PREFIX` | `--statsd-prefix` | `mtg` | Which bucket prefix we should use. For example, if you set `mtg`, then metric `traffic.ingress` would be send as `mtg.traffic.ingress`. |
| `MTG_STATSD_TAGS_FORMAT` | `--statsd-tags-format` | | Which tags format we should use. By default, we are using default vanilla statsd tags format but if you want to send directly to InfluxDB or Datadog, please specify it there. Possible options are `influxdb` and `datadog`. | | `MTG_STATSD_TAGS_FORMAT` | `--statsd-tags-format` | | Which tags format we should use. By default, we are using default vanilla statsd tags format but if you want to send directly to InfluxDB or Datadog, please specify it there. Possible options are `influxdb` and `datadog`. |
| `MTG_STATSD_TAGS` | `--statsd-tags` | | Which tags should we send to statsd with our metrics. Please specify them as `key=value` pairs. | | `MTG_STATSD_TAGS` | `--statsd-tags` | | Which tags should we send to statsd with our metrics. Please specify them as `key=value` pairs. |
| `MTG_BUFFER_WRITE` | `-w`, `--write-buffer` | `64KB` | The size of TCP write buffer in bytes. Write buffer is the buffer for messages which are going from client to Telegram. | | `MTG_BUFFER_WRITE` | `-w`, `--write-buffer` | `32KB` | The size of TCP write buffer in bytes. Write buffer is the buffer for messages which are going from client to Telegram. |
| `MTG_BUFFER_READ` | `-r`, `--read-buffer` | `128KB` | The size of TCP read buffer in bytes. Read buffer is the buffer for messages from Telegram to client. | | `MTG_BUFFER_READ` | `-r`, `--read-buffer` | `32KB` | The size of TCP read buffer in bytes. Read buffer is the buffer for messages from Telegram to client. |
| `MTG_ANTIREPLAY_MAXSIZE` | `--anti-replay-max-size` | `128MB` | Max size of antireplay cache. | | `MTG_ANTIREPLAY_MAXSIZE` | `--anti-replay-max-size` | `128MB` | Max size of antireplay cache. |
| `MTG_CLOAK_PORT` | `--cloak-port` | `443` | Which port we should use to connect to cloaked host in FakeTLS mode. | | `MTG_CLOAK_PORT` | `--cloak-port` | `443` | Which port we should use to connect to cloaked host in FakeTLS mode. |
| `MTG_MULTIPLEX_PERCONNECTION` | `--multiplex-per-connection` | `50` | How many client connections can share a single Telegram connection in adtag mode | | `MTG_MULTIPLEX_PERCONNECTION` | `--multiplex-per-connection` | `50` | How many client connections can share a single Telegram connection in adtag mode |
| `MTG_NTP_SERVERS` | `--ntp-server` | default pool | A list of NTP servers to use. |
| `MTG_PREFER_DIRECT_IP` | `--prefer-ip` | `ipv6` | Which IP protocol to prefer if possible. Works mostly in direct mode. |
Usually you want to modify only read/write buffer sizes. If you feel Usually you want to modify only read/write buffer sizes. If you feel
that proxy is slow, try to increase both sizes giving more priority to that proxy is slow, try to increase both sizes giving more priority to
+6 -1
View File
@@ -63,7 +63,12 @@ func (c *ClientProtocol) tlsHandshake(conn io.ReadWriter) error {
return fmt.Errorf("cannot read initial record: %w", err) return fmt.Errorf("cannot read initial record: %w", err)
} }
clientHello, err := tlstypes.ParseClientHello(helloRecord.Data.Bytes()) buf := acquireBytesBuffer()
defer releaseBytesBuffer(buf)
helloRecord.Data.WriteBytes(buf)
clientHello, err := tlstypes.ParseClientHello(buf.Bytes())
if err != nil { if err != nil {
return fmt.Errorf("cannot parse client hello: %w", err) return fmt.Errorf("cannot parse client hello: %w", err)
} }
+11 -8
View File
@@ -28,15 +28,9 @@ func cloak(one, another io.ReadWriteCloser) {
wg.Add(2) wg.Add(2)
go func() { go cloakPipe(one, another, wg)
defer wg.Done()
io.Copy(one, another) // nolint: errcheck
}()
go func() { go cloakPipe(another, one, wg)
defer wg.Done()
io.Copy(another, one) // nolint: errcheck
}()
go func() { go func() {
wg.Wait() wg.Wait()
@@ -69,3 +63,12 @@ func cloak(one, another io.ReadWriteCloser) {
<-ctx.Done() <-ctx.Done()
} }
func cloakPipe(one io.Writer, another io.Reader, wg *sync.WaitGroup) {
defer wg.Done()
buf := acquireCloakBuffer()
defer releaseCloakBuffer(buf)
io.CopyBuffer(one, another, *buf) // nolint: errcheck
}
+39
View File
@@ -0,0 +1,39 @@
package faketls
import (
"bytes"
"sync"
)
const cloakBufferSize = 1024
var (
poolBytesBuffer = sync.Pool{
New: func() interface{} {
return &bytes.Buffer{}
},
}
poolCloakBuffer = sync.Pool{
New: func() interface{} {
rv := make([]byte, cloakBufferSize)
return &rv
},
}
)
func acquireBytesBuffer() *bytes.Buffer {
return poolBytesBuffer.Get().(*bytes.Buffer)
}
func acquireCloakBuffer() *[]byte {
return poolCloakBuffer.Get().(*[]byte)
}
func releaseBytesBuffer(buf *bytes.Buffer) {
buf.Reset()
poolBytesBuffer.Put(buf)
}
func releaseCloakBuffer(buf *[]byte) {
poolCloakBuffer.Put(buf)
}
+3 -3
View File
@@ -11,10 +11,10 @@ require (
github.com/prometheus/procfs v0.0.11 // indirect github.com/prometheus/procfs v0.0.11 // indirect
github.com/smira/go-statsd v1.3.1 github.com/smira/go-statsd v1.3.1
go.uber.org/zap v1.14.1 go.uber.org/zap v1.14.1
golang.org/x/crypto v0.0.0-20200317142112-1b76d66859c6 golang.org/x/crypto v0.0.0-20200323165209-0ec3e9974c59
golang.org/x/lint v0.0.0-20200302205851-738671d3881b // indirect golang.org/x/lint v0.0.0-20200302205851-738671d3881b // indirect
golang.org/x/net v0.0.0-20200319234117-63522dbf7eec // indirect golang.org/x/net v0.0.0-20200324143707-d3edc9973b7e // indirect
golang.org/x/sys v0.0.0-20200317113312-5766fd39f98d golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd
golang.org/x/tools v0.0.0-20200319210407-521f4a0cd458 // indirect golang.org/x/tools v0.0.0-20200319210407-521f4a0cd458 // indirect
gopkg.in/alecthomas/kingpin.v2 v2.2.6 gopkg.in/alecthomas/kingpin.v2 v2.2.6
honnef.co/go/tools v0.0.1-2020.1.3 // indirect honnef.co/go/tools v0.0.1-2020.1.3 // indirect
+6 -6
View File
@@ -123,8 +123,8 @@ golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2 h1:VklqNMn3ovrHsnt90Pveol
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
golang.org/x/crypto v0.0.0-20200317142112-1b76d66859c6 h1:TjszyFsQsyZNHwdVdZ5m7bjmreu0znc2kRYsEml9/Ww= golang.org/x/crypto v0.0.0-20200323165209-0ec3e9974c59 h1:3zb4D3T4G8jdExgVU/95+vQXfpEPiMdCaZgmGVxjNHM=
golang.org/x/crypto v0.0.0-20200317142112-1b76d66859c6/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= golang.org/x/crypto v0.0.0-20200323165209-0ec3e9974c59/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto=
golang.org/x/lint v0.0.0-20190930215403-16217165b5de h1:5hukYrvBGR8/eNkX5mdUezrA6JiaEZDtJb9Ei+1LlBs= golang.org/x/lint v0.0.0-20190930215403-16217165b5de h1:5hukYrvBGR8/eNkX5mdUezrA6JiaEZDtJb9Ei+1LlBs=
golang.org/x/lint v0.0.0-20190930215403-16217165b5de/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc= golang.org/x/lint v0.0.0-20190930215403-16217165b5de/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc=
golang.org/x/lint v0.0.0-20200302205851-738671d3881b h1:Wh+f8QHJXR411sJR8/vRBTZ7YapZaRvUcLFFJhusH0k= golang.org/x/lint v0.0.0-20200302205851-738671d3881b h1:Wh+f8QHJXR411sJR8/vRBTZ7YapZaRvUcLFFJhusH0k=
@@ -140,8 +140,8 @@ golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn
golang.org/x/net v0.0.0-20190613194153-d28f0bde5980/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20190613194153-d28f0bde5980/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/net v0.0.0-20200319234117-63522dbf7eec h1:w0SItUiQ4sBiXBAwWNkyu8Fu2Qpn/dtDIcoPkPDqjRw= golang.org/x/net v0.0.0-20200324143707-d3edc9973b7e h1:3G+cUijn7XD+S4eJFddp53Pv7+slrESplyjG25HgL+k=
golang.org/x/net v0.0.0-20200319234117-63522dbf7eec/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20200324143707-d3edc9973b7e/go.mod h1:qpuaurCH72eLCgpAm/N6yyVIVM9cpaDIP3A8BGJEC5A=
golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4 h1:YUO/7uOKsKeq9UokNS62b8FYywz3ker1l1vDZRCRefw= golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4 h1:YUO/7uOKsKeq9UokNS62b8FYywz3ker1l1vDZRCRefw=
golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
@@ -154,8 +154,8 @@ golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7w
golang.org/x/sys v0.0.0-20190422165155-953cdadca894/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190422165155-953cdadca894/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20200106162015-b016eb3dc98e/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200106162015-b016eb3dc98e/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20200122134326-e047566fdf82/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200122134326-e047566fdf82/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20200317113312-5766fd39f98d h1:62ap6LNOjDU6uGmKXHJbSfciMoV+FeI1sRXx/pLDL44= golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd h1:xhmwyvizuTgC2qz7ZlMluP20uW+C3Rm0FD/WLDX8884=
golang.org/x/sys v0.0.0-20200317113312-5766fd39f98d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/tools v0.0.0-20190311212946-11955173bddd/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs= golang.org/x/tools v0.0.0-20190311212946-11955173bddd/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs=
golang.org/x/tools v0.0.0-20190621195816-6e04913cbbac/go.mod h1:/rFqwRUd4F7ZHNgwSSTFct+R/Kf4OFW1sUzUTQQTgfc= golang.org/x/tools v0.0.0-20190621195816-6e04913cbbac/go.mod h1:/rFqwRUd4F7ZHNgwSSTFct+R/Kf4OFW1sUzUTQQTgfc=
+2 -2
View File
@@ -92,13 +92,13 @@ var (
"Write buffer size. You can think about it as a buffer from client to Telegram."). "Write buffer size. You can think about it as a buffer from client to Telegram.").
Short('w'). Short('w').
Envar("MTG_BUFFER_WRITE"). Envar("MTG_BUFFER_WRITE").
Default("64KB"). Default("32KB").
Bytes() Bytes()
runReadBufferSize = runCommand.Flag("read-buffer", runReadBufferSize = runCommand.Flag("read-buffer",
"Read buffer size. You can think about it as a buffer from Telegram to client."). "Read buffer size. You can think about it as a buffer from Telegram to client.").
Short('r'). Short('r').
Envar("MTG_BUFFER_READ"). Envar("MTG_BUFFER_READ").
Default("128KB"). Default("32KB").
Bytes() Bytes()
runTLSCloakPort = runCommand.Flag("cloak-port", runTLSCloakPort = runCommand.Flag("cloak-port",
"Port which should be used for host cloaking."). "Port which should be used for host cloaking.").
+4 -3
View File
@@ -11,7 +11,7 @@ import (
"github.com/9seconds/mtg/protocol" "github.com/9seconds/mtg/protocol"
) )
const directPipeBufferSize = 1024 * 1024 const directPipeBufferSize = 1024
func directConnection(request *protocol.TelegramRequest) error { func directConnection(request *protocol.TelegramRequest) error {
telegramConnRaw, err := obfuscated2.TelegramProtocol(request) telegramConnRaw, err := obfuscated2.TelegramProtocol(request)
@@ -42,8 +42,9 @@ func directPipe(dst io.WriteCloser, src io.ReadCloser, wg *sync.WaitGroup, logge
wg.Done() wg.Done()
}() }()
buf := make([]byte, directPipeBufferSize) buf := [directPipeBufferSize]byte{}
if _, err := io.CopyBuffer(dst, src, buf); err != nil {
if _, err := io.CopyBuffer(dst, src, buf[:]); err != nil {
logger.Debugw("Cannot pump sockets", "error", err) logger.Debugw("Cannot pump sockets", "error", err)
} }
} }
+1 -1
View File
@@ -25,7 +25,7 @@ func (c ClientHello) Digest() []byte {
} }
mac := hmac.New(sha256.New, config.C.Secret) mac := hmac.New(sha256.New, config.C.Secret)
mac.Write(rec.Bytes()) // nolint: errcheck rec.WriteBytes(mac)
computedDigest := mac.Sum(nil) computedDigest := mac.Sum(nil)
for i := range computedDigest { for i := range computedDigest {
+10 -3
View File
@@ -1,5 +1,7 @@
package tlstypes package tlstypes
import "io"
type RecordType uint8 type RecordType uint8
const ( const (
@@ -69,11 +71,16 @@ var (
) )
type Byter interface { type Byter interface {
Bytes() []byte WriteBytes(io.Writer)
Len() int
} }
type RawBytes []byte type RawBytes []byte
func (r RawBytes) Bytes() []byte { func (r RawBytes) WriteBytes(writer io.Writer) {
return []byte(r) writer.Write(r) // nolint: errcheck
}
func (r RawBytes) Len() int {
return len(r)
} }
+16 -9
View File
@@ -1,7 +1,7 @@
package tlstypes package tlstypes
import ( import (
"bytes" "io"
"github.com/9seconds/mtg/utils" "github.com/9seconds/mtg/utils"
) )
@@ -14,24 +14,31 @@ type Handshake struct {
Tail Byter Tail Byter
} }
func (h *Handshake) Bytes() []byte { func (h *Handshake) WriteBytes(writer io.Writer) {
buf := bytes.Buffer{} packetBuf := acquireBytesBuffer()
packetBuf := bytes.Buffer{} defer releaseBytesBuffer(packetBuf)
buf.WriteByte(byte(h.Type)) writer.Write([]byte{byte(h.Type)}) // nolint: errcheck
packetBuf.Write(h.Version.Bytes()) packetBuf.Write(h.Version.Bytes())
packetBuf.Write(h.Random[:]) packetBuf.Write(h.Random[:])
packetBuf.WriteByte(byte(len(h.SessionID))) packetBuf.WriteByte(byte(len(h.SessionID)))
packetBuf.Write(h.SessionID) packetBuf.Write(h.SessionID)
packetBuf.Write(h.Tail.Bytes()) h.Tail.WriteBytes(packetBuf)
sizeUint24 := utils.ToUint24(uint32(packetBuf.Len())) sizeUint24 := utils.ToUint24(uint32(packetBuf.Len()))
sizeUint24Bytes := sizeUint24[:] sizeUint24Bytes := sizeUint24[:]
sizeUint24Bytes[0], sizeUint24Bytes[2] = sizeUint24Bytes[2], sizeUint24Bytes[0] sizeUint24Bytes[0], sizeUint24Bytes[2] = sizeUint24Bytes[2], sizeUint24Bytes[0]
buf.Write(sizeUint24Bytes) writer.Write(sizeUint24Bytes) // nolint: errcheck
packetBuf.WriteTo(&buf) // nolint: errcheck packetBuf.WriteTo(writer) // nolint: errcheck
}
return buf.Bytes() func (h *Handshake) Len() int {
buf := acquireBytesBuffer()
defer releaseBytesBuffer(buf)
h.WriteBytes(buf)
return buf.Len()
} }
+23
View File
@@ -0,0 +1,23 @@
package tlstypes
import (
"bytes"
"sync"
)
var (
poolBytesBuffer = sync.Pool{
New: func() interface{} {
return &bytes.Buffer{}
},
}
)
func acquireBytesBuffer() *bytes.Buffer {
return poolBytesBuffer.Get().(*bytes.Buffer)
}
func releaseBytesBuffer(buf *bytes.Buffer) {
buf.Reset()
poolBytesBuffer.Put(buf)
}
+8 -9
View File
@@ -15,16 +15,15 @@ type Record struct {
Data Byter Data Byter
} }
func (r Record) Bytes() []byte { func (r Record) WriteBytes(writer io.Writer) {
buf := bytes.Buffer{} writer.Write([]byte{byte(r.Type)}) // nolint: errcheck
data := r.Data.Bytes() writer.Write(r.Version.Bytes()) // nolint: errcheck
binary.Write(writer, binary.BigEndian, uint16(r.Data.Len())) // nolint: errcheck
r.Data.WriteBytes(writer)
}
buf.WriteByte(byte(r.Type)) func (r Record) Len() int {
buf.Write(r.Version.Bytes()) return 1 + 2 + 2 + r.Data.Len()
binary.Write(&buf, binary.BigEndian, uint16(len(data))) // nolint: errcheck
buf.Write(data)
return buf.Bytes()
} }
func ReadRecord(reader io.Reader) (Record, error) { func ReadRecord(reader io.Reader) (Record, error) {
+6 -3
View File
@@ -20,20 +20,22 @@ type ServerHello struct {
} }
func (s ServerHello) WelcomePacket() []byte { func (s ServerHello) WelcomePacket() []byte {
buf := &bytes.Buffer{}
s.Random = [32]byte{} s.Random = [32]byte{}
rec := Record{ rec := Record{
Type: RecordTypeHandshake, Type: RecordTypeHandshake,
Version: Version12, Version: Version12,
Data: &s, Data: &s,
} }
buf := bytes.NewBuffer(rec.Bytes()) rec.WriteBytes(buf)
recChangeCipher := Record{ recChangeCipher := Record{
Type: RecordTypeChangeCipherSpec, Type: RecordTypeChangeCipherSpec,
Version: Version12, Version: Version12,
Data: RawBytes([]byte{0x01}), Data: RawBytes([]byte{0x01}),
} }
buf.Write(recChangeCipher.Bytes()) recChangeCipher.WriteBytes(buf)
hostCert := make([]byte, 1024+mrand.Intn(3092)) hostCert := make([]byte, 1024+mrand.Intn(3092))
rand.Read(hostCert) // nolint: errcheck rand.Read(hostCert) // nolint: errcheck
@@ -43,7 +45,8 @@ func (s ServerHello) WelcomePacket() []byte {
Version: Version12, Version: Version12,
Data: RawBytes(hostCert), Data: RawBytes(hostCert),
} }
buf.Write(recData.Bytes()) recData.WriteBytes(buf)
packet := buf.Bytes() packet := buf.Bytes()
mac := hmac.New(sha256.New, config.C.Secret) mac := hmac.New(sha256.New, config.C.Secret)
+11
View File
@@ -3,10 +3,13 @@ package utils
import ( import (
"fmt" "fmt"
"net" "net"
"time"
"github.com/9seconds/mtg/config" "github.com/9seconds/mtg/config"
) )
const tcpKeepAlivePingPeriod = 2 * time.Second
func InitTCP(conn net.Conn) error { func InitTCP(conn net.Conn) error {
tcpConn := conn.(*net.TCPConn) tcpConn := conn.(*net.TCPConn)
@@ -22,5 +25,13 @@ func InitTCP(conn net.Conn) error {
return fmt.Errorf("cannot set write buffer size: %w", err) return fmt.Errorf("cannot set write buffer size: %w", err)
} }
if err := tcpConn.SetKeepAlive(true); err != nil {
return fmt.Errorf("cannot enable keep-alive: %w", err)
}
if err := tcpConn.SetKeepAlivePeriod(tcpKeepAlivePingPeriod); err != nil {
return fmt.Errorf("cannot set keep-alive period: %w", err)
}
return nil return nil
} }
+5 -4
View File
@@ -42,7 +42,9 @@ type wrapperMtprotoFrame struct {
} }
func (w *wrapperMtprotoFrame) Read() (conntypes.Packet, error) { // nolint: funlen func (w *wrapperMtprotoFrame) Read() (conntypes.Packet, error) { // nolint: funlen
buf := &bytes.Buffer{} buf := acquireMtprotoFrameBytesBuffer()
defer releaseMtprotoFrameBytesBuffer(buf)
sum := crc32.NewIEEE() sum := crc32.NewIEEE()
writer := io.MultiWriter(buf, sum) writer := io.MultiWriter(buf, sum)
@@ -71,7 +73,6 @@ func (w *wrapperMtprotoFrame) Read() (conntypes.Packet, error) { // nolint: funl
} }
buf.Reset() buf.Reset()
buf.Grow(int(messageLength) - 4 - 4)
if _, err := io.CopyN(writer, w.parent, int64(messageLength)-4-4); err != nil { if _, err := io.CopyN(writer, w.parent, int64(messageLength)-4-4); err != nil {
return nil, fmt.Errorf("cannot read the message frame: %w", err) return nil, fmt.Errorf("cannot read the message frame: %w", err)
@@ -113,8 +114,8 @@ func (w *wrapperMtprotoFrame) Write(p conntypes.Packet) error {
messageLength := 4 + 4 + len(p) + 4 messageLength := 4 + 4 + len(p) + 4
paddingLength := (aes.BlockSize - messageLength%aes.BlockSize) % aes.BlockSize paddingLength := (aes.BlockSize - messageLength%aes.BlockSize) % aes.BlockSize
buf := &bytes.Buffer{} buf := acquireMtprotoFrameBytesBuffer()
buf.Grow(messageLength + paddingLength) defer releaseMtprotoFrameBytesBuffer(buf)
binary.Write(buf, binary.LittleEndian, uint32(messageLength)) // nolint: errcheck binary.Write(buf, binary.LittleEndian, uint32(messageLength)) // nolint: errcheck
binary.Write(buf, binary.LittleEndian, w.writeSeqNo) // nolint: errcheck binary.Write(buf, binary.LittleEndian, w.writeSeqNo) // nolint: errcheck
+23
View File
@@ -0,0 +1,23 @@
package packet
import (
"bytes"
"sync"
)
var (
poolMtprotoFrameBytesBuffer = sync.Pool{
New: func() interface{} {
return &bytes.Buffer{}
},
}
)
func acquireMtprotoFrameBytesBuffer() *bytes.Buffer {
return poolMtprotoFrameBytesBuffer.Get().(*bytes.Buffer)
}
func releaseMtprotoFrameBytesBuffer(buf *bytes.Buffer) {
buf.Reset()
poolMtprotoFrameBytesBuffer.Put(buf)
}
+3 -1
View File
@@ -88,7 +88,9 @@ func (w *wrapperClientAbridged) Write(packet conntypes.Packet, acks *conntypes.C
return nil return nil
case packetLength < clientAbridgedLargePacketLength: case packetLength < clientAbridgedLargePacketLength:
length24 := utils.ToUint24(uint32(packetLength)) length24 := utils.ToUint24(uint32(packetLength))
buf := bytes.Buffer{}
buf := acquireClientBytesBuffer()
defer releaseClientBytesBuffer(buf)
buf.WriteByte(byte(clientAbridgedSmallPacketLength)) buf.WriteByte(byte(clientAbridgedSmallPacketLength))
buf.Write(length24[:]) buf.Write(length24[:])
@@ -1,7 +1,6 @@
package packetack package packetack
import ( import (
"bytes"
"encoding/binary" "encoding/binary"
"fmt" "fmt"
"math/rand" "math/rand"
@@ -35,11 +34,13 @@ func (w *wrapperClientIntermediateSecure) Write(packet conntypes.Packet, acks *c
return nil return nil
} }
buf := bytes.Buffer{} buf := acquireClientBytesBuffer()
defer releaseClientBytesBuffer(buf)
paddingLength := rand.Intn(4) paddingLength := rand.Intn(4)
buf.Grow(4 + len(packet) + paddingLength) buf.Grow(4 + len(packet) + paddingLength)
binary.Write(&buf, binary.LittleEndian, uint32(len(packet)+paddingLength)) // nolint: errcheck binary.Write(buf, binary.LittleEndian, uint32(len(packet)+paddingLength)) // nolint: errcheck
buf.Write(packet) buf.Write(packet)
buf.Write(make([]byte, paddingLength)) buf.Write(make([]byte, paddingLength))
+23
View File
@@ -0,0 +1,23 @@
package packetack
import (
"bytes"
"sync"
)
var (
poolClientBytesBuffer = sync.Pool{
New: func() interface{} {
return &bytes.Buffer{}
},
}
)
func acquireClientBytesBuffer() *bytes.Buffer {
return poolClientBytesBuffer.Get().(*bytes.Buffer)
}
func releaseClientBytesBuffer(buf *bytes.Buffer) {
buf.Reset()
poolClientBytesBuffer.Put(buf)
}
+2 -1
View File
@@ -23,8 +23,8 @@ type wrapperProxy struct {
func (w *wrapperProxy) Write(packet conntypes.Packet, acks *conntypes.ConnectionAcks) error { func (w *wrapperProxy) Write(packet conntypes.Packet, acks *conntypes.ConnectionAcks) error {
buf := bytes.Buffer{} buf := bytes.Buffer{}
flags := w.flags flags := w.flags
if acks.Quick { if acks.Quick {
flags |= rpc.ProxyRequestFlagsQuickAck flags |= rpc.ProxyRequestFlagsQuickAck
} }
@@ -43,6 +43,7 @@ func (w *wrapperProxy) Write(packet conntypes.Packet, acks *conntypes.Connection
buf.WriteByte(byte(len(config.C.AdTag))) buf.WriteByte(byte(len(config.C.AdTag)))
buf.Write(config.C.AdTag) buf.Write(config.C.AdTag)
buf.Write(make([]byte, (4-buf.Len()%4)%4)) buf.Write(make([]byte, (4-buf.Len()%4)%4))
buf.Grow(len(packet))
buf.Write(packet) buf.Write(packet)
return w.proxy.Write(buf.Bytes()) return w.proxy.Write(buf.Bytes())
+13 -3
View File
@@ -1,6 +1,7 @@
package stream package stream
import ( import (
"bytes"
"errors" "errors"
"fmt" "fmt"
"net" "net"
@@ -39,13 +40,19 @@ func (w *wrapperFakeTLS) WriteTimeout(p []byte, timeout time.Duration) (int, err
func (w *wrapperFakeTLS) write(p []byte, writeFunc func([]byte) (int, error)) (int, error) { func (w *wrapperFakeTLS) write(p []byte, writeFunc func([]byte) (int, error)) (int, error) {
sum := 0 sum := 0
buf := acquireBytesBuffer()
defer releaseBytesBuffer(buf)
for _, v := range tlstypes.MakeRecords(p) { for _, v := range tlstypes.MakeRecords(p) {
_, err := writeFunc(v.Bytes()) buf.Reset()
v.WriteBytes(buf)
_, err := writeFunc(buf.Bytes())
if err != nil { if err != nil {
return sum, err return sum, err
} }
sum += len(v.Data.Bytes()) sum += v.Data.Len()
} }
return sum, nil return sum, nil
@@ -86,7 +93,10 @@ func NewFakeTLS(socket conntypes.StreamReadWriteCloser) conntypes.StreamReadWrit
switch rec.Type { switch rec.Type {
case tlstypes.RecordTypeChangeCipherSpec: case tlstypes.RecordTypeChangeCipherSpec:
case tlstypes.RecordTypeApplicationData: case tlstypes.RecordTypeApplicationData:
return rec.Data.Bytes(), nil buf := &bytes.Buffer{}
rec.Data.WriteBytes(buf)
return buf.Bytes(), nil
default: default:
return nil, fmt.Errorf("unsupported record type %v", rec.Type) return nil, fmt.Errorf("unsupported record type %v", rec.Type)
} }
+3 -2
View File
@@ -1,7 +1,6 @@
package stream package stream
import ( import (
"bytes"
"crypto/aes" "crypto/aes"
"crypto/cipher" "crypto/cipher"
"crypto/md5" // nolint: gosec "crypto/md5" // nolint: gosec
@@ -54,7 +53,9 @@ func mtprotoDeriveKeys(purpose mtprotoCipherPurpose,
resp *rpc.NonceResponse, resp *rpc.NonceResponse,
client, remote *net.TCPAddr, client, remote *net.TCPAddr,
secret []byte) ([]byte, []byte) { secret []byte) ([]byte, []byte) {
message := bytes.Buffer{} message := acquireBytesBuffer()
defer releaseBytesBuffer(message)
message.Write(resp.Nonce) // nolint: gosec message.Write(resp.Nonce) // nolint: gosec
message.Write(req.Nonce) // nolint: gosec message.Write(req.Nonce) // nolint: gosec
message.Write(req.CryptoTS) // nolint: gosec message.Write(req.CryptoTS) // nolint: gosec
+14 -4
View File
@@ -40,16 +40,26 @@ func (w *wrapperObfuscated2) Read(p []byte) (int, error) {
} }
func (w *wrapperObfuscated2) WriteTimeout(p []byte, timeout time.Duration) (int, error) { func (w *wrapperObfuscated2) WriteTimeout(p []byte, timeout time.Duration) (int, error) {
buf := make([]byte, len(p)) buffer := acquireBytesBuffer()
copy(buf, p) defer releaseBytesBuffer(buffer)
buffer.Write(p)
buf := buffer.Bytes()
w.encryptor.XORKeyStream(buf, buf) w.encryptor.XORKeyStream(buf, buf)
return w.parent.WriteTimeout(buf, timeout) return w.parent.WriteTimeout(buf, timeout)
} }
func (w *wrapperObfuscated2) Write(p []byte) (int, error) { func (w *wrapperObfuscated2) Write(p []byte) (int, error) {
buf := make([]byte, len(p)) buffer := acquireBytesBuffer()
copy(buf, p) defer releaseBytesBuffer(buffer)
buffer.Write(p)
buf := buffer.Bytes()
w.encryptor.XORKeyStream(buf, buf) w.encryptor.XORKeyStream(buf, buf)
return w.parent.Write(buf) return w.parent.Write(buf)
+23
View File
@@ -0,0 +1,23 @@
package stream
import (
"bytes"
"sync"
)
var (
poolBytesBuffer = sync.Pool{
New: func() interface{} {
return &bytes.Buffer{}
},
}
)
func acquireBytesBuffer() *bytes.Buffer {
return poolBytesBuffer.Get().(*bytes.Buffer)
}
func releaseBytesBuffer(buf *bytes.Buffer) {
buf.Reset()
poolBytesBuffer.Put(buf)
}