Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6b8d3f7d47 | ||
|
|
26e0fdf605 | ||
|
|
3653ec19cd |
@@ -125,8 +125,7 @@ go build
|
|||||||
|
|
||||||
## 参与贡献
|
## 参与贡献
|
||||||
TODO或ISSUE里如果有你擅长的领域,或者你有特别好的主意,可以加入OpenP2P项目,贡献你的代码。待项目茁壮成长后,你们就是知名开源项目的主要代码贡献者,岂不快哉。
|
TODO或ISSUE里如果有你擅长的领域,或者你有特别好的主意,可以加入OpenP2P项目,贡献你的代码。待项目茁壮成长后,你们就是知名开源项目的主要代码贡献者,岂不快哉。
|
||||||
## 商业合作
|
|
||||||
它是一个中国人发起的项目,更懂国内网络环境,更懂用户需求,更好的企业级支持
|
|
||||||
## 技术交流
|
## 技术交流
|
||||||
QQ群:16947733
|
QQ群:16947733
|
||||||
邮箱:openp2p.cn@gmail.com tenderiron@139.com
|
邮箱:openp2p.cn@gmail.com tenderiron@139.com
|
||||||
|
|||||||
@@ -48,6 +48,7 @@ Download on local and remote computers and double-click to run, one-click instal
|
|||||||

|

|
||||||
|
|
||||||

|

|
||||||
|
|
||||||
### 3.New P2PApp
|
### 3.New P2PApp
|
||||||
|
|
||||||

|

