mirror of
https://github.com/ScuroNeko/mtg.git
synced 2026-08-31 18:54:02 +03:00
Refactoring
This commit is contained in:
+11
-47
@@ -2,7 +2,7 @@ package server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"strconv"
|
||||
"sync"
|
||||
@@ -13,6 +13,8 @@ import (
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
const bufferSize = 4096
|
||||
|
||||
type Server struct {
|
||||
ip net.IP
|
||||
port int
|
||||
@@ -79,42 +81,8 @@ func (s *Server) accept(conn net.Conn) {
|
||||
|
||||
wait := &sync.WaitGroup{}
|
||||
wait.Add(2)
|
||||
go func() {
|
||||
defer wait.Done()
|
||||
buf := make([]byte, 128)
|
||||
|
||||
for {
|
||||
fmt.Println("client loop")
|
||||
n, err := clientConn.Read(buf)
|
||||
if err != nil {
|
||||
fmt.Println("client read error", err)
|
||||
return
|
||||
}
|
||||
_, err = tgConn.Write(buf[:n])
|
||||
if err != nil {
|
||||
fmt.Println("tgConn write error", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
go func() {
|
||||
defer wait.Done()
|
||||
buf := make([]byte, 128)
|
||||
|
||||
for {
|
||||
fmt.Println("tg loop")
|
||||
n, err := tgConn.Read(buf)
|
||||
if err != nil {
|
||||
fmt.Println("tgConn read error", err)
|
||||
return
|
||||
}
|
||||
_, err = clientConn.Write(buf[:n])
|
||||
if err != nil {
|
||||
fmt.Println("client write error", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
go s.pipe(wait, clientConn, tgConn)
|
||||
go s.pipe(wait, tgConn, clientConn)
|
||||
<-ctx.Done()
|
||||
wait.Wait()
|
||||
|
||||
@@ -143,12 +111,6 @@ func (s *Server) getClientStream(conn net.Conn, ctx context.Context, cancel cont
|
||||
wConn := newCipherReadWriteCloser(conn, obfs2)
|
||||
|
||||
return wConn, dc, nil
|
||||
|
||||
// cipherConn := newCipherReadWriteCloser(conn, obfs2)
|
||||
// ctxConn := newCtxReadWriteCloser(cipherConn, ctx, cancel)
|
||||
// // logConn := newLogReadWriteCloser(ctxConn, s.logger, socketID, "client")
|
||||
|
||||
// return ctxConn, dc, nil
|
||||
}
|
||||
|
||||
func (s *Server) getTelegramStream(dc int16, ctx context.Context, cancel context.CancelFunc, socketID string) (*CipherReadWriteCloser, error) {
|
||||
@@ -165,12 +127,14 @@ func (s *Server) getTelegramStream(dc int16, ctx context.Context, cancel context
|
||||
wConn := newCipherReadWriteCloser(socket, obfs2)
|
||||
|
||||
return wConn, nil
|
||||
}
|
||||
|
||||
// cipherConn := newCipherReadWriteCloser(socket, obfs2)
|
||||
// ctxConn := newCtxReadWriteCloser(cipherConn, ctx, cancel)
|
||||
// logConn := newLogReadWriteCloser(ctxConn, s.logger, socketID, "telegram")
|
||||
func (s *Server) pipe(wait *sync.WaitGroup, reader io.Reader, writer io.Writer) {
|
||||
defer wait.Done()
|
||||
|
||||
buf := make([]byte, bufferSize)
|
||||
io.CopyBuffer(writer, reader, buf)
|
||||
|
||||
// return logConn, nil
|
||||
}
|
||||
|
||||
func NewServer(ip net.IP, port int, secret []byte, logger *zap.SugaredLogger) *Server {
|
||||
|
||||
Reference in New Issue
Block a user