mirror of
https://github.com/ScuroNeko/mtg.git
synced 2026-09-01 16:01:55 +03:00
Rework logging
This commit is contained in:
+6
-14
@@ -6,6 +6,8 @@ import (
|
|||||||
"crypto/cipher"
|
"crypto/cipher"
|
||||||
"net"
|
"net"
|
||||||
|
|
||||||
|
"go.uber.org/zap"
|
||||||
|
|
||||||
"github.com/9seconds/mtg/utils"
|
"github.com/9seconds/mtg/utils"
|
||||||
"github.com/juju/errors"
|
"github.com/juju/errors"
|
||||||
)
|
)
|
||||||
@@ -13,6 +15,7 @@ import (
|
|||||||
type BlockCipher struct {
|
type BlockCipher struct {
|
||||||
buf *bytes.Buffer
|
buf *bytes.Buffer
|
||||||
|
|
||||||
|
logger *zap.SugaredLogger
|
||||||
conn StreamReadWriteCloser
|
conn StreamReadWriteCloser
|
||||||
encryptor cipher.BlockMode
|
encryptor cipher.BlockMode
|
||||||
decryptor cipher.BlockMode
|
decryptor cipher.BlockMode
|
||||||
@@ -60,20 +63,8 @@ func (b *BlockCipher) Write(p []byte) (int, error) {
|
|||||||
return b.conn.Write(encrypted)
|
return b.conn.Write(encrypted)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *BlockCipher) LogDebug(msg string, data ...interface{}) {
|
func (b *BlockCipher) Logger() *zap.SugaredLogger {
|
||||||
b.conn.LogDebug(msg, data...)
|
return b.logger
|
||||||
}
|
|
||||||
|
|
||||||
func (b *BlockCipher) LogInfo(msg string, data ...interface{}) {
|
|
||||||
b.conn.LogInfo(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (b *BlockCipher) LogWarn(msg string, data ...interface{}) {
|
|
||||||
b.conn.LogWarn(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (b *BlockCipher) LogError(msg string, data ...interface{}) {
|
|
||||||
b.conn.LogError(msg, data...)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *BlockCipher) LocalAddr() *net.TCPAddr {
|
func (b *BlockCipher) LocalAddr() *net.TCPAddr {
|
||||||
@@ -92,6 +83,7 @@ func NewBlockCipher(conn StreamReadWriteCloser, encryptor, decryptor cipher.Bloc
|
|||||||
return &BlockCipher{
|
return &BlockCipher{
|
||||||
buf: &bytes.Buffer{},
|
buf: &bytes.Buffer{},
|
||||||
conn: conn,
|
conn: conn,
|
||||||
|
logger: conn.Logger().Named("block-cipher"),
|
||||||
encryptor: encryptor,
|
encryptor: encryptor,
|
||||||
decryptor: decryptor,
|
decryptor: decryptor,
|
||||||
}
|
}
|
||||||
|
|||||||
+4
-16
@@ -57,7 +57,7 @@ func (c *Conn) Read(p []byte) (int, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (c *Conn) Close() error {
|
func (c *Conn) Close() error {
|
||||||
defer c.LogDebug("Closed connection")
|
defer c.logger.Debugw("Closed connection")
|
||||||
return c.conn.Close()
|
return c.conn.Close()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -80,20 +80,8 @@ func (c *Conn) RemoteAddr() *net.TCPAddr {
|
|||||||
return c.conn.RemoteAddr().(*net.TCPAddr)
|
return c.conn.RemoteAddr().(*net.TCPAddr)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Conn) LogDebug(msg string, data ...interface{}) {
|
func (c *Conn) Logger() *zap.SugaredLogger {
|
||||||
c.logger.Debugw(msg, data...)
|
return c.logger
|
||||||
}
|
|
||||||
|
|
||||||
func (c *Conn) LogInfo(msg string, data ...interface{}) {
|
|
||||||
c.logger.Infow(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *Conn) LogWarn(msg string, data ...interface{}) {
|
|
||||||
c.logger.Warnw(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *Conn) LogError(msg string, data ...interface{}) {
|
|
||||||
c.logger.Errorw(msg, data...)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewConn(conn net.Conn, connID string, purpose ConnPurpose, publicIPv4, publicIPv6 net.IP) StreamReadWriteCloser {
|
func NewConn(conn net.Conn, connID string, purpose ConnPurpose, publicIPv4, publicIPv6 net.IP) StreamReadWriteCloser {
|
||||||
@@ -102,7 +90,7 @@ func NewConn(conn net.Conn, connID string, purpose ConnPurpose, publicIPv4, publ
|
|||||||
"local_address", conn.LocalAddr(),
|
"local_address", conn.LocalAddr(),
|
||||||
"remote_address", conn.RemoteAddr(),
|
"remote_address", conn.RemoteAddr(),
|
||||||
"purpose", purpose,
|
"purpose", purpose,
|
||||||
)
|
).Named("conn")
|
||||||
|
|
||||||
wrapper := Conn{
|
wrapper := Conn{
|
||||||
logger: logger,
|
logger: logger,
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
|
|
||||||
"github.com/juju/errors"
|
"github.com/juju/errors"
|
||||||
|
"go.uber.org/zap"
|
||||||
|
|
||||||
"github.com/9seconds/mtg/mtproto"
|
"github.com/9seconds/mtg/mtproto"
|
||||||
"github.com/9seconds/mtg/utils"
|
"github.com/9seconds/mtg/utils"
|
||||||
@@ -18,15 +19,16 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type MTProtoAbridged struct {
|
type MTProtoAbridged struct {
|
||||||
conn StreamReadWriteCloser
|
conn StreamReadWriteCloser
|
||||||
opts *mtproto.ConnectionOpts
|
opts *mtproto.ConnectionOpts
|
||||||
|
logger *zap.SugaredLogger
|
||||||
|
|
||||||
readCounter uint32
|
readCounter uint32
|
||||||
writeCounter uint32
|
writeCounter uint32
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoAbridged) Read() ([]byte, error) {
|
func (m *MTProtoAbridged) Read() ([]byte, error) {
|
||||||
m.LogDebug("Read packet",
|
m.logger.Debugw("Read packet",
|
||||||
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
||||||
"quick_ack", m.opts.ReadHacks.QuickAck,
|
"quick_ack", m.opts.ReadHacks.QuickAck,
|
||||||
"counter", m.readCounter,
|
"counter", m.readCounter,
|
||||||
@@ -41,7 +43,7 @@ func (m *MTProtoAbridged) Read() ([]byte, error) {
|
|||||||
msgLength := uint8(buf.Bytes()[0])
|
msgLength := uint8(buf.Bytes()[0])
|
||||||
buf.Reset()
|
buf.Reset()
|
||||||
|
|
||||||
m.LogDebug("Packet first byte",
|
m.logger.Debugw("Packet first byte",
|
||||||
"byte", msgLength,
|
"byte", msgLength,
|
||||||
"counter", m.readCounter,
|
"counter", m.readCounter,
|
||||||
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
||||||
@@ -64,7 +66,7 @@ func (m *MTProtoAbridged) Read() ([]byte, error) {
|
|||||||
}
|
}
|
||||||
msgLength32 *= 4
|
msgLength32 *= 4
|
||||||
|
|
||||||
m.LogDebug("Packet length",
|
m.logger.Debugw("Packet length",
|
||||||
"length", msgLength32,
|
"length", msgLength32,
|
||||||
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
||||||
"quick_ack", m.opts.ReadHacks.QuickAck,
|
"quick_ack", m.opts.ReadHacks.QuickAck,
|
||||||
@@ -83,7 +85,7 @@ func (m *MTProtoAbridged) Read() ([]byte, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoAbridged) Write(p []byte) (int, error) {
|
func (m *MTProtoAbridged) Write(p []byte) (int, error) {
|
||||||
m.LogDebug("Write packet",
|
m.logger.Debugw("Write packet",
|
||||||
"length", len(p),
|
"length", len(p),
|
||||||
"simple_ack", m.opts.WriteHacks.SimpleAck,
|
"simple_ack", m.opts.WriteHacks.SimpleAck,
|
||||||
"quick_ack", m.opts.WriteHacks.QuickAck,
|
"quick_ack", m.opts.WriteHacks.QuickAck,
|
||||||
@@ -124,24 +126,8 @@ func (m *MTProtoAbridged) Write(p []byte) (int, error) {
|
|||||||
return 0, errors.Errorf("Packet is too big %d", len(p))
|
return 0, errors.Errorf("Packet is too big %d", len(p))
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoAbridged) LogDebug(msg string, data ...interface{}) {
|
func (m *MTProtoAbridged) Logger() *zap.SugaredLogger {
|
||||||
data = append(data, []interface{}{"type", "abridged"}...)
|
return m.logger
|
||||||
m.conn.LogDebug(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoAbridged) LogInfo(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "abridged"}...)
|
|
||||||
m.conn.LogInfo(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoAbridged) LogWarn(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "abridged"}...)
|
|
||||||
m.conn.LogWarn(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoAbridged) LogError(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "abridged"}...)
|
|
||||||
m.conn.LogError(msg, data...)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoAbridged) LocalAddr() *net.TCPAddr {
|
func (m *MTProtoAbridged) LocalAddr() *net.TCPAddr {
|
||||||
@@ -158,7 +144,8 @@ func (m *MTProtoAbridged) Close() error {
|
|||||||
|
|
||||||
func NewMTProtoAbridged(conn StreamReadWriteCloser, opts *mtproto.ConnectionOpts) PacketReadWriteCloser {
|
func NewMTProtoAbridged(conn StreamReadWriteCloser, opts *mtproto.ConnectionOpts) PacketReadWriteCloser {
|
||||||
return &MTProtoAbridged{
|
return &MTProtoAbridged{
|
||||||
conn: conn,
|
conn: conn,
|
||||||
opts: opts,
|
opts: opts,
|
||||||
|
logger: conn.Logger().Named("mtproto-abridged"),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+10
-22
@@ -10,6 +10,7 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
|
|
||||||
"github.com/juju/errors"
|
"github.com/juju/errors"
|
||||||
|
"go.uber.org/zap"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -20,7 +21,9 @@ const (
|
|||||||
var mtprotoFramePadding = []byte{0x04, 0x00, 0x00, 0x00}
|
var mtprotoFramePadding = []byte{0x04, 0x00, 0x00, 0x00}
|
||||||
|
|
||||||
type MTProtoFrame struct {
|
type MTProtoFrame struct {
|
||||||
conn StreamReadWriteCloser
|
conn StreamReadWriteCloser
|
||||||
|
logger *zap.SugaredLogger
|
||||||
|
|
||||||
readSeqNo int32
|
readSeqNo int32
|
||||||
writeSeqNo int32
|
writeSeqNo int32
|
||||||
}
|
}
|
||||||
@@ -42,7 +45,7 @@ func (m *MTProtoFrame) Read() ([]byte, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
messageLength := binary.LittleEndian.Uint32(buf.Bytes())
|
messageLength := binary.LittleEndian.Uint32(buf.Bytes())
|
||||||
m.LogDebug("Read MTProto frame",
|
m.logger.Debugw("Read MTProto frame",
|
||||||
"messageLength", messageLength,
|
"messageLength", messageLength,
|
||||||
"sequence_number", m.readSeqNo,
|
"sequence_number", m.readSeqNo,
|
||||||
)
|
)
|
||||||
@@ -75,7 +78,7 @@ func (m *MTProtoFrame) Read() ([]byte, error) {
|
|||||||
return nil, errors.Errorf("CRC32 checksum mismatch. Wait for %d, got %d", sum.Sum32(), checksum)
|
return nil, errors.Errorf("CRC32 checksum mismatch. Wait for %d, got %d", sum.Sum32(), checksum)
|
||||||
}
|
}
|
||||||
|
|
||||||
m.LogDebug("Read MTProto frame",
|
m.logger.Debugw("Read MTProto frame",
|
||||||
"messageLength", messageLength,
|
"messageLength", messageLength,
|
||||||
"sequence_number", m.readSeqNo,
|
"sequence_number", m.readSeqNo,
|
||||||
"dataLength", len(data),
|
"dataLength", len(data),
|
||||||
@@ -101,7 +104,7 @@ func (m *MTProtoFrame) Write(p []byte) (int, error) {
|
|||||||
binary.Write(buf, binary.LittleEndian, checksum)
|
binary.Write(buf, binary.LittleEndian, checksum)
|
||||||
buf.Write(bytes.Repeat(mtprotoFramePadding, paddingLength/4))
|
buf.Write(bytes.Repeat(mtprotoFramePadding, paddingLength/4))
|
||||||
|
|
||||||
m.LogDebug("Write MTProto frame",
|
m.logger.Debugw("Write MTProto frame",
|
||||||
"length", len(p),
|
"length", len(p),
|
||||||
"sequence_number", m.writeSeqNo,
|
"sequence_number", m.writeSeqNo,
|
||||||
"crc32", checksum,
|
"crc32", checksum,
|
||||||
@@ -114,24 +117,8 @@ func (m *MTProtoFrame) Write(p []byte) (int, error) {
|
|||||||
return len(p), err
|
return len(p), err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoFrame) LogDebug(msg string, data ...interface{}) {
|
func (m *MTProtoFrame) Logger() *zap.SugaredLogger {
|
||||||
data = append(data, []interface{}{"type", "frame"}...)
|
return m.logger
|
||||||
m.conn.LogDebug(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoFrame) LogInfo(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "frame"}...)
|
|
||||||
m.conn.LogInfo(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoFrame) LogWarn(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "frame"}...)
|
|
||||||
m.conn.LogWarn(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoFrame) LogError(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "frame"}...)
|
|
||||||
m.conn.LogError(msg, data...)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoFrame) LocalAddr() *net.TCPAddr {
|
func (m *MTProtoFrame) LocalAddr() *net.TCPAddr {
|
||||||
@@ -149,6 +136,7 @@ func (m *MTProtoFrame) Close() error {
|
|||||||
func NewMTProtoFrame(conn StreamReadWriteCloser, seqNo int32) PacketReadWriteCloser {
|
func NewMTProtoFrame(conn StreamReadWriteCloser, seqNo int32) PacketReadWriteCloser {
|
||||||
return &MTProtoFrame{
|
return &MTProtoFrame{
|
||||||
conn: conn,
|
conn: conn,
|
||||||
|
logger: conn.Logger().Named("mtproto-frame"),
|
||||||
readSeqNo: seqNo,
|
readSeqNo: seqNo,
|
||||||
writeSeqNo: seqNo,
|
writeSeqNo: seqNo,
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -6,22 +6,25 @@ import (
|
|||||||
"io"
|
"io"
|
||||||
"net"
|
"net"
|
||||||
|
|
||||||
"github.com/9seconds/mtg/mtproto"
|
|
||||||
"github.com/juju/errors"
|
"github.com/juju/errors"
|
||||||
|
"go.uber.org/zap"
|
||||||
|
|
||||||
|
"github.com/9seconds/mtg/mtproto"
|
||||||
)
|
)
|
||||||
|
|
||||||
const mtprotoIntermediateQuickAckLength = 0x80000000
|
const mtprotoIntermediateQuickAckLength = 0x80000000
|
||||||
|
|
||||||
type MTProtoIntermediate struct {
|
type MTProtoIntermediate struct {
|
||||||
conn StreamReadWriteCloser
|
conn StreamReadWriteCloser
|
||||||
opts *mtproto.ConnectionOpts
|
opts *mtproto.ConnectionOpts
|
||||||
|
logger *zap.SugaredLogger
|
||||||
|
|
||||||
readCounter uint32
|
readCounter uint32
|
||||||
writeCounter uint32
|
writeCounter uint32
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoIntermediate) Read() ([]byte, error) {
|
func (m *MTProtoIntermediate) Read() ([]byte, error) {
|
||||||
m.LogDebug("Read packet",
|
m.logger.Debugw("Read packet",
|
||||||
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
||||||
"quick_ack", m.opts.ReadHacks.QuickAck,
|
"quick_ack", m.opts.ReadHacks.QuickAck,
|
||||||
"counter", m.readCounter,
|
"counter", m.readCounter,
|
||||||
@@ -35,7 +38,7 @@ func (m *MTProtoIntermediate) Read() ([]byte, error) {
|
|||||||
}
|
}
|
||||||
length := binary.LittleEndian.Uint32(buf.Bytes())
|
length := binary.LittleEndian.Uint32(buf.Bytes())
|
||||||
|
|
||||||
m.LogDebug("Packet message length",
|
m.logger.Debugw("Packet message length",
|
||||||
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
"simple_ack", m.opts.ReadHacks.SimpleAck,
|
||||||
"quick_ack", m.opts.ReadHacks.QuickAck,
|
"quick_ack", m.opts.ReadHacks.QuickAck,
|
||||||
"counter", m.readCounter,
|
"counter", m.readCounter,
|
||||||
@@ -62,7 +65,7 @@ func (m *MTProtoIntermediate) Read() ([]byte, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoIntermediate) Write(p []byte) (int, error) {
|
func (m *MTProtoIntermediate) Write(p []byte) (int, error) {
|
||||||
m.LogDebug("Write packet",
|
m.logger.Debugw("Write packet",
|
||||||
"simple_ack", m.opts.WriteHacks.SimpleAck,
|
"simple_ack", m.opts.WriteHacks.SimpleAck,
|
||||||
"quick_ack", m.opts.WriteHacks.QuickAck,
|
"quick_ack", m.opts.WriteHacks.QuickAck,
|
||||||
"counter", m.writeCounter,
|
"counter", m.writeCounter,
|
||||||
@@ -79,24 +82,8 @@ func (m *MTProtoIntermediate) Write(p []byte) (int, error) {
|
|||||||
return m.conn.Write(append(length[:], p...))
|
return m.conn.Write(append(length[:], p...))
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoIntermediate) LogDebug(msg string, data ...interface{}) {
|
func (m *MTProtoIntermediate) Logger() *zap.SugaredLogger {
|
||||||
data = append(data, []interface{}{"type", "intermediate"}...)
|
return m.logger
|
||||||
m.conn.LogDebug(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoIntermediate) LogInfo(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "intermediate"}...)
|
|
||||||
m.conn.LogInfo(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoIntermediate) LogWarn(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "intermediate"}...)
|
|
||||||
m.conn.LogWarn(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoIntermediate) LogError(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "intermediate"}...)
|
|
||||||
m.conn.LogError(msg, data...)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoIntermediate) LocalAddr() *net.TCPAddr {
|
func (m *MTProtoIntermediate) LocalAddr() *net.TCPAddr {
|
||||||
@@ -113,7 +100,8 @@ func (m *MTProtoIntermediate) Close() error {
|
|||||||
|
|
||||||
func NewMTProtoIntermediate(conn StreamReadWriteCloser, opts *mtproto.ConnectionOpts) PacketReadWriteCloser {
|
func NewMTProtoIntermediate(conn StreamReadWriteCloser, opts *mtproto.ConnectionOpts) PacketReadWriteCloser {
|
||||||
return &MTProtoIntermediate{
|
return &MTProtoIntermediate{
|
||||||
conn: conn,
|
conn: conn,
|
||||||
opts: opts,
|
logger: conn.Logger().Named("mtproto-intermediate"),
|
||||||
|
opts: opts,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+13
-26
@@ -5,21 +5,23 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
|
|
||||||
"github.com/juju/errors"
|
"github.com/juju/errors"
|
||||||
|
"go.uber.org/zap"
|
||||||
|
|
||||||
"github.com/9seconds/mtg/mtproto"
|
"github.com/9seconds/mtg/mtproto"
|
||||||
"github.com/9seconds/mtg/mtproto/rpc"
|
"github.com/9seconds/mtg/mtproto/rpc"
|
||||||
)
|
)
|
||||||
|
|
||||||
type MTProtoProxy struct {
|
type MTProtoProxy struct {
|
||||||
conn PacketReadWriteCloser
|
conn PacketReadWriteCloser
|
||||||
req *rpc.ProxyRequest
|
req *rpc.ProxyRequest
|
||||||
|
logger *zap.SugaredLogger
|
||||||
|
|
||||||
readCounter uint32
|
readCounter uint32
|
||||||
writeCounter uint32
|
writeCounter uint32
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoProxy) Read() ([]byte, error) {
|
func (m *MTProtoProxy) Read() ([]byte, error) {
|
||||||
m.LogDebug("Read packet",
|
m.logger.Debugw("Read packet",
|
||||||
"counter", m.readCounter,
|
"counter", m.readCounter,
|
||||||
"simple_ack", m.req.Options.WriteHacks.SimpleAck,
|
"simple_ack", m.req.Options.WriteHacks.SimpleAck,
|
||||||
"quick_ack", m.req.Options.WriteHacks.QuickAck,
|
"quick_ack", m.req.Options.WriteHacks.QuickAck,
|
||||||
@@ -29,7 +31,7 @@ func (m *MTProtoProxy) Read() ([]byte, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, errors.Annotate(err, "Cannot read packet")
|
return nil, errors.Annotate(err, "Cannot read packet")
|
||||||
}
|
}
|
||||||
m.LogDebug("Read packet length",
|
m.logger.Debugw("Read packet length",
|
||||||
"counter", m.readCounter,
|
"counter", m.readCounter,
|
||||||
"simple_ack", m.req.Options.WriteHacks.SimpleAck,
|
"simple_ack", m.req.Options.WriteHacks.SimpleAck,
|
||||||
"quick_ack", m.req.Options.WriteHacks.QuickAck,
|
"quick_ack", m.req.Options.WriteHacks.QuickAck,
|
||||||
@@ -41,7 +43,7 @@ func (m *MTProtoProxy) Read() ([]byte, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
tag, packet := packet[:4], packet[4:]
|
tag, packet := packet[:4], packet[4:]
|
||||||
m.LogDebug("Read RPC tag",
|
m.logger.Debugw("Read RPC tag",
|
||||||
"counter", m.readCounter,
|
"counter", m.readCounter,
|
||||||
"simple_ack", m.req.Options.WriteHacks.SimpleAck,
|
"simple_ack", m.req.Options.WriteHacks.SimpleAck,
|
||||||
"quick_ack", m.req.Options.WriteHacks.QuickAck,
|
"quick_ack", m.req.Options.WriteHacks.QuickAck,
|
||||||
@@ -82,7 +84,7 @@ func (m *MTProtoProxy) readCloseExt(data []byte) ([]byte, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoProxy) Write(p []byte) (int, error) {
|
func (m *MTProtoProxy) Write(p []byte) (int, error) {
|
||||||
m.LogDebug("Write packet",
|
m.logger.Debugw("Write packet",
|
||||||
"length", len(p),
|
"length", len(p),
|
||||||
"counter", m.writeCounter,
|
"counter", m.writeCounter,
|
||||||
"simple_ack", m.req.Options.ReadHacks.SimpleAck,
|
"simple_ack", m.req.Options.ReadHacks.SimpleAck,
|
||||||
@@ -97,24 +99,8 @@ func (m *MTProtoProxy) Write(p []byte) (int, error) {
|
|||||||
return len(p), nil
|
return len(p), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoProxy) LogDebug(msg string, data ...interface{}) {
|
func (m *MTProtoProxy) Logger() *zap.SugaredLogger {
|
||||||
data = append(data, []interface{}{"type", "proxy"}...)
|
return m.logger
|
||||||
m.conn.LogDebug(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoProxy) LogInfo(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "proxy"}...)
|
|
||||||
m.conn.LogInfo(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoProxy) LogWarn(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "proxy"}...)
|
|
||||||
m.conn.LogWarn(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (m *MTProtoProxy) LogError(msg string, data ...interface{}) {
|
|
||||||
data = append(data, []interface{}{"type", "proxy"}...)
|
|
||||||
m.conn.LogError(msg, data...)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MTProtoProxy) LocalAddr() *net.TCPAddr {
|
func (m *MTProtoProxy) LocalAddr() *net.TCPAddr {
|
||||||
@@ -136,7 +122,8 @@ func NewMTProtoProxy(conn PacketReadWriteCloser, connOpts *mtproto.ConnectionOpt
|
|||||||
}
|
}
|
||||||
|
|
||||||
return &MTProtoProxy{
|
return &MTProtoProxy{
|
||||||
conn: conn,
|
conn: conn,
|
||||||
req: req,
|
logger: conn.Logger().Named("mtproto-proxy"),
|
||||||
|
req: req,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,12 +5,14 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
|
|
||||||
"github.com/juju/errors"
|
"github.com/juju/errors"
|
||||||
|
"go.uber.org/zap"
|
||||||
)
|
)
|
||||||
|
|
||||||
type StreamCipher struct {
|
type StreamCipher struct {
|
||||||
encryptor cipher.Stream
|
encryptor cipher.Stream
|
||||||
decryptor cipher.Stream
|
decryptor cipher.Stream
|
||||||
conn StreamReadWriteCloser
|
conn StreamReadWriteCloser
|
||||||
|
logger *zap.SugaredLogger
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *StreamCipher) Read(p []byte) (int, error) {
|
func (s *StreamCipher) Read(p []byte) (int, error) {
|
||||||
@@ -30,20 +32,8 @@ func (s *StreamCipher) Write(p []byte) (int, error) {
|
|||||||
return s.conn.Write(encrypted)
|
return s.conn.Write(encrypted)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *StreamCipher) LogDebug(msg string, data ...interface{}) {
|
func (s *StreamCipher) Logger() *zap.SugaredLogger {
|
||||||
s.conn.LogDebug(msg, data...)
|
return s.logger
|
||||||
}
|
|
||||||
|
|
||||||
func (s *StreamCipher) LogInfo(msg string, data ...interface{}) {
|
|
||||||
s.conn.LogInfo(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (s *StreamCipher) LogWarn(msg string, data ...interface{}) {
|
|
||||||
s.conn.LogWarn(msg, data...)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (s *StreamCipher) LogError(msg string, data ...interface{}) {
|
|
||||||
s.conn.LogError(msg, data...)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *StreamCipher) LocalAddr() *net.TCPAddr {
|
func (s *StreamCipher) LocalAddr() *net.TCPAddr {
|
||||||
@@ -61,6 +51,7 @@ func (s *StreamCipher) Close() error {
|
|||||||
func NewStreamCipher(conn StreamReadWriteCloser, encryptor, decryptor cipher.Stream) StreamReadWriteCloser {
|
func NewStreamCipher(conn StreamReadWriteCloser, encryptor, decryptor cipher.Stream) StreamReadWriteCloser {
|
||||||
return &StreamCipher{
|
return &StreamCipher{
|
||||||
conn: conn,
|
conn: conn,
|
||||||
|
logger: conn.Logger().Named("stream-cipher"),
|
||||||
encryptor: encryptor,
|
encryptor: encryptor,
|
||||||
decryptor: decryptor,
|
decryptor: decryptor,
|
||||||
}
|
}
|
||||||
|
|||||||
+3
-5
@@ -3,14 +3,12 @@ package wrappers
|
|||||||
import (
|
import (
|
||||||
"io"
|
"io"
|
||||||
"net"
|
"net"
|
||||||
|
|
||||||
|
"go.uber.org/zap"
|
||||||
)
|
)
|
||||||
|
|
||||||
type Wrap interface {
|
type Wrap interface {
|
||||||
LogDebug(msg string, data ...interface{})
|
Logger() *zap.SugaredLogger
|
||||||
LogInfo(msg string, data ...interface{})
|
|
||||||
LogWarn(msg string, data ...interface{})
|
|
||||||
LogError(msg string, data ...interface{})
|
|
||||||
|
|
||||||
LocalAddr() *net.TCPAddr
|
LocalAddr() *net.TCPAddr
|
||||||
RemoteAddr() *net.TCPAddr
|
RemoteAddr() *net.TCPAddr
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user