From 7f01c03cb76fa90630d6d8156c578e311cda105b Mon Sep 17 00:00:00 2001 From: 9seconds Date: Tue, 5 Jun 2018 10:56:14 +0300 Subject: [PATCH] Refactor trafficrwc to wrappers --- proxy/server.go | 4 ++-- wrappers/timeoutrwc.go | 40 +++++++++++++++++++++++++++++++ {proxy => wrappers}/trafficrwc.go | 4 ++-- 3 files changed, 44 insertions(+), 4 deletions(-) create mode 100644 wrappers/timeoutrwc.go rename {proxy => wrappers}/trafficrwc.go (85%) diff --git a/proxy/server.go b/proxy/server.go index 43a9be2..26fcfec 100644 --- a/proxy/server.go +++ b/proxy/server.go @@ -113,7 +113,7 @@ func (s *Server) makeSocketID() string { func (s *Server) getClientStream(ctx context.Context, cancel context.CancelFunc, conn net.Conn, socketID string) (io.ReadWriteCloser, int16, error) { wConn := wrappers.NewTimeoutRWC(conn, s.readTimeout, s.writeTimeout) - wConn = newTrafficReadWriteCloser(wConn, s.stats.addIncomingTraffic, s.stats.addOutgoingTraffic) + wConn = wrappers.NewTrafficRWC(wConn, s.stats.addIncomingTraffic, s.stats.addOutgoingTraffic) frame, err := obfuscated2.ExtractFrame(wConn) if err != nil { return nil, 0, errors.Annotate(err, "Cannot create client stream") @@ -137,7 +137,7 @@ func (s *Server) getTelegramStream(ctx context.Context, cancel context.CancelFun return nil, errors.Annotate(err, "Cannot dial") } wConn := wrappers.NewTimeoutRWC(socket, s.readTimeout, s.writeTimeout) - wConn = newTrafficReadWriteCloser(wConn, s.stats.addIncomingTraffic, s.stats.addOutgoingTraffic) + wConn = wrappers.NewTrafficRWC(wConn, s.stats.addIncomingTraffic, s.stats.addOutgoingTraffic) obfs2, frame := obfuscated2.MakeTelegramObfuscated2Frame() if n, err := socket.Write(frame); err != nil || n != len(frame) { diff --git a/wrappers/timeoutrwc.go b/wrappers/timeoutrwc.go new file mode 100644 index 0000000..b83236c --- /dev/null +++ b/wrappers/timeoutrwc.go @@ -0,0 +1,40 @@ +package wrappers + +import ( + "io" + "net" + "time" +) + +// TimeoutReadWriteCloser sets timeouts for read/write into underlying +// network connection. +type TimeoutReadWriteCloser struct { + conn net.Conn + readTimeout time.Duration + writeTimeout time.Duration +} + +// Read reads from connection +func (t *TimeoutReadWriteCloser) Read(p []byte) (int, error) { + t.conn.SetReadDeadline(time.Now().Add(t.readTimeout)) // nolint: errcheck, gas + return t.conn.Read(p) +} + +// Write writes into connection. +func (t *TimeoutReadWriteCloser) Write(p []byte) (int, error) { + t.conn.SetWriteDeadline(time.Now().Add(t.writeTimeout)) // nolint: errcheck, gas + return t.conn.Write(p) +} + +// Close closes underlying connection. +func (t *TimeoutReadWriteCloser) Close() error { + return t.conn.Close() +} + +func NewTimeoutRWC(conn net.Conn, readTimeout, writeTimeout time.Duration) io.ReadWriteCloser { + return &TimeoutReadWriteCloser{ + conn: conn, + readTimeout: readTimeout, + writeTimeout: writeTimeout, + } +} diff --git a/proxy/trafficrwc.go b/wrappers/trafficrwc.go similarity index 85% rename from proxy/trafficrwc.go rename to wrappers/trafficrwc.go index a6be861..207addd 100644 --- a/proxy/trafficrwc.go +++ b/wrappers/trafficrwc.go @@ -1,4 +1,4 @@ -package proxy +package wrappers import "io" @@ -29,7 +29,7 @@ func (t *TrafficReadWriteCloser) Close() error { return t.conn.Close() } -func newTrafficReadWriteCloser(conn io.ReadWriteCloser, readCallback, writeCallback func(int)) io.ReadWriteCloser { +func NewTrafficRWC(conn io.ReadWriteCloser, readCallback, writeCallback func(int)) io.ReadWriteCloser { return &TrafficReadWriteCloser{ conn: conn, readCallback: readCallback,