|
||||||
|
|||||||
@@ -17,6 +17,12 @@
|
|||||||
>* -sharebandwidth: 作为共享节点时提供带宽,默认10mbps. 如果是光纤大带宽,设置越大效果越好. 0表示不共享,该节点只在私有的P2P网络使用。不加入共享的P2P网络,这样也意味着无法使用别人的共享节点
|
>* -sharebandwidth: 作为共享节点时提供带宽,默认10mbps. 如果是光纤大带宽,设置越大效果越好. 0表示不共享,该节点只在私有的P2P网络使用。不加入共享的P2P网络,这样也意味着无法使用别人的共享节点
|
||||||
>* -loglevel: 需要查看更多调试日志,设置0;默认是1
|
>* -loglevel: 需要查看更多调试日志,设置0;默认是1
|
||||||
|
|
||||||
|
### 在docker容器里运行openp2p
|
||||||
|
我们暂时还没提供官方docker镜像,你可以在随便一个容器里运行
|
||||||
|
```
|
||||||
|
nohup ./openp2p -d -node OFFICEPC1 -token TOKEN &
|
||||||
|
#这里由于一般的镜像都精简过,install系统服务会失败,所以使用直接daemon模式后台运行
|
||||||
|
```
|
||||||
## 连接
|
## 连接
|
||||||
```
|
```
|
||||||
./openp2p -d -node HOMEPC123 -token TOKEN -appname OfficeWindowsRemote -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389
|
./openp2p -d -node HOMEPC123 -token TOKEN -appname OfficeWindowsRemote -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389
|
||||||
|
|||||||
@@ -19,6 +19,13 @@ Or
|
|||||||
>* -sharebandwidth: Provides bandwidth when used as a shared node, the default is 10mbps. If it is a large bandwidth of optical fiber, the larger the setting, the better the effect. 0 means not shared, the node is only used in a private P2P network. Do not join the shared P2P network, which also means that you CAN NOT use other people’s shared nodes
|
>* -sharebandwidth: Provides bandwidth when used as a shared node, the default is 10mbps. If it is a large bandwidth of optical fiber, the larger the setting, the better the effect. 0 means not shared, the node is only used in a private P2P network. Do not join the shared P2P network, which also means that you CAN NOT use other people’s shared nodes
|
||||||
>* -loglevel: Need to view more debug logs, set 0; the default is 1
|
>* -loglevel: Need to view more debug logs, set 0; the default is 1
|
||||||
|
|
||||||
|
### Run in Docker container
|
||||||
|
We don't provide official docker image yet, you can run it in any container
|
||||||
|
```
|
||||||
|
nohup ./openp2p -d -node OFFICEPC1 -token TOKEN &
|
||||||
|
# Since many docker images have been simplified, the install system service will fail, so the daemon mode is used to run in the background
|
||||||
|
```
|
||||||
|
|
||||||
## Connect
|
## Connect
|
||||||
```
|
```
|
||||||
./openp2p -d -node HOMEPC123 -token TOKEN -appname OfficeWindowsRemote -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389
|
./openp2p -d -node HOMEPC123 -token TOKEN -appname OfficeWindowsRemote -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389
|
||||||
|
|||||||
@@ -4,7 +4,6 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"flag"
|
"flag"
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
"os"
|
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
@@ -30,15 +29,18 @@ type AppConfig struct {
|
|||||||
peerConeNatPort int
|
peerConeNatPort int
|
||||||
retryNum int
|
retryNum int
|
||||||
retryTime time.Time
|
retryTime time.Time
|
||||||
|
nextRetryTime time.Time
|
||||||
shareBandwidth int
|
shareBandwidth int
|
||||||
|
errMsg string
|
||||||
|
connectTime time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
// TODO: add loglevel, maxlogfilesize
|
// TODO: add loglevel, maxlogfilesize
|
||||||
type Config struct {
|
type Config struct {
|
||||||
Network NetworkConfig `json:"network"`
|
Network NetworkConfig `json:"network"`
|
||||||
Apps []AppConfig `json:"apps"`
|
Apps []*AppConfig `json:"apps"`
|
||||||
LogLevel int
|
LogLevel int
|
||||||
|
daemonMode bool
|
||||||
mtx sync.Mutex
|
mtx sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -48,27 +50,29 @@ func (c *Config) switchApp(app AppConfig, enabled int) {
|
|||||||
for i := 0; i < len(c.Apps); i++ {
|
for i := 0; i < len(c.Apps); i++ {
|
||||||
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
||||||
c.Apps[i].Enabled = enabled
|
c.Apps[i].Enabled = enabled
|
||||||
|
c.Apps[i].retryNum = 0
|
||||||
|
c.Apps[i].nextRetryTime = time.Now()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) add(app AppConfig, force bool) {
|
func (c *Config) add(app AppConfig, override bool) {
|
||||||
c.mtx.Lock()
|
c.mtx.Lock()
|
||||||
defer c.mtx.Unlock()
|
defer c.mtx.Unlock()
|
||||||
if app.SrcPort == 0 || app.DstPort == 0 {
|
if app.SrcPort == 0 || app.DstPort == 0 {
|
||||||
gLog.Println(LevelERROR, "invalid app ", app)
|
gLog.Println(LevelERROR, "invalid app ", app)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
if override {
|
||||||
for i := 0; i < len(c.Apps); i++ {
|
for i := 0; i < len(c.Apps); i++ {
|
||||||
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
||||||
if force {
|
c.Apps[i] = &app // override it
|
||||||
c.Apps[i] = app
|
|
||||||
}
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
c.Apps = append(c.Apps, app)
|
}
|
||||||
|
c.Apps = append(c.Apps, &app)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) delete(app AppConfig) {
|
func (c *Config) delete(app AppConfig) {
|
||||||
@@ -112,6 +116,27 @@ func (c *Config) load() error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c *Config) setToken(token uint64) {
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
|
c.Network.Token = token
|
||||||
|
}
|
||||||
|
func (c *Config) setUser(user string) {
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
|
c.Network.User = user
|
||||||
|
}
|
||||||
|
func (c *Config) setNode(node string) {
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
|
c.Network.Node = node
|
||||||
|
}
|
||||||
|
func (c *Config) setShareBandwidth(bw int) {
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
|
c.Network.ShareBandwidth = bw
|
||||||
|
}
|
||||||
|
|
||||||
type NetworkConfig struct {
|
type NetworkConfig struct {
|
||||||
// local info
|
// local info
|
||||||
Token uint64
|
Token uint64
|
||||||
@@ -142,11 +167,9 @@ func parseParams() {
|
|||||||
srcPort := flag.Int("srcport", 0, "source port ")
|
srcPort := flag.Int("srcport", 0, "source port ")
|
||||||
protocol := flag.String("protocol", "tcp", "tcp or udp")
|
protocol := flag.String("protocol", "tcp", "tcp or udp")
|
||||||
appName := flag.String("appname", "", "app name")
|
appName := flag.String("appname", "", "app name")
|
||||||
flag.Bool("noshare", false, "deprecated. uses -sharebandwidth 0") // Deprecated, rm later
|
shareBandwidth := flag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private network no limit")
|
||||||
shareBandwidth := flag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private node no limit")
|
|
||||||
flag.Bool("f", false, "deprecated. config file") // Deprecated, rm later
|
|
||||||
daemonMode := flag.Bool("d", false, "daemonMode")
|
daemonMode := flag.Bool("d", false, "daemonMode")
|
||||||
flag.Bool("bydaemon", false, "start by daemon") // Deprecated, rm later
|
notVerbose := flag.Bool("nv", false, "not log console")
|
||||||
logLevel := flag.Int("loglevel", 1, "0:debug 1:info 2:warn 3:error")
|
logLevel := flag.Int("loglevel", 1, "0:debug 1:info 2:warn 3:error")
|
||||||
flag.Parse()
|
flag.Parse()
|
||||||
|
|
||||||
@@ -161,8 +184,8 @@ func parseParams() {
|
|||||||
if config.SrcPort != 0 {
|
if config.SrcPort != 0 {
|
||||||
gConf.add(config, true)
|
gConf.add(config, true)
|
||||||
}
|
}
|
||||||
gConf.mtx.Lock()
|
// gConf.mtx.Lock() // when calling this func it's single-thread no lock
|
||||||
|
gConf.daemonMode = *daemonMode
|
||||||
// spec paramters in commandline will always be used
|
// spec paramters in commandline will always be used
|
||||||
flag.Visit(func(f *flag.Flag) {
|
flag.Visit(func(f *flag.Flag) {
|
||||||
if f.Name == "sharebandwidth" {
|
if f.Name == "sharebandwidth" {
|
||||||
@@ -203,11 +226,9 @@ func parseParams() {
|
|||||||
gConf.Network.UDPPort1 = 27182
|
gConf.Network.UDPPort1 = 27182
|
||||||
gConf.Network.UDPPort2 = 27183
|
gConf.Network.UDPPort2 = 27183
|
||||||
gLog.setLevel(LogLevel(gConf.LogLevel))
|
gLog.setLevel(LogLevel(gConf.LogLevel))
|
||||||
gConf.mtx.Unlock()
|
if *notVerbose {
|
||||||
|
gLog.setMode(LogFile)
|
||||||
|
}
|
||||||
|
// gConf.mtx.Unlock()
|
||||||
gConf.save()
|
gConf.save()
|
||||||
if *daemonMode {
|
|
||||||
d := daemon{}
|
|
||||||
d.run()
|
|
||||||
os.Exit(0)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -64,10 +64,18 @@ func (d *daemon) run() {
|
|||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
args = append(args, "-nv")
|
||||||
for {
|
for {
|
||||||
// start worker
|
// start worker
|
||||||
|
tmpDump := filepath.Join("log", "dump.log.tmp")
|
||||||
|
dumpFile := filepath.Join("log", "dump.log")
|
||||||
|
f, err := os.Create(filepath.Join(tmpDump))
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "start worker error:%s", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
gLog.Println(LevelINFO, "start worker process, args:", args)
|
gLog.Println(LevelINFO, "start worker process, args:", args)
|
||||||
execSpec := &os.ProcAttr{Files: []*os.File{os.Stdin, os.Stdout, os.Stderr}}
|
execSpec := &os.ProcAttr{Env: append(os.Environ(), "GOTRACEBACK=crash"), Files: []*os.File{os.Stdin, os.Stdout, f}}
|
||||||
p, err := os.StartProcess(binPath, args, execSpec)
|
p, err := os.StartProcess(binPath, args, execSpec)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Printf(LevelERROR, "start worker error:%s", err)
|
gLog.Printf(LevelERROR, "start worker error:%s", err)
|
||||||
@@ -75,6 +83,12 @@ func (d *daemon) run() {
|
|||||||
}
|
}
|
||||||
d.proc = p
|
d.proc = p
|
||||||
_, _ = p.Wait()
|
_, _ = p.Wait()
|
||||||
|
f.Close()
|
||||||
|
time.Sleep(time.Second)
|
||||||
|
err = os.Rename(tmpDump, dumpFile)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "rename dump error:%s", err)
|
||||||
|
}
|
||||||
if !d.running {
|
if !d.running {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -110,9 +124,22 @@ func (d *daemon) Control(ctrlComm string, exeAbsPath string, args []string) erro
|
|||||||
// listen and build p2papp:
|
// listen and build p2papp:
|
||||||
// ./openp2p install -node hhd1207-222 -token YOUR-TOKEN -sharebandwidth 0 -peernode hhdhome-n1 -dstip 127.0.0.1 -dstport 50022 -protocol tcp -srcport 22
|
// ./openp2p install -node hhd1207-222 -token YOUR-TOKEN -sharebandwidth 0 -peernode hhdhome-n1 -dstip 127.0.0.1 -dstport 50022 -protocol tcp -srcport 22
|
||||||
func install() {
|
func install() {
|
||||||
|
gLog.Println(LevelINFO, "openp2p start. version: ", OpenP2PVersion)
|
||||||
|
gLog.Println(LevelINFO, "Contact: QQ Group: 16947733, Email: [email protected]")
|
||||||
gLog.Println(LevelINFO, "install start")
|
gLog.Println(LevelINFO, "install start")
|
||||||
defer gLog.Println(LevelINFO, "install end")
|
defer gLog.Println(LevelINFO, "install end")
|
||||||
// auto uninstall
|
// auto uninstall
|
||||||
|
err := os.MkdirAll(defaultInstallPath, 0775)
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "MkdirAll %s error:%s", defaultInstallPath, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
err = os.Chdir(defaultInstallPath)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Println(LevelERROR, "cd error:", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
uninstall()
|
uninstall()
|
||||||
// save config file
|
// save config file
|
||||||
@@ -127,23 +154,25 @@ func install() {
|
|||||||
srcPort := installFlag.Int("srcport", 0, "source port ")
|
srcPort := installFlag.Int("srcport", 0, "source port ")
|
||||||
protocol := installFlag.String("protocol", "tcp", "tcp or udp")
|
protocol := installFlag.String("protocol", "tcp", "tcp or udp")
|
||||||
appName := flag.String("appname", "", "app name")
|
appName := flag.String("appname", "", "app name")
|
||||||
installFlag.Bool("noshare", false, "deprecated. uses -sharebandwidth 0")
|
shareBandwidth := installFlag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private network no limit")
|
||||||
shareBandwidth := installFlag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private node no limit")
|
|
||||||
logLevel := installFlag.Int("loglevel", 1, "0:debug 1:info 2:warn 3:error")
|
logLevel := installFlag.Int("loglevel", 1, "0:debug 1:info 2:warn 3:error")
|
||||||
installFlag.Parse(os.Args[2:])
|
installFlag.Parse(os.Args[2:])
|
||||||
if *node != "" && len(*node) < 8 {
|
|
||||||
gLog.Println(LevelERROR, ErrNodeTooShort)
|
|
||||||
os.Exit(9)
|
|
||||||
}
|
|
||||||
if *node == "" { // if node name not set. use os.Hostname
|
|
||||||
hostname := defaultNodeName()
|
|
||||||
node = &hostname
|
|
||||||
}
|
|
||||||
gConf.load() // load old config. otherwise will clear all apps
|
gConf.load() // load old config. otherwise will clear all apps
|
||||||
gConf.LogLevel = *logLevel
|
gConf.LogLevel = *logLevel
|
||||||
gConf.Network.ServerHost = *serverHost
|
gConf.Network.ServerHost = *serverHost
|
||||||
gConf.Network.Token = *token
|
gConf.Network.Token = *token
|
||||||
|
if *node != "" {
|
||||||
|
if len(*node) < 8 {
|
||||||
|
gLog.Println(LevelERROR, ErrNodeTooShort)
|
||||||
|
os.Exit(9)
|
||||||
|
}
|
||||||
gConf.Network.Node = *node
|
gConf.Network.Node = *node
|
||||||
|
} else {
|
||||||
|
if gConf.Network.Node == "" { // if node name not set. use os.Hostname
|
||||||
|
gConf.Network.Node = defaultNodeName()
|
||||||
|
}
|
||||||
|
}
|
||||||
gConf.Network.ServerPort = 27183
|
gConf.Network.ServerPort = 27183
|
||||||
gConf.Network.UDPPort1 = 27182
|
gConf.Network.UDPPort1 = 27182
|
||||||
gConf.Network.UDPPort2 = 27183
|
gConf.Network.UDPPort2 = 27183
|
||||||
@@ -158,16 +187,6 @@ func install() {
|
|||||||
if config.SrcPort != 0 {
|
if config.SrcPort != 0 {
|
||||||
gConf.add(config, true)
|
gConf.add(config, true)
|
||||||
}
|
}
|
||||||
err := os.MkdirAll(defaultInstallPath, 0775)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Printf(LevelERROR, "MkdirAll %s error:%s", defaultInstallPath, err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
err = os.Chdir(defaultInstallPath)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Println(LevelERROR, "cd error:", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
gConf.save()
|
gConf.save()
|
||||||
targetPath := filepath.Join(defaultInstallPath, defaultBinName)
|
targetPath := filepath.Join(defaultInstallPath, defaultBinName)
|
||||||
d := daemon{}
|
d := daemon{}
|
||||||
@@ -195,7 +214,6 @@ func install() {
|
|||||||
dst.Close()
|
dst.Close()
|
||||||
|
|
||||||
// install system service
|
// install system service
|
||||||
// args := []string{""}
|
|
||||||
gLog.Println(LevelINFO, "targetPath:", targetPath)
|
gLog.Println(LevelINFO, "targetPath:", targetPath)
|
||||||
err = d.Control("install", targetPath, []string{"-d"})
|
err = d.Control("install", targetPath, []string{"-d"})
|
||||||
if err == nil {
|
if err == nil {
|
||||||
|
|||||||
|
Before Width: | Height: | Size: 65 KiB After Width: | Height: | Size: 19 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 8.4 KiB |
|
Before Width: | Height: | Size: 50 KiB After Width: | Height: | Size: 15 KiB |
|
Before Width: | Height: | Size: 14 KiB After Width: | Height: | Size: 6.8 KiB |
@@ -110,30 +110,36 @@ func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
|
|||||||
case MsgPushReportApps:
|
case MsgPushReportApps:
|
||||||
gLog.Println(LevelINFO, "MsgPushReportApps")
|
gLog.Println(LevelINFO, "MsgPushReportApps")
|
||||||
req := ReportApps{}
|
req := ReportApps{}
|
||||||
// TODO: add the retrying apps
|
|
||||||
gConf.mtx.Lock()
|
gConf.mtx.Lock()
|
||||||
defer gConf.mtx.Unlock()
|
defer gConf.mtx.Unlock()
|
||||||
for _, config := range gConf.Apps {
|
for _, config := range gConf.Apps {
|
||||||
appActive := 0
|
appActive := 0
|
||||||
|
relayNode := ""
|
||||||
|
relayMode := ""
|
||||||
i, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
i, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
if ok {
|
if ok {
|
||||||
app := i.(*p2pApp)
|
app := i.(*p2pApp)
|
||||||
if app.isActive() {
|
if app.isActive() {
|
||||||
appActive = 1
|
appActive = 1
|
||||||
}
|
}
|
||||||
|
relayNode = app.relayNode
|
||||||
|
relayMode = app.relayMode
|
||||||
}
|
}
|
||||||
appInfo := AppInfo{
|
appInfo := AppInfo{
|
||||||
AppName: config.AppName,
|
AppName: config.AppName,
|
||||||
|
Error: config.errMsg,
|
||||||
Protocol: config.Protocol,
|
Protocol: config.Protocol,
|
||||||
SrcPort: config.SrcPort,
|
SrcPort: config.SrcPort,
|
||||||
// RelayNode: relayNode,
|
RelayNode: relayNode,
|
||||||
|
RelayMode: relayMode,
|
||||||
PeerNode: config.PeerNode,
|
PeerNode: config.PeerNode,
|
||||||
DstHost: config.DstHost,
|
DstHost: config.DstHost,
|
||||||
DstPort: config.DstPort,
|
DstPort: config.DstPort,
|
||||||
PeerUser: config.PeerUser,
|
PeerUser: config.PeerUser,
|
||||||
PeerIP: config.peerIP,
|
PeerIP: config.peerIP,
|
||||||
PeerNatType: config.peerNatType,
|
PeerNatType: config.peerNatType,
|
||||||
RetryTime: config.retryTime.String(),
|
RetryTime: config.retryTime.Local().Format("2006-01-02T15:04:05-0700"),
|
||||||
|
ConnectTime: config.connectTime.Local().Format("2006-01-02T15:04:05-0700"),
|
||||||
IsActive: appActive,
|
IsActive: appActive,
|
||||||
Enabled: config.Enabled,
|
Enabled: config.Enabled,
|
||||||
}
|
}
|
||||||
@@ -176,10 +182,8 @@ func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
|
|||||||
gLog.Printf(LevelERROR, "wrong MsgPushEditNode:%s %s", err, string(msg[openP2PHeaderSize:]))
|
gLog.Printf(LevelERROR, "wrong MsgPushEditNode:%s %s", err, string(msg[openP2PHeaderSize:]))
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
gConf.mtx.Lock()
|
gConf.setNode(req.NewName)
|
||||||
gConf.Network.Node = req.NewName
|
gConf.setShareBandwidth(req.Bandwidth)
|
||||||
gConf.Network.ShareBandwidth = req.Bandwidth
|
|
||||||
gConf.mtx.Unlock()
|
|
||||||
gConf.save()
|
gConf.save()
|
||||||
// TODO: hot reload
|
// TODO: hot reload
|
||||||
os.Exit(0)
|
os.Exit(0)
|
||||||
|
|||||||
@@ -97,6 +97,11 @@ func (vl *V8log) setLevel(level LogLevel) {
|
|||||||
defer vl.mtx.Unlock()
|
defer vl.mtx.Unlock()
|
||||||
vl.level = level
|
vl.level = level
|
||||||
}
|
}
|
||||||
|
func (vl *V8log) setMode(mode int) {
|
||||||
|
vl.mtx.Lock()
|
||||||
|
defer vl.mtx.Unlock()
|
||||||
|
vl.mode = mode
|
||||||
|
}
|
||||||
|
|
||||||
func (vl *V8log) checkFile() {
|
func (vl *V8log) checkFile() {
|
||||||
if vl.maxLogSize <= 0 {
|
if vl.maxLogSize <= 0 {
|
||||||
@@ -110,10 +115,10 @@ func (vl *V8log) checkFile() {
|
|||||||
for l, logFile := range vl.files {
|
for l, logFile := range vl.files {
|
||||||
f, e := logFile.Stat()
|
f, e := logFile.Stat()
|
||||||
if e != nil {
|
if e != nil {
|
||||||
break
|
continue
|
||||||
}
|
}
|
||||||
if f.Size() <= vl.maxLogSize {
|
if f.Size() <= vl.maxLogSize {
|
||||||
break
|
continue
|
||||||
}
|
}
|
||||||
logFile.Close()
|
logFile.Close()
|
||||||
fname := f.Name()
|
fname := f.Name()
|
||||||
|
|||||||
@@ -13,7 +13,6 @@ func main() {
|
|||||||
binDir := filepath.Dir(os.Args[0])
|
binDir := filepath.Dir(os.Args[0])
|
||||||
os.Chdir(binDir) // for system service
|
os.Chdir(binDir) // for system service
|
||||||
gLog = InitLogger(binDir, "openp2p", LevelDEBUG, 1024*1024, LogFileAndConsole)
|
gLog = InitLogger(binDir, "openp2p", LevelDEBUG, 1024*1024, LogFileAndConsole)
|
||||||
gLog.Println(LevelINFO, "openp2p start. version: ", OpenP2PVersion)
|
|
||||||
// TODO: install sub command, deamon process
|
// TODO: install sub command, deamon process
|
||||||
if len(os.Args) > 1 {
|
if len(os.Args) > 1 {
|
||||||
switch os.Args[1] {
|
switch os.Args[1] {
|
||||||
@@ -41,8 +40,16 @@ func main() {
|
|||||||
} else {
|
} else {
|
||||||
installByFilename()
|
installByFilename()
|
||||||
}
|
}
|
||||||
|
|
||||||
parseParams()
|
parseParams()
|
||||||
|
gLog.Println(LevelINFO, "openp2p start. version: ", OpenP2PVersion)
|
||||||
|
gLog.Println(LevelINFO, "Contact: QQ Group: 16947733, Email: [email protected]")
|
||||||
|
|
||||||
|
if gConf.daemonMode {
|
||||||
|
d := daemon{}
|
||||||
|
d.run()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
gLog.Println(LevelINFO, &gConf)
|
gLog.Println(LevelINFO, &gConf)
|
||||||
setFirewall()
|
setFirewall()
|
||||||
network := P2PNetworkInstance(&gConf.Network)
|
network := P2PNetworkInstance(&gConf.Network)
|
||||||
|
|||||||
@@ -0,0 +1,150 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"encoding/binary"
|
||||||
|
"errors"
|
||||||
|
"net"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
var ErrDeadlineExceeded error = &DeadlineExceededError{}
|
||||||
|
|
||||||
|
// DeadlineExceededError is returned for an expired deadline.
|
||||||
|
type DeadlineExceededError struct{}
|
||||||
|
|
||||||
|
// Implement the net.Error interface.
|
||||||
|
// The string is "i/o timeout" because that is what was returned
|
||||||
|
// by earlier Go versions. Changing it may break programs that
|
||||||
|
// match on error strings.
|
||||||
|
func (e *DeadlineExceededError) Error() string { return "i/o timeout" }
|
||||||
|
func (e *DeadlineExceededError) Timeout() bool { return true }
|
||||||
|
func (e *DeadlineExceededError) Temporary() bool { return true }
|
||||||
|
|
||||||
|
// implement io.Writer
|
||||||
|
type overlayConn struct {
|
||||||
|
tunnel *P2PTunnel
|
||||||
|
connTCP net.Conn
|
||||||
|
id uint64
|
||||||
|
rtid uint64
|
||||||
|
running bool
|
||||||
|
isClient bool
|
||||||
|
appID uint64
|
||||||
|
appKey uint64
|
||||||
|
appKeyBytes []byte
|
||||||
|
// for udp
|
||||||
|
connUDP *net.UDPConn
|
||||||
|
remoteAddr net.Addr
|
||||||
|
udpRelayData chan []byte
|
||||||
|
lastReadUDPTs time.Time
|
||||||
|
}
|
||||||
|
|
||||||
|
func (oConn *overlayConn) run() {
|
||||||
|
gLog.Printf(LevelDEBUG, "%d overlayConn run start", oConn.id)
|
||||||
|
defer gLog.Printf(LevelDEBUG, "%d overlayConn run end", oConn.id)
|
||||||
|
oConn.running = true
|
||||||
|
oConn.lastReadUDPTs = time.Now()
|
||||||
|
buffer := make([]byte, ReadBuffLen+PaddingSize)
|
||||||
|
readBuf := buffer[:ReadBuffLen]
|
||||||
|
encryptData := make([]byte, ReadBuffLen+PaddingSize) // 16 bytes for padding
|
||||||
|
tunnelHead := new(bytes.Buffer)
|
||||||
|
relayHead := new(bytes.Buffer)
|
||||||
|
binary.Write(relayHead, binary.LittleEndian, oConn.rtid)
|
||||||
|
binary.Write(tunnelHead, binary.LittleEndian, oConn.id)
|
||||||
|
for oConn.running && oConn.tunnel.isRuning() {
|
||||||
|
buff, dataLen, err := oConn.Read(readBuf)
|
||||||
|
if err != nil {
|
||||||
|
if ne, ok := err.(net.Error); ok && ne.Timeout() {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
// overlay tcp connection normal close, debug log
|
||||||
|
gLog.Printf(LevelDEBUG, "overlayConn %d read error:%s,close it", oConn.id, err)
|
||||||
|
break
|
||||||
|
}
|
||||||
|
payload := buff[:dataLen]
|
||||||
|
if oConn.appKey != 0 {
|
||||||
|
payload, _ = encryptBytes(oConn.appKeyBytes, encryptData, buffer[:dataLen], dataLen)
|
||||||
|
}
|
||||||
|
writeBytes := append(tunnelHead.Bytes(), payload...)
|
||||||
|
if oConn.rtid == 0 {
|
||||||
|
oConn.tunnel.conn.WriteBytes(MsgP2P, MsgOverlayData, writeBytes)
|
||||||
|
gLog.Printf(LevelDEBUG, "write overlay data to %d:%d bodylen=%d", oConn.rtid, oConn.id, len(writeBytes))
|
||||||
|
} else {
|
||||||
|
// write raley data
|
||||||
|
all := append(relayHead.Bytes(), encodeHeader(MsgP2P, MsgOverlayData, uint32(len(writeBytes)))...)
|
||||||
|
all = append(all, writeBytes...)
|
||||||
|
oConn.tunnel.conn.WriteBytes(MsgP2P, MsgRelayData, all)
|
||||||
|
gLog.Printf(LevelDEBUG, "write relay data to %d:%d bodylen=%d", oConn.rtid, oConn.id, len(writeBytes))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if oConn.connTCP != nil {
|
||||||
|
oConn.connTCP.Close()
|
||||||
|
}
|
||||||
|
if oConn.connUDP != nil {
|
||||||
|
oConn.connUDP.Close()
|
||||||
|
}
|
||||||
|
oConn.tunnel.overlayConns.Delete(oConn.id)
|
||||||
|
// notify peer disconnect
|
||||||
|
if oConn.isClient {
|
||||||
|
req := OverlayDisconnectReq{ID: oConn.id}
|
||||||
|
if oConn.rtid == 0 {
|
||||||
|
oConn.tunnel.conn.WriteMessage(MsgP2P, MsgOverlayDisconnectReq, &req)
|
||||||
|
} else {
|
||||||
|
// write relay data
|
||||||
|
msg, _ := newMessage(MsgP2P, MsgOverlayDisconnectReq, &req)
|
||||||
|
msgWithHead := append(relayHead.Bytes(), msg...)
|
||||||
|
oConn.tunnel.conn.WriteBytes(MsgP2P, MsgRelayData, msgWithHead)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (oConn *overlayConn) Read(reuseBuff []byte) (buff []byte, n int, err error) {
|
||||||
|
if oConn.connUDP != nil {
|
||||||
|
if time.Now().After(oConn.lastReadUDPTs.Add(time.Minute * 5)) {
|
||||||
|
err = errors.New("udp close")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if oConn.remoteAddr != nil { // as server
|
||||||
|
select {
|
||||||
|
case buff = <-oConn.udpRelayData:
|
||||||
|
n = len(buff)
|
||||||
|
oConn.lastReadUDPTs = time.Now()
|
||||||
|
case <-time.After(time.Second * 10):
|
||||||
|
err = ErrDeadlineExceeded
|
||||||
|
}
|
||||||
|
} else { // as client
|
||||||
|
oConn.connUDP.SetReadDeadline(time.Now().Add(5 * time.Second))
|
||||||
|
n, _, err = oConn.connUDP.ReadFrom(reuseBuff)
|
||||||
|
if err == nil {
|
||||||
|
oConn.lastReadUDPTs = time.Now()
|
||||||
|
}
|
||||||
|
buff = reuseBuff
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
oConn.connTCP.SetReadDeadline(time.Now().Add(time.Second * 5))
|
||||||
|
n, err = oConn.connTCP.Read(reuseBuff)
|
||||||
|
buff = reuseBuff
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// calling by p2pTunnel
|
||||||
|
func (oConn *overlayConn) Write(buff []byte) (n int, err error) {
|
||||||
|
// add mutex when multi-thread calling
|
||||||
|
if oConn.connUDP != nil {
|
||||||
|
if oConn.remoteAddr == nil {
|
||||||
|
n, err = oConn.connUDP.Write(buff)
|
||||||
|
} else {
|
||||||
|
n, err = oConn.connUDP.WriteTo(buff, oConn.remoteAddr)
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
oConn.running = false
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
n, err = oConn.connTCP.Write(buff)
|
||||||
|
if err != nil {
|
||||||
|
oConn.running = false
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
@@ -1,85 +0,0 @@
|
|||||||
package main
|
|
||||||
|
|
||||||
import (
|
|
||||||
"bytes"
|
|
||||||
"encoding/binary"
|
|
||||||
"net"
|
|
||||||
"time"
|
|
||||||
)
|
|
||||||
|
|
||||||
// implement io.Writer
|
|
||||||
type overlayTCP struct {
|
|
||||||
tunnel *P2PTunnel
|
|
||||||
conn net.Conn
|
|
||||||
id uint64
|
|
||||||
rtid uint64
|
|
||||||
running bool
|
|
||||||
isClient bool
|
|
||||||
appID uint64
|
|
||||||
appKey uint64
|
|
||||||
appKeyBytes []byte
|
|
||||||
}
|
|
||||||
|
|
||||||
func (otcp *overlayTCP) run() {
|
|
||||||
gLog.Printf(LevelDEBUG, "%d overlayTCP run start", otcp.id)
|
|
||||||
defer gLog.Printf(LevelDEBUG, "%d overlayTCP run end", otcp.id)
|
|
||||||
otcp.running = true
|
|
||||||
buffer := make([]byte, ReadBuffLen+PaddingSize)
|
|
||||||
readBuf := buffer[:ReadBuffLen]
|
|
||||||
encryptData := make([]byte, ReadBuffLen+PaddingSize) // 16 bytes for padding
|
|
||||||
tunnelHead := new(bytes.Buffer)
|
|
||||||
relayHead := new(bytes.Buffer)
|
|
||||||
binary.Write(relayHead, binary.LittleEndian, otcp.rtid)
|
|
||||||
binary.Write(tunnelHead, binary.LittleEndian, otcp.id)
|
|
||||||
for otcp.running && otcp.tunnel.isRuning() {
|
|
||||||
otcp.conn.SetReadDeadline(time.Now().Add(time.Second * 5))
|
|
||||||
dataLen, err := otcp.conn.Read(readBuf)
|
|
||||||
if err != nil {
|
|
||||||
if ne, ok := err.(net.Error); ok && ne.Timeout() {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
// overlay tcp connection normal close, debug log
|
|
||||||
gLog.Printf(LevelDEBUG, "overlayTCP %d read error:%s,close it", otcp.id, err)
|
|
||||||
break
|
|
||||||
} else {
|
|
||||||
payload := readBuf[:dataLen]
|
|
||||||
if otcp.appKey != 0 {
|
|
||||||
payload, _ = encryptBytes(otcp.appKeyBytes, encryptData, buffer[:dataLen], dataLen)
|
|
||||||
}
|
|
||||||
writeBytes := append(tunnelHead.Bytes(), payload...)
|
|
||||||
if otcp.rtid == 0 {
|
|
||||||
otcp.tunnel.conn.WriteBytes(MsgP2P, MsgOverlayData, writeBytes)
|
|
||||||
} else {
|
|
||||||
// write raley data
|
|
||||||
all := append(relayHead.Bytes(), encodeHeader(MsgP2P, MsgOverlayData, uint32(len(writeBytes)))...)
|
|
||||||
all = append(all, writeBytes...)
|
|
||||||
otcp.tunnel.conn.WriteBytes(MsgP2P, MsgRelayData, all)
|
|
||||||
gLog.Printf(LevelDEBUG, "write relay data to %d:%d bodylen=%d", otcp.rtid, otcp.id, len(writeBytes))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
otcp.conn.Close()
|
|
||||||
otcp.tunnel.overlayConns.Delete(otcp.id)
|
|
||||||
// notify peer disconnect
|
|
||||||
if otcp.isClient {
|
|
||||||
req := OverlayDisconnectReq{ID: otcp.id}
|
|
||||||
if otcp.rtid == 0 {
|
|
||||||
otcp.tunnel.conn.WriteMessage(MsgP2P, MsgOverlayDisconnectReq, &req)
|
|
||||||
} else {
|
|
||||||
// write relay data
|
|
||||||
msg, _ := newMessage(MsgP2P, MsgOverlayDisconnectReq, &req)
|
|
||||||
msgWithHead := append(relayHead.Bytes(), msg...)
|
|
||||||
otcp.tunnel.conn.WriteBytes(MsgP2P, MsgRelayData, msgWithHead)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// calling by p2pTunnel
|
|
||||||
func (otcp *overlayTCP) Write(buff []byte) (n int, err error) {
|
|
||||||
// add mutex when multi-thread calling
|
|
||||||
n, err = otcp.conn.Write(buff)
|
|
||||||
if err != nil {
|
|
||||||
otcp.tunnel.overlayConns.Delete(otcp.id)
|
|
||||||
}
|
|
||||||
return
|
|
||||||
}
|
|
||||||
@@ -6,6 +6,8 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
"net"
|
"net"
|
||||||
|
"strconv"
|
||||||
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
@@ -13,9 +15,11 @@ import (
|
|||||||
type p2pApp struct {
|
type p2pApp struct {
|
||||||
config AppConfig
|
config AppConfig
|
||||||
listener net.Listener
|
listener net.Listener
|
||||||
|
listenerUDP *net.UDPConn
|
||||||
tunnel *P2PTunnel
|
tunnel *P2PTunnel
|
||||||
rtid uint64
|
rtid uint64
|
||||||
relayNode string
|
relayNode string
|
||||||
|
relayMode string
|
||||||
hbTime time.Time
|
hbTime time.Time
|
||||||
hbMtx sync.Mutex
|
hbMtx sync.Mutex
|
||||||
running bool
|
running bool
|
||||||
@@ -44,38 +48,42 @@ func (app *p2pApp) updateHeartbeat() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (app *p2pApp) listenTCP() error {
|
func (app *p2pApp) listenTCP() error {
|
||||||
|
gLog.Printf(LevelDEBUG, "tcp accept on port %d start", app.config.SrcPort)
|
||||||
|
defer gLog.Printf(LevelDEBUG, "tcp accept on port %d end", app.config.SrcPort)
|
||||||
var err error
|
var err error
|
||||||
app.listener, err = net.Listen("tcp4", fmt.Sprintf("0.0.0.0:%d", app.config.SrcPort))
|
app.listener, err = net.Listen("tcp4", fmt.Sprintf("0.0.0.0:%d", app.config.SrcPort))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Printf(LevelERROR, "listen error:%s", err)
|
gLog.Printf(LevelERROR, "listen error:%s", err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
for {
|
for app.running {
|
||||||
conn, err := app.listener.Accept()
|
conn, err := app.listener.Accept()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
if app.running {
|
||||||
gLog.Printf(LevelERROR, "%d accept error:%s", app.tunnel.id, err)
|
gLog.Printf(LevelERROR, "%d accept error:%s", app.tunnel.id, err)
|
||||||
|
}
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
otcp := overlayTCP{
|
oConn := overlayConn{
|
||||||
tunnel: app.tunnel,
|
tunnel: app.tunnel,
|
||||||
conn: conn,
|
connTCP: conn,
|
||||||
id: rand.Uint64(),
|
id: rand.Uint64(),
|
||||||
isClient: true,
|
isClient: true,
|
||||||
rtid: app.rtid,
|
rtid: app.rtid,
|
||||||
appID: app.id,
|
appID: app.id,
|
||||||
appKey: app.key,
|
appKey: app.key,
|
||||||
}
|
}
|
||||||
// calc key bytes for encrypt
|
// pre-calc key bytes for encrypt
|
||||||
if otcp.appKey != 0 {
|
if oConn.appKey != 0 {
|
||||||
encryptKey := make([]byte, AESKeySize)
|
encryptKey := make([]byte, AESKeySize)
|
||||||
binary.LittleEndian.PutUint64(encryptKey, otcp.appKey)
|
binary.LittleEndian.PutUint64(encryptKey, oConn.appKey)
|
||||||
binary.LittleEndian.PutUint64(encryptKey[8:], otcp.appKey)
|
binary.LittleEndian.PutUint64(encryptKey[8:], oConn.appKey)
|
||||||
otcp.appKeyBytes = encryptKey
|
oConn.appKeyBytes = encryptKey
|
||||||
}
|
}
|
||||||
app.tunnel.overlayConns.Store(otcp.id, &otcp)
|
app.tunnel.overlayConns.Store(oConn.id, &oConn)
|
||||||
gLog.Printf(LevelDEBUG, "Accept overlayID:%d", otcp.id)
|
gLog.Printf(LevelDEBUG, "Accept TCP overlayID:%d", oConn.id)
|
||||||
// tell peer connect
|
// tell peer connect
|
||||||
req := OverlayConnectReq{ID: otcp.id,
|
req := OverlayConnectReq{ID: oConn.id,
|
||||||
Token: app.tunnel.pn.config.Token,
|
Token: app.tunnel.pn.config.Token,
|
||||||
DstIP: app.config.DstHost,
|
DstIP: app.config.DstHost,
|
||||||
DstPort: app.config.DstPort,
|
DstPort: app.config.DstPort,
|
||||||
@@ -92,29 +100,117 @@ func (app *p2pApp) listenTCP() error {
|
|||||||
msgWithHead := append(relayHead.Bytes(), msg...)
|
msgWithHead := append(relayHead.Bytes(), msg...)
|
||||||
app.tunnel.conn.WriteBytes(MsgP2P, MsgRelayData, msgWithHead)
|
app.tunnel.conn.WriteBytes(MsgP2P, MsgRelayData, msgWithHead)
|
||||||
}
|
}
|
||||||
|
go oConn.run()
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
go otcp.run()
|
func (app *p2pApp) listenUDP() error {
|
||||||
|
gLog.Printf(LevelDEBUG, "udp accept on port %d start", app.config.SrcPort)
|
||||||
|
defer gLog.Printf(LevelDEBUG, "udp accept on port %d end", app.config.SrcPort)
|
||||||
|
var err error
|
||||||
|
app.listenerUDP, err = net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4zero, Port: app.config.SrcPort})
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "listen error:%s", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
buffer := make([]byte, 64*1024)
|
||||||
|
udpID := make([]byte, 8)
|
||||||
|
for {
|
||||||
|
app.listenerUDP.SetReadDeadline(time.Now().Add(time.Second * 10))
|
||||||
|
len, remoteAddr, err := app.listenerUDP.ReadFrom(buffer)
|
||||||
|
if err != nil {
|
||||||
|
if ne, ok := err.(net.Error); ok && ne.Timeout() {
|
||||||
|
continue
|
||||||
|
} else {
|
||||||
|
gLog.Printf(LevelERROR, "udp read failed:%s", err)
|
||||||
|
break
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
b := bytes.Buffer{}
|
||||||
|
b.Write(buffer[:len])
|
||||||
|
// load from app.tunnel.overlayConns by remoteAddr error, new udp connection
|
||||||
|
remoteIP := strings.Split(remoteAddr.String(), ":")[0]
|
||||||
|
port, _ := strconv.Atoi(strings.Split(remoteAddr.String(), ":")[1])
|
||||||
|
a := net.ParseIP(remoteIP)
|
||||||
|
udpID[0] = a[0]
|
||||||
|
udpID[1] = a[1]
|
||||||
|
udpID[2] = a[2]
|
||||||
|
udpID[3] = a[3]
|
||||||
|
udpID[4] = byte(port)
|
||||||
|
udpID[5] = byte(port >> 8)
|
||||||
|
id := binary.LittleEndian.Uint64(udpID)
|
||||||
|
s, ok := app.tunnel.overlayConns.Load(id)
|
||||||
|
if !ok {
|
||||||
|
oConn := overlayConn{
|
||||||
|
tunnel: app.tunnel,
|
||||||
|
connUDP: app.listenerUDP,
|
||||||
|
remoteAddr: remoteAddr,
|
||||||
|
udpRelayData: make(chan []byte, 1000),
|
||||||
|
id: id,
|
||||||
|
isClient: true,
|
||||||
|
rtid: app.rtid,
|
||||||
|
appID: app.id,
|
||||||
|
appKey: app.key,
|
||||||
|
}
|
||||||
|
// calc key bytes for encrypt
|
||||||
|
if oConn.appKey != 0 {
|
||||||
|
encryptKey := make([]byte, AESKeySize)
|
||||||
|
binary.LittleEndian.PutUint64(encryptKey, oConn.appKey)
|
||||||
|
binary.LittleEndian.PutUint64(encryptKey[8:], oConn.appKey)
|
||||||
|
oConn.appKeyBytes = encryptKey
|
||||||
|
}
|
||||||
|
app.tunnel.overlayConns.Store(oConn.id, &oConn)
|
||||||
|
gLog.Printf(LevelDEBUG, "Accept UDP overlayID:%d", oConn.id)
|
||||||
|
// tell peer connect
|
||||||
|
req := OverlayConnectReq{ID: oConn.id,
|
||||||
|
Token: app.tunnel.pn.config.Token,
|
||||||
|
DstIP: app.config.DstHost,
|
||||||
|
DstPort: app.config.DstPort,
|
||||||
|
Protocol: app.config.Protocol,
|
||||||
|
AppID: app.id,
|
||||||
|
}
|
||||||
|
if app.rtid == 0 {
|
||||||
|
app.tunnel.conn.WriteMessage(MsgP2P, MsgOverlayConnectReq, &req)
|
||||||
|
} else {
|
||||||
|
req.RelayTunnelID = app.tunnel.id
|
||||||
|
relayHead := new(bytes.Buffer)
|
||||||
|
binary.Write(relayHead, binary.LittleEndian, app.rtid)
|
||||||
|
msg, _ := newMessage(MsgP2P, MsgOverlayConnectReq, &req)
|
||||||
|
msgWithHead := append(relayHead.Bytes(), msg...)
|
||||||
|
app.tunnel.conn.WriteBytes(MsgP2P, MsgRelayData, msgWithHead)
|
||||||
|
}
|
||||||
|
go oConn.run()
|
||||||
|
oConn.udpRelayData <- b.Bytes()
|
||||||
|
}
|
||||||
|
|
||||||
|
// load from app.tunnel.overlayConns by remoteAddr ok, write relay data
|
||||||
|
overlayConn, ok := s.(*overlayConn)
|
||||||
|
if !ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
overlayConn.udpRelayData <- b.Bytes()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (app *p2pApp) listen() error {
|
func (app *p2pApp) listen() error {
|
||||||
gLog.Printf(LevelINFO, "LISTEN ON PORT %d START", app.config.SrcPort)
|
gLog.Printf(LevelINFO, "LISTEN ON PORT %s:%d START", app.config.Protocol, app.config.SrcPort)
|
||||||
defer gLog.Printf(LevelINFO, "LISTEN ON PORT %d START", app.config.SrcPort)
|
defer gLog.Printf(LevelINFO, "LISTEN ON PORT %s:%d END", app.config.Protocol, app.config.SrcPort)
|
||||||
app.wg.Add(1)
|
app.wg.Add(1)
|
||||||
defer app.wg.Done()
|
defer app.wg.Done()
|
||||||
app.running = true
|
app.running = true
|
||||||
if app.rtid != 0 {
|
if app.rtid != 0 {
|
||||||
go app.relayHeartbeatLoop()
|
go app.relayHeartbeatLoop()
|
||||||
}
|
}
|
||||||
for app.running {
|
for app.tunnel.isRuning() && app.running {
|
||||||
if app.config.Protocol == "udp" {
|
if app.config.Protocol == "udp" {
|
||||||
app.listenTCP()
|
app.listenUDP()
|
||||||
} else {
|
} else {
|
||||||
app.listenTCP()
|
app.listenTCP()
|
||||||
}
|
}
|
||||||
time.Sleep(time.Second * 5)
|
time.Sleep(time.Second * 10)
|
||||||
// TODO: listen UDP
|
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -124,6 +220,9 @@ func (app *p2pApp) close() {
|
|||||||
if app.listener != nil {
|
if app.listener != nil {
|
||||||
app.listener.Close()
|
app.listener.Close()
|
||||||
}
|
}
|
||||||
|
if app.listenerUDP != nil {
|
||||||
|
app.listenerUDP.Close()
|
||||||
|
}
|
||||||
if app.tunnel != nil {
|
if app.tunnel != nil {
|
||||||
app.tunnel.closeOverlayConns(app.id)
|
app.tunnel.closeOverlayConns(app.id)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"math"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
"net"
|
"net"
|
||||||
"net/url"
|
"net/url"
|
||||||
@@ -93,9 +94,14 @@ func (pn *P2PNetwork) Connect(timeout int) bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) runAll() {
|
func (pn *P2PNetwork) runAll() {
|
||||||
gConf.mtx.Lock()
|
gConf.mtx.Lock() // lock for copy gConf.Apps and the modification of config(it's pointer)
|
||||||
defer gConf.mtx.Unlock()
|
defer gConf.mtx.Unlock()
|
||||||
for _, config := range gConf.Apps {
|
allApps := gConf.Apps // read a copy, other thread will modify the gConf.Apps
|
||||||
|
|
||||||
|
for _, config := range allApps {
|
||||||
|
if config.nextRetryTime.After(time.Now()) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
if config.Enabled == 0 {
|
if config.Enabled == 0 {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
@@ -103,37 +109,43 @@ func (pn *P2PNetwork) runAll() {
|
|||||||
config.AppName = fmt.Sprintf("%s%d", config.Protocol, config.SrcPort)
|
config.AppName = fmt.Sprintf("%s%d", config.Protocol, config.SrcPort)
|
||||||
}
|
}
|
||||||
appExist := false
|
appExist := false
|
||||||
appActive := false
|
var appID uint64
|
||||||
i, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
i, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
if ok {
|
if ok {
|
||||||
app := i.(*p2pApp)
|
app := i.(*p2pApp)
|
||||||
appExist = true
|
appExist = true
|
||||||
|
appID = app.id
|
||||||
if app.isActive() {
|
if app.isActive() {
|
||||||
appActive = true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if appExist && appActive {
|
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if appExist && !appActive {
|
}
|
||||||
gLog.Printf(LevelINFO, "detect app %s disconnect, reconnecting...", config.AppName)
|
if appExist {
|
||||||
pn.DeleteApp(config)
|
pn.DeleteApp(*config)
|
||||||
if config.retryTime.Add(time.Minute * 15).Before(time.Now()) {
|
}
|
||||||
|
if config.retryNum > 0 {
|
||||||
|
gLog.Printf(LevelINFO, "detect app %s(%d) disconnect, reconnecting the %d times...", config.AppName, appID, config.retryNum)
|
||||||
|
if time.Now().Add(-time.Minute * 15).After(config.retryTime) { // normal lasts 15min
|
||||||
config.retryNum = 0
|
config.retryNum = 0
|
||||||
}
|
}
|
||||||
|
}
|
||||||
config.retryNum++
|
config.retryNum++
|
||||||
config.retryTime = time.Now()
|
config.retryTime = time.Now()
|
||||||
if config.retryNum > MaxRetry {
|
increase := math.Pow(1.3, float64(config.retryNum))
|
||||||
gLog.Printf(LevelERROR, "app %s%d retry more than %d times, exit.", config.Protocol, config.SrcPort, MaxRetry)
|
if increase > 900 {
|
||||||
continue
|
increase = 900
|
||||||
}
|
}
|
||||||
|
config.nextRetryTime = time.Now().Add(time.Second * time.Duration(increase)) // exponential increase retry time. 1.3^x
|
||||||
|
config.connectTime = time.Now()
|
||||||
|
gConf.mtx.Unlock() // AddApp will take a period of time
|
||||||
|
err := pn.AddApp(*config)
|
||||||
|
gConf.mtx.Lock()
|
||||||
|
if err != nil {
|
||||||
|
config.errMsg = err.Error()
|
||||||
}
|
}
|
||||||
go pn.AddApp(config)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
func (pn *P2PNetwork) autorunApp() {
|
func (pn *P2PNetwork) autorunApp() {
|
||||||
gLog.Println(LevelINFO, "autorunApp start")
|
gLog.Println(LevelINFO, "autorunApp start")
|
||||||
// TODO: use gConf to check reconnect
|
|
||||||
for pn.running {
|
for pn.running {
|
||||||
time.Sleep(time.Second)
|
time.Sleep(time.Second)
|
||||||
if !pn.online {
|
if !pn.online {
|
||||||
@@ -145,23 +157,23 @@ func (pn *P2PNetwork) autorunApp() {
|
|||||||
gLog.Println(LevelINFO, "autorunApp end")
|
gLog.Println(LevelINFO, "autorunApp end")
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) addRelayTunnel(config AppConfig, appid uint64, appkey uint64) (*P2PTunnel, uint64, error) {
|
func (pn *P2PNetwork) addRelayTunnel(config AppConfig, appid uint64, appkey uint64) (*P2PTunnel, uint64, string, error) {
|
||||||
gLog.Printf(LevelINFO, "addRelayTunnel to %s start", config.PeerNode)
|
gLog.Printf(LevelINFO, "addRelayTunnel to %s start", config.PeerNode)
|
||||||
defer gLog.Printf(LevelINFO, "addRelayTunnel to %s end", config.PeerNode)
|
defer gLog.Printf(LevelINFO, "addRelayTunnel to %s end", config.PeerNode)
|
||||||
pn.write(MsgRelay, MsgRelayNodeReq, &RelayNodeReq{config.PeerNode})
|
pn.write(MsgRelay, MsgRelayNodeReq, &RelayNodeReq{config.PeerNode})
|
||||||
head, body := pn.read("", MsgRelay, MsgRelayNodeRsp, time.Second*10)
|
head, body := pn.read("", MsgRelay, MsgRelayNodeRsp, time.Second*10)
|
||||||
if head == nil {
|
if head == nil {
|
||||||
return nil, 0, errors.New("read MsgRelayNodeRsp error")
|
return nil, 0, "", errors.New("read MsgRelayNodeRsp error")
|
||||||
}
|
}
|
||||||
rsp := RelayNodeRsp{}
|
rsp := RelayNodeRsp{}
|
||||||
err := json.Unmarshal(body, &rsp)
|
err := json.Unmarshal(body, &rsp)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Printf(LevelERROR, "wrong RelayNodeRsp:%s", err)
|
gLog.Printf(LevelERROR, "wrong RelayNodeRsp:%s", err)
|
||||||
return nil, 0, errors.New("unmarshal MsgRelayNodeRsp error")
|
return nil, 0, "", errors.New("unmarshal MsgRelayNodeRsp error")
|
||||||
}
|
}
|
||||||
if rsp.RelayName == "" || rsp.RelayToken == 0 {
|
if rsp.RelayName == "" || rsp.RelayToken == 0 {
|
||||||
gLog.Printf(LevelERROR, "MsgRelayNodeReq error")
|
gLog.Printf(LevelERROR, "MsgRelayNodeReq error")
|
||||||
return nil, 0, errors.New("MsgRelayNodeReq error")
|
return nil, 0, "", errors.New("MsgRelayNodeReq error")
|
||||||
}
|
}
|
||||||
gLog.Printf(LevelINFO, "got relay node:%s", rsp.RelayName)
|
gLog.Printf(LevelINFO, "got relay node:%s", rsp.RelayName)
|
||||||
relayConfig := config
|
relayConfig := config
|
||||||
@@ -170,7 +182,7 @@ func (pn *P2PNetwork) addRelayTunnel(config AppConfig, appid uint64, appkey uint
|
|||||||
t, err := pn.addDirectTunnel(relayConfig, 0)
|
t, err := pn.addDirectTunnel(relayConfig, 0)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "direct connect error:", err)
|
gLog.Println(LevelERROR, "direct connect error:", err)
|
||||||
return nil, 0, err
|
return nil, 0, "", err
|
||||||
}
|
}
|
||||||
// notify peer addRelayTunnel
|
// notify peer addRelayTunnel
|
||||||
req := AddRelayTunnelReq{
|
req := AddRelayTunnelReq{
|
||||||
@@ -187,17 +199,18 @@ func (pn *P2PNetwork) addRelayTunnel(config AppConfig, appid uint64, appkey uint
|
|||||||
head, body = pn.read(config.PeerNode, MsgPush, MsgPushAddRelayTunnelRsp, PeerAddRelayTimeount) // TODO: const value
|
head, body = pn.read(config.PeerNode, MsgPush, MsgPushAddRelayTunnelRsp, PeerAddRelayTimeount) // TODO: const value
|
||||||
if head == nil {
|
if head == nil {
|
||||||
gLog.Printf(LevelERROR, "read MsgPushAddRelayTunnelRsp error")
|
gLog.Printf(LevelERROR, "read MsgPushAddRelayTunnelRsp error")
|
||||||
return nil, 0, errors.New("read MsgPushAddRelayTunnelRsp error")
|
return nil, 0, "", errors.New("read MsgPushAddRelayTunnelRsp error")
|
||||||
}
|
}
|
||||||
rspID := TunnelMsg{}
|
rspID := TunnelMsg{}
|
||||||
err = json.Unmarshal(body, &rspID)
|
err = json.Unmarshal(body, &rspID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Printf(LevelERROR, "wrong RelayNodeRsp:%s", err)
|
gLog.Printf(LevelERROR, "wrong RelayNodeRsp:%s", err)
|
||||||
return nil, 0, errors.New("unmarshal MsgRelayNodeRsp error")
|
return nil, 0, "", errors.New("unmarshal MsgRelayNodeRsp error")
|
||||||
}
|
}
|
||||||
return t, rspID.ID, err
|
return t, rspID.ID, rsp.Mode, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// use *AppConfig to save status
|
||||||
func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
||||||
gLog.Printf(LevelINFO, "addApp %s to %s:%s:%d start", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
gLog.Printf(LevelINFO, "addApp %s to %s:%s:%d start", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
||||||
defer gLog.Printf(LevelINFO, "addApp %s to %s:%s:%d end", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
defer gLog.Printf(LevelINFO, "addApp %s to %s:%s:%d end", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
||||||
@@ -215,24 +228,26 @@ func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
|||||||
}
|
}
|
||||||
appID := rand.Uint64()
|
appID := rand.Uint64()
|
||||||
appKey := uint64(0)
|
appKey := uint64(0)
|
||||||
t, err := pn.addDirectTunnel(config, 0)
|
|
||||||
var rtid uint64
|
var rtid uint64
|
||||||
relayNode := ""
|
relayNode := ""
|
||||||
|
relayMode := ""
|
||||||
peerNatType := NATUnknown
|
peerNatType := NATUnknown
|
||||||
peerIP := ""
|
peerIP := ""
|
||||||
errMsg := ""
|
errMsg := ""
|
||||||
if err != nil && err == ErrorHandshake {
|
t, err := pn.addDirectTunnel(config, 0)
|
||||||
gLog.Println(LevelERROR, "direct connect failed, try to relay")
|
|
||||||
appKey = rand.Uint64()
|
|
||||||
t, rtid, err = pn.addRelayTunnel(config, appID, appKey)
|
|
||||||
if t != nil {
|
|
||||||
relayNode = t.config.PeerNode
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if t != nil {
|
if t != nil {
|
||||||
peerNatType = t.config.peerNatType
|
peerNatType = t.config.peerNatType
|
||||||
peerIP = t.config.peerIP
|
peerIP = t.config.peerIP
|
||||||
}
|
}
|
||||||
|
if err != nil && err == ErrorHandshake {
|
||||||
|
gLog.Println(LevelERROR, "direct connect failed, try to relay")
|
||||||
|
appKey = rand.Uint64()
|
||||||
|
t, rtid, relayMode, err = pn.addRelayTunnel(config, appID, appKey)
|
||||||
|
if t != nil {
|
||||||
|
relayNode = t.config.PeerNode
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
errMsg = err.Error()
|
errMsg = err.Error()
|
||||||
}
|
}
|
||||||
@@ -261,6 +276,7 @@ func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
|||||||
config: config,
|
config: config,
|
||||||
rtid: rtid,
|
rtid: rtid,
|
||||||
relayNode: relayNode,
|
relayNode: relayNode,
|
||||||
|
relayMode: relayMode,
|
||||||
hbTime: time.Now()}
|
hbTime: time.Now()}
|
||||||
pn.apps.Store(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort), &app)
|
pn.apps.Store(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort), &app)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
@@ -452,10 +468,8 @@ func (pn *P2PNetwork) handleMessage(t int, msg []byte) {
|
|||||||
pn.serverTs = rsp.Ts
|
pn.serverTs = rsp.Ts
|
||||||
pn.config.Token = rsp.Token
|
pn.config.Token = rsp.Token
|
||||||
pn.config.User = rsp.User
|
pn.config.User = rsp.User
|
||||||
gConf.mtx.Lock()
|
gConf.setToken(rsp.Token)
|
||||||
gConf.Network.Token = rsp.Token
|
gConf.setUser(rsp.User)
|
||||||
gConf.Network.User = rsp.User
|
|
||||||
gConf.mtx.Unlock()
|
|
||||||
gConf.save()
|
gConf.save()
|
||||||
pn.localTs = time.Now().Unix()
|
pn.localTs = time.Now().Unix()
|
||||||
gLog.Printf(LevelINFO, "login ok. user=%s,Server ts=%d, local ts=%d", rsp.User, rsp.Ts, pn.localTs)
|
gLog.Printf(LevelINFO, "login ok. user=%s,Server ts=%d, local ts=%d", rsp.User, rsp.Ts, pn.localTs)
|
||||||
|
|||||||
@@ -276,7 +276,7 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
gLog.Printf(LevelDEBUG, "%d tunnel not found overlay connection %d", t.id, overlayID)
|
gLog.Printf(LevelDEBUG, "%d tunnel not found overlay connection %d", t.id, overlayID)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
overlayConn, ok := s.(*overlayTCP)
|
overlayConn, ok := s.(*overlayConn)
|
||||||
if !ok {
|
if !ok {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
@@ -333,32 +333,34 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
|
|
||||||
overlayID := req.ID
|
overlayID := req.ID
|
||||||
gLog.Printf(LevelDEBUG, "App:%d overlayID:%d connect %+v", req.AppID, overlayID, req)
|
gLog.Printf(LevelDEBUG, "App:%d overlayID:%d connect %+v", req.AppID, overlayID, req)
|
||||||
if req.Protocol == "tcp" {
|
oConn := overlayConn{
|
||||||
conn, err := net.DialTimeout("tcp", fmt.Sprintf("%s:%d", req.DstIP, req.DstPort), time.Second*5)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Println(LevelERROR, err)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
otcp := overlayTCP{
|
|
||||||
tunnel: t,
|
tunnel: t,
|
||||||
conn: conn,
|
|
||||||
id: overlayID,
|
id: overlayID,
|
||||||
isClient: false,
|
isClient: false,
|
||||||
rtid: req.RelayTunnelID,
|
rtid: req.RelayTunnelID,
|
||||||
appID: req.AppID,
|
appID: req.AppID,
|
||||||
appKey: GetKey(req.AppID),
|
appKey: GetKey(req.AppID),
|
||||||
}
|
}
|
||||||
// calc key bytes for encrypt
|
if req.Protocol == "udp" {
|
||||||
if otcp.appKey != 0 {
|
oConn.connUDP, err = net.DialUDP("udp", nil, &net.UDPAddr{IP: net.ParseIP(req.DstIP), Port: req.DstPort})
|
||||||
encryptKey := make([]byte, 16)
|
} else {
|
||||||
binary.LittleEndian.PutUint64(encryptKey, otcp.appKey)
|
oConn.connTCP, err = net.DialTimeout("tcp", fmt.Sprintf("%s:%d", req.DstIP, req.DstPort), time.Second*5)
|
||||||
binary.LittleEndian.PutUint64(encryptKey[8:], otcp.appKey)
|
}
|
||||||
otcp.appKeyBytes = encryptKey
|
if err != nil {
|
||||||
|
gLog.Println(LevelERROR, err)
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
t.overlayConns.Store(otcp.id, &otcp)
|
// calc key bytes for encrypt
|
||||||
go otcp.run()
|
if oConn.appKey != 0 {
|
||||||
|
encryptKey := make([]byte, 16)
|
||||||
|
binary.LittleEndian.PutUint64(encryptKey, oConn.appKey)
|
||||||
|
binary.LittleEndian.PutUint64(encryptKey[8:], oConn.appKey)
|
||||||
|
oConn.appKeyBytes = encryptKey
|
||||||
}
|
}
|
||||||
|
|
||||||
|
t.overlayConns.Store(oConn.id, &oConn)
|
||||||
|
go oConn.run()
|
||||||
case MsgOverlayDisconnectReq:
|
case MsgOverlayDisconnectReq:
|
||||||
req := OverlayDisconnectReq{}
|
req := OverlayDisconnectReq{}
|
||||||
err := json.Unmarshal(body, &req)
|
err := json.Unmarshal(body, &req)
|
||||||
@@ -370,8 +372,8 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
gLog.Printf(LevelDEBUG, "%d disconnect overlay connection %d", t.id, overlayID)
|
gLog.Printf(LevelDEBUG, "%d disconnect overlay connection %d", t.id, overlayID)
|
||||||
i, ok := t.overlayConns.Load(overlayID)
|
i, ok := t.overlayConns.Load(overlayID)
|
||||||
if ok {
|
if ok {
|
||||||
otcp := i.(*overlayTCP)
|
oConn := i.(*overlayConn)
|
||||||
otcp.running = false
|
oConn.running = false
|
||||||
}
|
}
|
||||||
default:
|
default:
|
||||||
}
|
}
|
||||||
@@ -408,9 +410,16 @@ func (t *P2PTunnel) listen() error {
|
|||||||
|
|
||||||
func (t *P2PTunnel) closeOverlayConns(appID uint64) {
|
func (t *P2PTunnel) closeOverlayConns(appID uint64) {
|
||||||
t.overlayConns.Range(func(_, i interface{}) bool {
|
t.overlayConns.Range(func(_, i interface{}) bool {
|
||||||
otcp := i.(*overlayTCP)
|
oConn := i.(*overlayConn)
|
||||||
if otcp.appID == appID {
|
if oConn.appID == appID {
|
||||||
otcp.conn.Close()
|
if oConn.connTCP != nil {
|
||||||
|
oConn.connTCP.Close()
|
||||||
|
oConn.connTCP = nil
|
||||||
|
}
|
||||||
|
if oConn.connUDP != nil {
|
||||||
|
oConn.connUDP.Close()
|
||||||
|
oConn.connUDP = nil
|
||||||
|
}
|
||||||
}
|
}
|
||||||
return true
|
return true
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -10,7 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
const OpenP2PVersion = "1.1.0"
|
const OpenP2PVersion = "1.4.2"
|
||||||
const ProducnName string = "openp2p"
|
const ProducnName string = "openp2p"
|
||||||
|
|
||||||
type openP2PHeader struct {
|
type openP2PHeader struct {
|
||||||
@@ -117,7 +117,7 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
ReadBuffLen = 1024
|
ReadBuffLen = 4096 // for UDP maybe not enough
|
||||||
NetworkHeartbeatTime = time.Second * 30 // TODO: server no response hb, save flow
|
NetworkHeartbeatTime = time.Second * 30 // TODO: server no response hb, save flow
|
||||||
TunnelHeartbeatTime = time.Second * 15
|
TunnelHeartbeatTime = time.Second * 15
|
||||||
TunnelIdleTimeout = time.Minute
|
TunnelIdleTimeout = time.Minute
|
||||||
@@ -134,6 +134,7 @@ const (
|
|||||||
PublicIPEchoTimeout = time.Second * 3
|
PublicIPEchoTimeout = time.Second * 3
|
||||||
NatTestTimeout = time.Second * 10
|
NatTestTimeout = time.Second * 10
|
||||||
ClientAPITimeout = time.Second * 10
|
ClientAPITimeout = time.Second * 10
|
||||||
|
MaxDirectTry = 5
|
||||||
)
|
)
|
||||||
|
|
||||||
// NATNone has public ip
|
// NATNone has public ip
|
||||||
@@ -236,6 +237,7 @@ type RelayNodeReq struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type RelayNodeRsp struct {
|
type RelayNodeRsp struct {
|
||||||
|
Mode string `json:"mode,omitempty"` // private,public
|
||||||
RelayName string `json:"relayName,omitempty"`
|
RelayName string `json:"relayName,omitempty"`
|
||||||
RelayToken uint64 `json:"relayToken,omitempty"`
|
RelayToken uint64 `json:"relayToken,omitempty"`
|
||||||
}
|
}
|
||||||
@@ -269,7 +271,7 @@ type ReportConnect struct {
|
|||||||
NatType int `json:"natType,omitempty"`
|
NatType int `json:"natType,omitempty"`
|
||||||
PeerNode string `json:"peerNode,omitempty"`
|
PeerNode string `json:"peerNode,omitempty"`
|
||||||
DstPort int `json:"dstPort,omitempty"`
|
DstPort int `json:"dstPort,omitempty"`
|
||||||
DstHost string `json:"dsdtHost,omitempty"`
|
DstHost string `json:"dstHost,omitempty"`
|
||||||
PeerUser string `json:"peerUser,omitempty"`
|
PeerUser string `json:"peerUser,omitempty"`
|
||||||
PeerNatType int `json:"peerNatType,omitempty"`
|
PeerNatType int `json:"peerNatType,omitempty"`
|
||||||
PeerIP string `json:"peerIP,omitempty"`
|
PeerIP string `json:"peerIP,omitempty"`
|
||||||
@@ -294,8 +296,10 @@ type AppInfo struct {
|
|||||||
PeerIP string `json:"peerIP,omitempty"`
|
PeerIP string `json:"peerIP,omitempty"`
|
||||||
ShareBandwidth int `json:"shareBandWidth,omitempty"`
|
ShareBandwidth int `json:"shareBandWidth,omitempty"`
|
||||||
RelayNode string `json:"relayNode,omitempty"`
|
RelayNode string `json:"relayNode,omitempty"`
|
||||||
|
RelayMode string `json:"relayMode,omitempty"`
|
||||||
Version string `json:"version,omitempty"`
|
Version string `json:"version,omitempty"`
|
||||||
RetryTime string `json:"retryTime,omitempty"`
|
RetryTime string `json:"retryTime,omitempty"`
|
||||||
|
ConnectTime string `json:"connectTime,omitempty"`
|
||||||
IsActive int `json:"isActive,omitempty"`
|
IsActive int `json:"isActive,omitempty"`
|
||||||
Enabled int `json:"enabled,omitempty"`
|
Enabled int `json:"enabled,omitempty"`
|
||||||
}
|
}
|
||||||
|
|||||||