mirror of
https://github.com/xtaci/kcptun.git
synced 2024-04-21 12:32:32 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
09aea3056a | ||
|
|
0788aa8260 | ||
|
|
2f479b788c | ||
|
|
f3f78460a4 |
+25
-7
@@ -12,6 +12,7 @@ import (
|
||||
"golang.org/x/crypto/pbkdf2"
|
||||
|
||||
"github.com/golang/snappy"
|
||||
"github.com/pkg/errors"
|
||||
"github.com/urfave/cli"
|
||||
kcp "github.com/xtaci/kcp-go"
|
||||
"github.com/xtaci/smux"
|
||||
@@ -301,9 +302,11 @@ func main() {
|
||||
smuxConfig := smux.DefaultConfig()
|
||||
smuxConfig.MaxReceiveBuffer = config.SockBuf
|
||||
|
||||
createConn := func() *smux.Session {
|
||||
createConn := func() (*smux.Session, error) {
|
||||
kcpconn, err := kcp.DialWithOptions(config.RemoteAddr, block, config.DataShard, config.ParityShard)
|
||||
checkError(err)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "createConn()")
|
||||
}
|
||||
kcpconn.SetStreamMode(true)
|
||||
kcpconn.SetNoDelay(config.NoDelay, config.Interval, config.Resend, config.NoCongestion)
|
||||
kcpconn.SetWindowSize(config.SndWnd, config.RcvWnd)
|
||||
@@ -328,8 +331,21 @@ func main() {
|
||||
} else {
|
||||
session, err = smux.Client(newCompStream(kcpconn), smuxConfig)
|
||||
}
|
||||
checkError(err)
|
||||
return session
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "createConn()")
|
||||
}
|
||||
return session, nil
|
||||
}
|
||||
|
||||
// wait until a connection is ready
|
||||
waitConn := func() *smux.Session {
|
||||
for {
|
||||
if session, err := createConn(); err == nil {
|
||||
return session
|
||||
} else {
|
||||
time.Sleep(time.Second)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
numconn := uint16(config.Conn)
|
||||
@@ -339,7 +355,9 @@ func main() {
|
||||
}, numconn)
|
||||
|
||||
for k := range muxes {
|
||||
muxes[k].session = createConn()
|
||||
sess, err := createConn()
|
||||
checkError(err)
|
||||
muxes[k].session = sess
|
||||
muxes[k].ttl = time.Now().Add(time.Duration(config.AutoExpire) * time.Second)
|
||||
}
|
||||
|
||||
@@ -361,7 +379,7 @@ func main() {
|
||||
// do auto expiration
|
||||
if config.AutoExpire > 0 && time.Now().After(muxes[idx].ttl) {
|
||||
chScavenger <- muxes[idx].session
|
||||
muxes[idx].session = createConn()
|
||||
muxes[idx].session = waitConn()
|
||||
muxes[idx].ttl = time.Now().Add(time.Duration(config.AutoExpire) * time.Second)
|
||||
}
|
||||
|
||||
@@ -369,7 +387,7 @@ func main() {
|
||||
p2, err := muxes[idx].session.OpenStream()
|
||||
if err != nil { // mux failure
|
||||
chScavenger <- muxes[idx].session
|
||||
muxes[idx].session = createConn()
|
||||
muxes[idx].session = waitConn()
|
||||
muxes[idx].ttl = time.Now().Add(time.Duration(config.AutoExpire) * time.Second)
|
||||
goto OPEN_P2
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user