Compare commits

...
5 Commits
4 changed files with 34 additions and 22 deletions
+3 -6
View File
@@ -5,6 +5,7 @@ import (
"os"
)
// Config for client
type Config struct {
LocalAddr string `json:"localaddr"`
RemoteAddr string `json:"remoteaddr"`
@@ -29,16 +30,12 @@ type Config struct {
KeepAlive int `json:"keepalive"`
}
func parseJsonConfig(config *Config, path string) error {
func parseJSONConfig(config *Config, path string) error {
file, err := os.Open(path) // For read access.
if err != nil {
return err
}
defer file.Close()
if err = json.NewDecoder(file).Decode(config); err != nil {
return err
}
return err
return json.NewDecoder(file).Decode(config)
}
+27 -9
View File
@@ -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"
@@ -232,7 +233,7 @@ func main() {
config.KeepAlive = c.Int("keepalive")
if c.String("c") != "" {
err := parseJsonConfig(&config, c.String("c"))
err := parseJSONConfig(&config, c.String("c"))
checkError(err)
}
@@ -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
}
@@ -401,7 +419,7 @@ func scavenger(ch chan *smux.Session) {
var newList []scavengeSession
for k := range sessionList {
s := sessionList[k]
if s.session.NumStreams() == 0 || s.session.IsClosed() || time.Now().Sub(s.ttl) > maxScavengeTTL {
if s.session.NumStreams() == 0 || s.session.IsClosed() || time.Since(s.ttl) > maxScavengeTTL {
log.Println("session scavenged")
s.session.Close()
} else {
+3 -6
View File
@@ -5,6 +5,7 @@ import (
"os"
)
// Config for server
type Config struct {
Listen string `json:"listen"`
Target string `json:"target"`
@@ -27,16 +28,12 @@ type Config struct {
KeepAlive int `json:"keepalive"`
}
func parseJsonConfig(config *Config, path string) error {
func parseJSONConfig(config *Config, path string) error {
file, err := os.Open(path) // For read access.
if err != nil {
return err
}
defer file.Close()
if err = json.NewDecoder(file).Decode(config); err != nil {
return err
}
return err
return json.NewDecoder(file).Decode(config)
}
+1 -1
View File
@@ -254,7 +254,7 @@ func main() {
if c.String("c") != "" {
//Now only support json config file
err := parseJsonConfig(&config, c.String("c"))
err := parseJSONConfig(&config, c.String("c"))
checkError(err)
}