mirror of
https://github.com/wweir/sower.git
synced 2024-04-21 12:42:15 +00:00
51 lines
974 B
Go
51 lines
974 B
Go
package util
|
|
|
|
import (
|
|
"io"
|
|
"net"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/pkg/errors"
|
|
)
|
|
|
|
func RelayTo(conn net.Conn, addr string) (dur time.Duration, err error) {
|
|
if _, _, err := net.SplitHostPort(addr); err != nil {
|
|
addr = net.JoinHostPort(addr, "80")
|
|
}
|
|
|
|
start := time.Now()
|
|
rc, err := net.DialTimeout("tcp", addr, 5*time.Second)
|
|
if err != nil {
|
|
return time.Since(start), errors.WithStack(err)
|
|
}
|
|
defer rc.Close()
|
|
|
|
Relay(conn, rc)
|
|
return time.Since(start), nil
|
|
}
|
|
|
|
func Relay(conn1, conn2 net.Conn) {
|
|
wg := &sync.WaitGroup{}
|
|
exitFlag := new(int32)
|
|
wg.Add(2)
|
|
go redirect(conn2, conn1, wg, exitFlag)
|
|
redirect(conn1, conn2, wg, exitFlag)
|
|
wg.Wait()
|
|
}
|
|
func redirect(dst, src net.Conn, wg *sync.WaitGroup, exitFlag *int32) {
|
|
|
|
// io.Copy(dst, io.TeeReader(src, os.Stdout))
|
|
io.Copy(dst, src)
|
|
|
|
if atomic.CompareAndSwapInt32(exitFlag, 0, 1) {
|
|
// wakeup blocked goroutine
|
|
now := time.Now()
|
|
src.SetDeadline(now)
|
|
dst.SetDeadline(now)
|
|
}
|
|
|
|
wg.Done()
|
|
}
|