Merge pull request #144 from 9seconds/tune-buffers

Add possibility to scale middleproxy buffers
This commit is contained in:
Sergey Arkhipov
2020-03-31 09:51:40 +03:00
committed by GitHub
5 changed files with 54 additions and 9 deletions
+1 -1
View File
@@ -50,7 +50,7 @@ func Proxy() error { // nolint: funlen
zap.S().Debugw("Configuration", "config", config.Printable())
if len(config.C.AdTag) > 0 {
if config.C.MiddleProxyMode() {
zap.S().Infow("Use middle proxy connection to Telegram")
diff, err := ntp.Fetch()
+47
View File
@@ -6,6 +6,7 @@ import (
"encoding/json"
"errors"
"fmt"
"math"
"net"
"github.com/alecthomas/units"
@@ -104,6 +105,52 @@ type Config struct {
AdTag []byte `json:"adtag"`
}
func (c *Config) ClientReadBuffer() int {
return c.ReadBuffer
}
func (c *Config) ClientWriteBuffer() int {
return c.WriteBuffer
}
func (c *Config) MiddleProxyMode() bool {
return len(c.AdTag) > 0
}
func (c *Config) ProxyReadBuffer() int {
value := c.ReadBuffer
if c.MiddleProxyMode() {
value = c.adjustProxyValue(value)
}
return value
}
func (c *Config) ProxyWriteBuffer() int {
value := c.WriteBuffer
if c.MiddleProxyMode() {
value = c.adjustProxyValue(value)
}
return value
}
func (c *Config) adjustProxyValue(value int) int {
if c.MultiplexPerConnection == 0 {
return value
}
fvalue := float64(value)
newValue := fvalue * 2 * math.Log(float64(c.MultiplexPerConnection))
newValue = math.Ceil(newValue)
newValue = math.Max(fvalue, newValue)
return int(newValue)
}
type Opt struct {
Option OptionType
Value interface{}
+2 -2
View File
@@ -51,7 +51,7 @@ func (p *Proxy) accept(conn net.Conn) {
connID := conntypes.NewConnID()
logger := p.Logger.With("connection_id", connID)
if err := utils.InitTCP(conn); err != nil {
if err := utils.InitTCP(conn, config.C.ClientReadBuffer(), config.C.ClientWriteBuffer()); err != nil {
logger.Errorw("Cannot initialize client TCP connection", "error", err)
return
}
@@ -90,7 +90,7 @@ func (p *Proxy) accept(conn net.Conn) {
err = nil
if len(config.C.AdTag) > 0 {
if config.C.MiddleProxyMode() {
middleConnection(req)
} else {
err = directConnection(req)
+1 -1
View File
@@ -37,7 +37,7 @@ func (b *baseTelegram) dial(dc conntypes.DC,
continue
}
if err := utils.InitTCP(conn); err != nil {
if err := utils.InitTCP(conn, config.C.ProxyReadBuffer(), config.C.ProxyWriteBuffer()); err != nil {
b.logger.Infow("Cannot initialize TCP socket", "address", addr, "error", err)
continue
}
+3 -5
View File
@@ -4,24 +4,22 @@ import (
"fmt"
"net"
"time"
"github.com/9seconds/mtg/config"
)
const tcpKeepAlivePingPeriod = 2 * time.Second
func InitTCP(conn net.Conn) error {
func InitTCP(conn net.Conn, readBufferSize int, writeBufferSize int) error {
tcpConn := conn.(*net.TCPConn)
if err := tcpConn.SetNoDelay(true); err != nil {
return fmt.Errorf("cannot set TCP_NO_DELAY: %w", err)
}
if err := tcpConn.SetReadBuffer(config.C.ReadBuffer); err != nil {
if err := tcpConn.SetReadBuffer(readBufferSize); err != nil {
return fmt.Errorf("cannot set read buffer size: %w", err)
}
if err := tcpConn.SetWriteBuffer(config.C.WriteBuffer); err != nil {
if err := tcpConn.SetWriteBuffer(writeBufferSize); err != nil {
return fmt.Errorf("cannot set write buffer size: %w", err)
}