Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
029d69869f |
+1
-1
@@ -45,7 +45,7 @@ P2P直连可以让你的设备跑满带宽。不论你的设备在任何网络
|
|||||||

|

|
||||||
3. 在家里下载最新的OpenP2P,解压出来,在命令行执行
|
3. 在家里下载最新的OpenP2P,解压出来,在命令行执行
|
||||||
```
|
```
|
||||||
openp2p.exe -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 --peernode OFFICEPC1 --dstip 127.0.0.1 --dstport 3389 --srcport 23389 --protocol tcp
|
openp2p.exe -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 -appname WindowsRemote -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389 -protocol tcp
|
||||||
```
|
```
|
||||||
> :warning: **切记将标记大写的参数改成自己的**
|
> :warning: **切记将标记大写的参数改成自己的**
|
||||||
|
|
||||||
|
|||||||
@@ -56,7 +56,7 @@ Under the outbreak of covid-19 pandemic, surely remote work becomes a fundamenta
|
|||||||
|
|
||||||
3. Download OpenP2P on your home device,unzip and execute below command line.
|
3. Download OpenP2P on your home device,unzip and execute below command line.
|
||||||
```
|
```
|
||||||
openp2p.exe -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 --peernode OFFICEPC1 --dstip 127.0.0.1 --dstport 22 --srcport 22022 --protocol tcp
|
openp2p.exe -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 -appname OfficeSSH -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 22 -srcport 22022 -protocol tcp
|
||||||
```
|
```
|
||||||
|
|
||||||
> :warning: **Must change the parameters marked in UPPERCASE to your own**
|
> :warning: **Must change the parameters marked in UPPERCASE to your own**
|
||||||
|
|||||||
+3
-3
@@ -14,17 +14,17 @@
|
|||||||
>* -node: 独一无二的节点名字,唯一标识
|
>* -node: 独一无二的节点名字,唯一标识
|
||||||
>* -user: 独一无二的用户名字,该节点属于这个user
|
>* -user: 独一无二的用户名字,该节点属于这个user
|
||||||
>* -password: 密码
|
>* -password: 密码
|
||||||
>* -sharebandwidth: 作为共享节点时提供带宽,默认10mbps. 如果是光纤大带宽,设置越大效果越好
|
>* -sharebandwidth: 作为共享节点时提供带宽,默认10mbps. 如果是光纤大带宽,设置越大效果越好. -1表示不共享,该节点只在私有的P2P网络使用。不加入共享的P2P网络,这样也意味着无法使用别人的共享节点
|
||||||
>* -loglevel: 需要查看更多调试日志,设置0;默认是1
|
>* -loglevel: 需要查看更多调试日志,设置0;默认是1
|
||||||
>* -noshare: 不共享,该节点只在私有的P2P网络使用。不加入共享的P2P网络,这样也意味着无法使用别人的共享节点
|
|
||||||
|
|
||||||
## 连接
|
## 连接
|
||||||
```
|
```
|
||||||
./openp2p -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389 -protocol tcp
|
./openp2p -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 -appname OfficeWindowsRemote -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389 -protocol tcp
|
||||||
使用配置文件,建立多个P2PApp
|
使用配置文件,建立多个P2PApp
|
||||||
./openp2p -d -f
|
./openp2p -d -f
|
||||||
./openp2p -f
|
./openp2p -f
|
||||||
```
|
```
|
||||||
|
>* -appname: 这个P2P应用名字
|
||||||
>* -peernode: 目标节点名字
|
>* -peernode: 目标节点名字
|
||||||
>* -dstip: 目标服务地址,默认本机127.0.0.1
|
>* -dstip: 目标服务地址,默认本机127.0.0.1
|
||||||
>* -dstport: 目标服务端口,常见的如windows远程桌面3389,Linux ssh 22
|
>* -dstport: 目标服务端口,常见的如windows远程桌面3389,Linux ssh 22
|
||||||
|
|||||||
@@ -15,17 +15,17 @@ Or
|
|||||||
>* -node: Unique node name, unique identification
|
>* -node: Unique node name, unique identification
|
||||||
>* -user: Unique user name, the node belongs to this user
|
>* -user: Unique user name, the node belongs to this user
|
||||||
>* -password: Password
|
>* -password: Password
|
||||||
>* -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
|
>* -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. -1 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
|
||||||
>* -noshare: 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
|
|
||||||
|
|
||||||
## Connect
|
## Connect
|
||||||
```
|
```
|
||||||
./openp2p -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389 -protocol tcp
|
./openp2p -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 -appname OfficeWindowsRemote -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389 -protocol tcp
|
||||||
Create multiple P2PApp by config file
|
Create multiple P2PApp by config file
|
||||||
./openp2p -d -f
|
./openp2p -d -f
|
||||||
./openp2p -f
|
./openp2p -f
|
||||||
```
|
```
|
||||||
|
>* -appname: This P2PApp name
|
||||||
>* -peernode: Target node name
|
>* -peernode: Target node name
|
||||||
>* -dstip: Target service address, default local 127.0.0.1
|
>* -dstip: Target service address, default local 127.0.0.1
|
||||||
>* -dstport: Target service port, such as windows remote desktop 3389, Linux ssh 22
|
>* -dstport: Target service port, such as windows remote desktop 3389, Linux ssh 22
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package main
|
|||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -10,6 +11,7 @@ var gConf Config
|
|||||||
|
|
||||||
type AppConfig struct {
|
type AppConfig struct {
|
||||||
// required
|
// required
|
||||||
|
AppName string
|
||||||
Protocol string
|
Protocol string
|
||||||
SrcPort int
|
SrcPort int
|
||||||
PeerNode string
|
PeerNode string
|
||||||
@@ -32,9 +34,13 @@ type Config struct {
|
|||||||
Network NetworkConfig `json:"network"`
|
Network NetworkConfig `json:"network"`
|
||||||
Apps []AppConfig `json:"apps"`
|
Apps []AppConfig `json:"apps"`
|
||||||
daemonMode bool
|
daemonMode bool
|
||||||
|
logLevel int
|
||||||
|
mtx sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) add(app AppConfig) {
|
func (c *Config) add(app AppConfig) {
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
if app.SrcPort == 0 || app.DstPort == 0 {
|
if app.SrcPort == 0 || app.DstPort == 0 {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -46,8 +52,24 @@ func (c *Config) add(app AppConfig) {
|
|||||||
c.Apps = append(c.Apps, app)
|
c.Apps = append(c.Apps, app)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c *Config) delete(app AppConfig) {
|
||||||
|
if app.SrcPort == 0 || app.DstPort == 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
|
for i := 0; i < len(c.Apps); i++ {
|
||||||
|
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
||||||
|
c.Apps = append(c.Apps[:i], c.Apps[i+1:]...)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (c *Config) save() {
|
func (c *Config) save() {
|
||||||
data, _ := json.MarshalIndent(c, "", "")
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
|
data, _ := json.MarshalIndent(c, "", " ")
|
||||||
err := ioutil.WriteFile("config.json", data, 0644)
|
err := ioutil.WriteFile("config.json", data, 0644)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "save config.json error:", err)
|
gLog.Println(LevelERROR, "save config.json error:", err)
|
||||||
@@ -55,6 +77,8 @@ func (c *Config) save() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) load() error {
|
func (c *Config) load() error {
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
data, err := ioutil.ReadFile("config.json")
|
data, err := ioutil.ReadFile("config.json")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "read config.json error:", err)
|
gLog.Println(LevelERROR, "read config.json error:", err)
|
||||||
@@ -72,7 +96,6 @@ type NetworkConfig struct {
|
|||||||
Node string
|
Node string
|
||||||
User string
|
User string
|
||||||
Password string
|
Password string
|
||||||
NoShare bool
|
|
||||||
localIP string
|
localIP string
|
||||||
ipv6 string
|
ipv6 string
|
||||||
hostName string
|
hostName string
|
||||||
|
|||||||
@@ -38,7 +38,6 @@ func (d *daemon) Stop(s service.Service) error {
|
|||||||
func (d *daemon) run() {
|
func (d *daemon) run() {
|
||||||
gLog.Println(LevelINFO, "daemon run start")
|
gLog.Println(LevelINFO, "daemon run start")
|
||||||
defer gLog.Println(LevelINFO, "daemon run end")
|
defer gLog.Println(LevelINFO, "daemon run end")
|
||||||
os.Chdir(filepath.Dir(os.Args[0])) // for system service
|
|
||||||
d.running = true
|
d.running = true
|
||||||
binPath, _ := os.Executable()
|
binPath, _ := os.Executable()
|
||||||
mydir, err := os.Getwd()
|
mydir, err := os.Getwd()
|
||||||
@@ -106,9 +105,9 @@ func (d *daemon) Control(ctrlComm string, exeAbsPath string, args []string) erro
|
|||||||
|
|
||||||
// examples:
|
// examples:
|
||||||
// listen:
|
// listen:
|
||||||
// ./openp2p install -node hhd1207-222 -user tenderiron -password 13760636579 -noshare
|
// ./openp2p install -node hhd1207-222 -user tenderiron -password 13760636579 -sharebandwidth 0
|
||||||
// listen and build p2papp:
|
// listen and build p2papp:
|
||||||
// ./openp2p install -node hhd1207-222 -user tenderiron -password 13760636579 -noshare -peernode hhdhome-n1 -dstip 127.0.0.1 -dstport 50022 -protocol tcp -srcport 22
|
// ./openp2p install -node hhd1207-222 -user tenderiron -password 13760636579 -sharebandwidth 0 -peernode hhdhome-n1 -dstip 127.0.0.1 -dstport 50022 -protocol tcp -srcport 22
|
||||||
func install() {
|
func install() {
|
||||||
gLog = InitLogger(filepath.Dir(os.Args[0]), "openp2p-install", LevelDEBUG, 1024*1024, LogConsole)
|
gLog = InitLogger(filepath.Dir(os.Args[0]), "openp2p-install", LevelDEBUG, 1024*1024, LogConsole)
|
||||||
// save config file
|
// save config file
|
||||||
@@ -125,11 +124,13 @@ func install() {
|
|||||||
dstPort := installFlag.Int("dstport", 0, "destination port ")
|
dstPort := installFlag.Int("dstport", 0, "destination port ")
|
||||||
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")
|
||||||
noShare := installFlag.Bool("noshare", false, "disable using the huge numbers of shared nodes in OpenP2P network, your connectivity will be weak. also this node will not shared with others")
|
appName := flag.String("appname", "", "app name")
|
||||||
|
installFlag.Bool("noshare", false, "deprecated. uses -sharebandwidth -1")
|
||||||
shareBandwidth := installFlag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private node 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:])
|
||||||
checkParams(*node, *user, *password)
|
checkParams(*node, *user, *password)
|
||||||
|
gConf.logLevel = *logLevel
|
||||||
gConf.Network.ServerHost = *serverHost
|
gConf.Network.ServerHost = *serverHost
|
||||||
gConf.Network.User = *user
|
gConf.Network.User = *user
|
||||||
gConf.Network.Node = *node
|
gConf.Network.Node = *node
|
||||||
@@ -137,7 +138,6 @@ func install() {
|
|||||||
gConf.Network.ServerPort = 27182
|
gConf.Network.ServerPort = 27182
|
||||||
gConf.Network.UDPPort1 = 27182
|
gConf.Network.UDPPort1 = 27182
|
||||||
gConf.Network.UDPPort2 = 27183
|
gConf.Network.UDPPort2 = 27183
|
||||||
gConf.Network.NoShare = *noShare
|
|
||||||
gConf.Network.ShareBandwidth = *shareBandwidth
|
gConf.Network.ShareBandwidth = *shareBandwidth
|
||||||
config := AppConfig{}
|
config := AppConfig{}
|
||||||
config.PeerNode = *peerNode
|
config.PeerNode = *peerNode
|
||||||
@@ -147,6 +147,7 @@ func install() {
|
|||||||
config.DstPort = *dstPort
|
config.DstPort = *dstPort
|
||||||
config.SrcPort = *srcPort
|
config.SrcPort = *srcPort
|
||||||
config.Protocol = *protocol
|
config.Protocol = *protocol
|
||||||
|
config.AppName = *appName
|
||||||
gConf.add(config)
|
gConf.add(config)
|
||||||
os.MkdirAll(defaultInstallPath, 0775)
|
os.MkdirAll(defaultInstallPath, 0775)
|
||||||
err := os.Chdir(defaultInstallPath)
|
err := os.Chdir(defaultInstallPath)
|
||||||
@@ -184,14 +185,14 @@ func install() {
|
|||||||
|
|
||||||
// args := []string{""}
|
// args := []string{""}
|
||||||
gLog.Println(LevelINFO, "targetPath:", targetPath)
|
gLog.Println(LevelINFO, "targetPath:", targetPath)
|
||||||
err = d.Control("install", targetPath, []string{"-d", "-f"})
|
err = d.Control("install", targetPath, []string{"-d"})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "install system service error:", err)
|
gLog.Println(LevelERROR, "install system service error:", err)
|
||||||
} else {
|
} else {
|
||||||
gLog.Println(LevelINFO, "install system service ok.")
|
gLog.Println(LevelINFO, "install system service ok.")
|
||||||
}
|
}
|
||||||
time.Sleep(time.Second * 2)
|
time.Sleep(time.Second * 2)
|
||||||
err = d.Control("start", targetPath, []string{"-d", "-f"})
|
err = d.Control("start", targetPath, []string{"-d"})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "start openp2p service error:", err)
|
gLog.Println(LevelERROR, "start openp2p service error:", err)
|
||||||
} else {
|
} else {
|
||||||
|
|||||||
+173
@@ -0,0 +1,173 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"encoding/binary"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
|
||||||
|
pushHead := PushHeader{}
|
||||||
|
err := binary.Read(bytes.NewReader(msg[openP2PHeaderSize:openP2PHeaderSize+PushHeaderSize]), binary.LittleEndian, &pushHead)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
gLog.Printf(LevelDEBUG, "handle push msg type:%d, push header:%+v", subType, pushHead)
|
||||||
|
switch subType {
|
||||||
|
case MsgPushConnectReq:
|
||||||
|
req := PushConnectReq{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong MsgPushConnectReq:%s", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
gLog.Printf(LevelINFO, "%s is connecting...", req.From)
|
||||||
|
gLog.Println(LevelDEBUG, "push connect response to ", req.From)
|
||||||
|
// verify token or name&password
|
||||||
|
if VerifyTOTP(req.Token, pn.config.User, pn.config.Password, time.Now().Unix()+(pn.serverTs-pn.localTs)) || // localTs may behind, auto adjust ts
|
||||||
|
VerifyTOTP(req.Token, pn.config.User, pn.config.Password, time.Now().Unix()) ||
|
||||||
|
(req.User == pn.config.User && req.Password == pn.config.Password) {
|
||||||
|
gLog.Printf(LevelINFO, "Access Granted\n")
|
||||||
|
config := AppConfig{}
|
||||||
|
config.peerNatType = req.NatType
|
||||||
|
config.peerConeNatPort = req.ConeNatPort
|
||||||
|
config.peerIP = req.FromIP
|
||||||
|
config.PeerNode = req.From
|
||||||
|
// share relay node will limit bandwidth
|
||||||
|
if req.User != pn.config.User || req.Password != pn.config.Password {
|
||||||
|
gLog.Printf(LevelINFO, "set share bandwidth %d mbps", pn.config.ShareBandwidth)
|
||||||
|
config.shareBandwidth = pn.config.ShareBandwidth
|
||||||
|
}
|
||||||
|
// go pn.AddTunnel(config, req.ID)
|
||||||
|
go pn.addDirectTunnel(config, req.ID)
|
||||||
|
break
|
||||||
|
}
|
||||||
|
gLog.Println(LevelERROR, "Access Denied:", req.From)
|
||||||
|
rsp := PushConnectRsp{
|
||||||
|
Error: 1,
|
||||||
|
Detail: fmt.Sprintf("connect to %s error: Access Denied", pn.config.Node),
|
||||||
|
To: req.From,
|
||||||
|
From: pn.config.Node,
|
||||||
|
}
|
||||||
|
pn.push(req.From, MsgPushConnectRsp, rsp)
|
||||||
|
case MsgPushRsp:
|
||||||
|
rsp := PushRsp{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize:], &rsp)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong pushRsp:%s", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if rsp.Error == 0 {
|
||||||
|
gLog.Printf(LevelDEBUG, "push ok, detail:%s", rsp.Detail)
|
||||||
|
} else {
|
||||||
|
gLog.Printf(LevelERROR, "push error:%d, detail:%s", rsp.Error, rsp.Detail)
|
||||||
|
}
|
||||||
|
case MsgPushAddRelayTunnelReq:
|
||||||
|
req := AddRelayTunnelReq{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong RelayNodeRsp:%s", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
config := AppConfig{}
|
||||||
|
config.PeerNode = req.RelayName
|
||||||
|
config.peerToken = req.RelayToken
|
||||||
|
// set user password, maybe the relay node is your private node
|
||||||
|
config.PeerUser = pn.config.User
|
||||||
|
config.PeerPassword = pn.config.Password
|
||||||
|
go func(r AddRelayTunnelReq) {
|
||||||
|
t, errDt := pn.addDirectTunnel(config, 0)
|
||||||
|
if errDt == nil {
|
||||||
|
// notify peer relay ready
|
||||||
|
msg := TunnelMsg{ID: t.id}
|
||||||
|
pn.push(r.From, MsgPushAddRelayTunnelRsp, msg)
|
||||||
|
SaveKey(req.AppID, req.AppKey)
|
||||||
|
}
|
||||||
|
|
||||||
|
}(req)
|
||||||
|
case MsgPushUpdate:
|
||||||
|
update()
|
||||||
|
if gConf.daemonMode {
|
||||||
|
os.Exit(0)
|
||||||
|
}
|
||||||
|
case MsgPushReportApps:
|
||||||
|
gLog.Println(LevelINFO, "MsgPushReportApps")
|
||||||
|
req := ReportApps{}
|
||||||
|
// TODO: add the retrying apps
|
||||||
|
gConf.mtx.Lock()
|
||||||
|
defer gConf.mtx.Unlock()
|
||||||
|
for _, config := range gConf.Apps {
|
||||||
|
appInfo := AppInfo{
|
||||||
|
AppName: config.AppName,
|
||||||
|
Protocol: config.Protocol,
|
||||||
|
SrcPort: config.SrcPort,
|
||||||
|
// RelayNode: relayNode,
|
||||||
|
PeerNode: config.PeerNode,
|
||||||
|
DstHost: config.DstHost,
|
||||||
|
DstPort: config.DstPort,
|
||||||
|
PeerUser: config.PeerUser,
|
||||||
|
PeerIP: config.peerIP,
|
||||||
|
PeerNatType: config.peerNatType,
|
||||||
|
RetryTime: config.retryTime.String(),
|
||||||
|
IsActive: 1,
|
||||||
|
}
|
||||||
|
req.Apps = append(req.Apps, appInfo)
|
||||||
|
}
|
||||||
|
// pn.apps.Range(func(_, i interface{}) bool {
|
||||||
|
// app := i.(*p2pApp)
|
||||||
|
// appInfo := AppInfo{
|
||||||
|
// AppName: app.config.AppName,
|
||||||
|
// Protocol: app.config.Protocol,
|
||||||
|
// SrcPort: app.config.SrcPort,
|
||||||
|
// RelayNode: app.relayNode,
|
||||||
|
// PeerNode: app.config.PeerNode,
|
||||||
|
// DstHost: app.config.DstHost,
|
||||||
|
// DstPort: app.config.DstPort,
|
||||||
|
// PeerUser: app.config.PeerUser,
|
||||||
|
// PeerIP: app.config.peerIP,
|
||||||
|
// PeerNatType: app.config.peerNatType,
|
||||||
|
// RetryTime: app.config.retryTime.String(),
|
||||||
|
// IsActive: 1,
|
||||||
|
// }
|
||||||
|
// req.Apps = append(req.Apps, appInfo)
|
||||||
|
// return true
|
||||||
|
// })
|
||||||
|
pn.write(MsgReport, MsgReportApps, &req)
|
||||||
|
case MsgPushEditApp:
|
||||||
|
gLog.Println(LevelINFO, "MsgPushEditApp")
|
||||||
|
newApp := AppInfo{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize:], &newApp)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong MsgPushEditApp:%s %s", err, string(msg[openP2PHeaderSize:]))
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
var config AppConfig
|
||||||
|
// protocol0+srcPort0 exist, delApp
|
||||||
|
config.AppName = newApp.AppName
|
||||||
|
config.Protocol = newApp.Protocol0
|
||||||
|
config.SrcPort = newApp.SrcPort0
|
||||||
|
config.PeerNode = newApp.PeerNode
|
||||||
|
config.DstHost = newApp.DstHost
|
||||||
|
config.DstPort = newApp.DstPort
|
||||||
|
|
||||||
|
gConf.delete(config)
|
||||||
|
// AddApp
|
||||||
|
config.Protocol = newApp.Protocol
|
||||||
|
config.SrcPort = newApp.SrcPort
|
||||||
|
gConf.add(config)
|
||||||
|
gConf.save()
|
||||||
|
pn.DeleteApp(config) // save quickly for the next request reportApplist
|
||||||
|
// autoReconnect will auto AddApp
|
||||||
|
// pn.AddApp(config)
|
||||||
|
// TODO: report result
|
||||||
|
default:
|
||||||
|
pn.msgMapMtx.Lock()
|
||||||
|
ch := pn.msgMap[pushHead.From]
|
||||||
|
pn.msgMapMtx.Unlock()
|
||||||
|
ch <- msg
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -43,7 +43,7 @@ func natTest(serverHost string, serverPort int, localPort int, echoPort int) (pu
|
|||||||
// testing for public ip
|
// testing for public ip
|
||||||
if echoPort != 0 {
|
if echoPort != 0 {
|
||||||
for {
|
for {
|
||||||
gLog.Printf(LevelINFO, "public ip test start %s:%d", natRsp.IP, echoPort)
|
gLog.Printf(LevelDEBUG, "public ip test start %s:%d", natRsp.IP, echoPort)
|
||||||
conn, err := net.ListenUDP("udp", nil)
|
conn, err := net.ListenUDP("udp", nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
break
|
break
|
||||||
@@ -60,10 +60,10 @@ func natTest(serverHost string, serverPort int, localPort int, echoPort int) (pu
|
|||||||
conn.SetReadDeadline(time.Now().Add(PublicIPEchoTimeout))
|
conn.SetReadDeadline(time.Now().Add(PublicIPEchoTimeout))
|
||||||
_, _, err = conn.ReadFromUDP(buf)
|
_, _, err = conn.ReadFromUDP(buf)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
gLog.Println(LevelINFO, "public ip:YES")
|
gLog.Println(LevelDEBUG, "public ip:YES")
|
||||||
natRsp.IsPublicIP = 1
|
natRsp.IsPublicIP = 1
|
||||||
} else {
|
} else {
|
||||||
gLog.Println(LevelINFO, "public ip:NO")
|
gLog.Println(LevelDEBUG, "public ip:NO")
|
||||||
}
|
}
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|||||||
+41
-85
@@ -11,7 +11,8 @@ import (
|
|||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
rand.Seed(time.Now().UnixNano())
|
rand.Seed(time.Now().UnixNano())
|
||||||
|
binDir := filepath.Dir(os.Args[0])
|
||||||
|
os.Chdir(binDir) // for system service
|
||||||
// TODO: install sub command, deamon process
|
// TODO: install sub command, deamon process
|
||||||
// groups := flag.String("groups", "", "you could join in several groups. like: GroupName1:Password1;GroupName2:Password2; group name 8-31 characters")
|
// groups := flag.String("groups", "", "you could join in several groups. like: GroupName1:Password1;GroupName2:Password2; group name 8-31 characters")
|
||||||
if len(os.Args) > 1 {
|
if len(os.Args) > 1 {
|
||||||
@@ -51,36 +52,15 @@ func main() {
|
|||||||
dstPort := flag.Int("dstport", 0, "destination port ")
|
dstPort := flag.Int("dstport", 0, "destination port ")
|
||||||
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")
|
||||||
noShare := flag.Bool("noshare", false, "disable using the huge numbers of shared nodes in OpenP2P network, your connectivity will be weak. also this node will not shared with others")
|
appName := flag.String("appname", "", "app name")
|
||||||
|
flag.Bool("noshare", false, "deprecated. uses -sharebandwidth -1")
|
||||||
shareBandwidth := flag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private node no limit")
|
shareBandwidth := flag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private node no limit")
|
||||||
configFile := flag.Bool("f", false, "config file")
|
flag.Bool("f", false, "deprecated. config file")
|
||||||
daemonMode := flag.Bool("d", false, "daemonMode")
|
daemonMode := flag.Bool("d", false, "daemonMode")
|
||||||
byDaemon := flag.Bool("bydaemon", false, "start by daemon")
|
byDaemon := flag.Bool("bydaemon", false, "start by daemon")
|
||||||
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()
|
||||||
|
|
||||||
gLog = InitLogger(filepath.Dir(os.Args[0]), "openp2p", LogLevel(*logLevel), 1024*1024, LogFileAndConsole)
|
|
||||||
gLog.Println(LevelINFO, "openp2p start. version: ", OpenP2PVersion)
|
|
||||||
if *daemonMode {
|
|
||||||
d := daemon{}
|
|
||||||
d.run()
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if !*configFile {
|
|
||||||
// validate cmd params
|
|
||||||
checkParams(*node, *user, *password)
|
|
||||||
if *peerNode != "" {
|
|
||||||
if *dstPort == 0 {
|
|
||||||
gLog.Println(LevelERROR, "dstPort not set")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if *srcPort == 0 {
|
|
||||||
gLog.Println(LevelERROR, "srcPort not set")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
config := AppConfig{}
|
config := AppConfig{}
|
||||||
config.PeerNode = *peerNode
|
config.PeerNode = *peerNode
|
||||||
config.PeerUser = *peerUser
|
config.PeerUser = *peerUser
|
||||||
@@ -89,73 +69,49 @@ func main() {
|
|||||||
config.DstPort = *dstPort
|
config.DstPort = *dstPort
|
||||||
config.SrcPort = *srcPort
|
config.SrcPort = *srcPort
|
||||||
config.Protocol = *protocol
|
config.Protocol = *protocol
|
||||||
gLog.Println(LevelINFO, config)
|
config.AppName = *appName
|
||||||
if *configFile {
|
// add command config first
|
||||||
if err := gConf.load(); err != nil {
|
gConf.add(config)
|
||||||
gLog.Println(LevelERROR, "load config error. exit.")
|
gConf.load()
|
||||||
return
|
gConf.mtx.Lock()
|
||||||
}
|
|
||||||
} else {
|
|
||||||
gConf.add(config)
|
|
||||||
gConf.Network = NetworkConfig{
|
|
||||||
Node: *node,
|
|
||||||
User: *user,
|
|
||||||
Password: *password,
|
|
||||||
NoShare: *noShare,
|
|
||||||
ServerHost: *serverHost,
|
|
||||||
ServerPort: 27182,
|
|
||||||
UDPPort1: 27182,
|
|
||||||
UDPPort2: 27183,
|
|
||||||
ipv6: "240e:3b7:621:def0:fda4:dd7f:36a1:2803", // TODO: detect real ipv6
|
|
||||||
ShareBandwidth: *shareBandwidth,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
// gConf.save() // not change config file
|
|
||||||
gConf.daemonMode = *byDaemon
|
|
||||||
|
|
||||||
gLog.Println(LevelINFO, gConf)
|
flag.Visit(func(f *flag.Flag) {
|
||||||
|
if f.Name == "sharebandwidth" {
|
||||||
|
gConf.Network.ShareBandwidth = *shareBandwidth
|
||||||
|
}
|
||||||
|
if f.Name == "node" {
|
||||||
|
gConf.Network.Node = *node
|
||||||
|
}
|
||||||
|
if f.Name == "user" {
|
||||||
|
gConf.Network.User = *user
|
||||||
|
}
|
||||||
|
if f.Name == "password" {
|
||||||
|
gConf.Network.Password = *password
|
||||||
|
}
|
||||||
|
if f.Name == "serverhost" {
|
||||||
|
gConf.Network.ServerHost = *serverHost
|
||||||
|
}
|
||||||
|
if f.Name == "loglevel" {
|
||||||
|
gConf.logLevel = *logLevel
|
||||||
|
}
|
||||||
|
})
|
||||||
|
gLog = InitLogger(binDir, "openp2p", LogLevel(gConf.logLevel), 1024*1024, LogFileAndConsole)
|
||||||
|
gLog.Println(LevelINFO, "openp2p start. version: ", OpenP2PVersion)
|
||||||
|
gConf.mtx.Unlock()
|
||||||
|
gConf.save()
|
||||||
|
gConf.daemonMode = *byDaemon
|
||||||
|
if *daemonMode {
|
||||||
|
d := daemon{}
|
||||||
|
d.run()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
gLog.Println(LevelINFO, &gConf)
|
||||||
setFirewall()
|
setFirewall()
|
||||||
network := P2PNetworkInstance(&gConf.Network)
|
network := P2PNetworkInstance(&gConf.Network)
|
||||||
if ok := network.Connect(30000); !ok {
|
if ok := network.Connect(30000); !ok {
|
||||||
gLog.Println(LevelERROR, "P2PNetwork login error")
|
gLog.Println(LevelERROR, "P2PNetwork login error")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
for _, app := range gConf.Apps {
|
|
||||||
// set default peer user password
|
|
||||||
if app.PeerPassword == "" {
|
|
||||||
app.PeerPassword = gConf.Network.Password
|
|
||||||
}
|
|
||||||
if app.PeerUser == "" {
|
|
||||||
app.PeerUser = gConf.Network.User
|
|
||||||
}
|
|
||||||
err := network.AddApp(app)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Println(LevelERROR, "addTunnel error")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// test
|
|
||||||
// go func() {
|
|
||||||
|
|
||||||
// time.Sleep(time.Second * 30)
|
|
||||||
// config := AppConfig{}
|
|
||||||
// config.PeerNode = *peerNode
|
|
||||||
// config.PeerUser = *peerUser
|
|
||||||
// config.PeerPassword = *peerPassword
|
|
||||||
// config.DstHost = *dstIP
|
|
||||||
// config.DstPort = *dstPort
|
|
||||||
// config.SrcPort = 32
|
|
||||||
// config.Protocol = *protocol
|
|
||||||
// network.AddApp(config)
|
|
||||||
// // time.Sleep(time.Second * 30)
|
|
||||||
// // network.DeleteTunnel(config)
|
|
||||||
// // time.Sleep(time.Second * 30)
|
|
||||||
// // network.DeleteTunnel(config)
|
|
||||||
// }()
|
|
||||||
|
|
||||||
// // TODO: http api
|
|
||||||
// api := ClientAPI{}
|
|
||||||
// go api.run()
|
|
||||||
gLog.Println(LevelINFO, "waiting for connection...")
|
gLog.Println(LevelINFO, "waiting for connection...")
|
||||||
forever := make(chan bool)
|
forever := make(chan bool)
|
||||||
<-forever
|
<-forever
|
||||||
|
|||||||
+2
-2
@@ -21,8 +21,8 @@ type overlayTCP struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (otcp *overlayTCP) run() {
|
func (otcp *overlayTCP) run() {
|
||||||
gLog.Printf(LevelINFO, "%d overlayTCP run start", otcp.id)
|
gLog.Printf(LevelDEBUG, "%d overlayTCP run start", otcp.id)
|
||||||
defer gLog.Printf(LevelINFO, "%d overlayTCP run end", otcp.id)
|
defer gLog.Printf(LevelDEBUG, "%d overlayTCP run end", otcp.id)
|
||||||
otcp.running = true
|
otcp.running = true
|
||||||
buffer := make([]byte, ReadBuffLen+PaddingSize)
|
buffer := make([]byte, ReadBuffLen+PaddingSize)
|
||||||
readBuf := buffer[:ReadBuffLen]
|
readBuf := buffer[:ReadBuffLen]
|
||||||
|
|||||||
@@ -11,16 +11,17 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type p2pApp struct {
|
type p2pApp struct {
|
||||||
config AppConfig
|
config AppConfig
|
||||||
listener net.Listener
|
listener net.Listener
|
||||||
tunnel *P2PTunnel
|
tunnel *P2PTunnel
|
||||||
rtid uint64
|
rtid uint64
|
||||||
hbTime time.Time
|
relayNode string
|
||||||
hbMtx sync.Mutex
|
hbTime time.Time
|
||||||
running bool
|
hbMtx sync.Mutex
|
||||||
id uint64
|
running bool
|
||||||
key uint64
|
id uint64
|
||||||
wg sync.WaitGroup
|
key uint64
|
||||||
|
wg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
||||||
func (app *p2pApp) isActive() bool {
|
func (app *p2pApp) isActive() bool {
|
||||||
@@ -72,7 +73,7 @@ func (app *p2pApp) listenTCP() error {
|
|||||||
otcp.appKeyBytes = encryptKey
|
otcp.appKeyBytes = encryptKey
|
||||||
}
|
}
|
||||||
app.tunnel.overlayConns.Store(otcp.id, &otcp)
|
app.tunnel.overlayConns.Store(otcp.id, &otcp)
|
||||||
gLog.Printf(LevelINFO, "Accept overlayID:%d", otcp.id)
|
gLog.Printf(LevelDEBUG, "Accept overlayID:%d", otcp.id)
|
||||||
// tell peer connect
|
// tell peer connect
|
||||||
req := OverlayConnectReq{ID: otcp.id,
|
req := OverlayConnectReq{ID: otcp.id,
|
||||||
User: app.config.PeerUser,
|
User: app.config.PeerUser,
|
||||||
|
|||||||
+84
-184
@@ -37,7 +37,7 @@ type P2PNetwork struct {
|
|||||||
msgMapMtx sync.Mutex
|
msgMapMtx sync.Mutex
|
||||||
config NetworkConfig
|
config NetworkConfig
|
||||||
allTunnels sync.Map
|
allTunnels sync.Map
|
||||||
apps sync.Map
|
apps sync.Map //key: protocol+srcport; value: p2pApp
|
||||||
limiter *BandwidthLimiter
|
limiter *BandwidthLimiter
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -63,7 +63,7 @@ func P2PNetworkInstance(config *NetworkConfig) *P2PNetwork {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) run() {
|
func (pn *P2PNetwork) run() {
|
||||||
go pn.autoReconnectApp()
|
go pn.autorunApp()
|
||||||
heartbeatTimer := time.NewTicker(NetworkHeartbeatTime)
|
heartbeatTimer := time.NewTicker(NetworkHeartbeatTime)
|
||||||
for pn.running {
|
for pn.running {
|
||||||
select {
|
select {
|
||||||
@@ -93,55 +93,61 @@ func (pn *P2PNetwork) Connect(timeout int) bool {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) autoReconnectApp() {
|
func (pn *P2PNetwork) runAll() {
|
||||||
gLog.Println(LevelINFO, "autoReconnectApp start")
|
gConf.mtx.Lock()
|
||||||
retryApps := make([]AppConfig, 0)
|
defer gConf.mtx.Unlock()
|
||||||
|
for _, config := range gConf.Apps {
|
||||||
|
// set default peer user password
|
||||||
|
if config.PeerPassword == "" {
|
||||||
|
config.PeerPassword = gConf.Network.Password
|
||||||
|
}
|
||||||
|
if config.PeerUser == "" {
|
||||||
|
config.PeerUser = gConf.Network.User
|
||||||
|
}
|
||||||
|
if config.AppName == "" {
|
||||||
|
config.AppName = fmt.Sprintf("%s%d", config.Protocol, config.SrcPort)
|
||||||
|
}
|
||||||
|
appExist := false
|
||||||
|
appActive := false
|
||||||
|
i, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
|
if ok {
|
||||||
|
app := i.(*p2pApp)
|
||||||
|
appExist = true
|
||||||
|
if app.isActive() {
|
||||||
|
appActive = true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if appExist && appActive {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if appExist && !appActive {
|
||||||
|
gLog.Printf(LevelINFO, "detect app %s disconnect, reconnecting...", config.AppName)
|
||||||
|
pn.DeleteApp(config)
|
||||||
|
if config.retryTime.Add(time.Minute * 15).Before(time.Now()) {
|
||||||
|
config.retryNum = 0
|
||||||
|
}
|
||||||
|
config.retryNum++
|
||||||
|
config.retryTime = time.Now()
|
||||||
|
if config.retryNum > MaxRetry {
|
||||||
|
gLog.Printf(LevelERROR, "app %s%d retry more than %d times, exit.", config.Protocol, config.SrcPort, MaxRetry)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
}
|
||||||
|
go pn.AddApp(config)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
func (pn *P2PNetwork) autorunApp() {
|
||||||
|
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 {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if len(retryApps) > 0 {
|
pn.runAll()
|
||||||
gLog.Printf(LevelINFO, "retryApps len=%d", len(retryApps))
|
time.Sleep(time.Second * 10)
|
||||||
thisRound := make([]AppConfig, 0)
|
|
||||||
for i := 0; i < len(retryApps); i++ {
|
|
||||||
// reset retryNum when running 15min continuously
|
|
||||||
if retryApps[i].retryTime.Add(time.Minute * 15).Before(time.Now()) {
|
|
||||||
retryApps[i].retryNum = 0
|
|
||||||
}
|
|
||||||
retryApps[i].retryNum++
|
|
||||||
retryApps[i].retryTime = time.Now()
|
|
||||||
if retryApps[i].retryNum > MaxRetry {
|
|
||||||
gLog.Printf(LevelERROR, "app %s%d retry more than %d times, exit.", retryApps[i].Protocol, retryApps[i].SrcPort, MaxRetry)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
pn.DeleteApp(retryApps[i])
|
|
||||||
if err := pn.AddApp(retryApps[i]); err != nil {
|
|
||||||
gLog.Printf(LevelERROR, "AddApp %s%d error:%s", retryApps[i].Protocol, retryApps[i].SrcPort, err)
|
|
||||||
thisRound = append(thisRound, retryApps[i])
|
|
||||||
time.Sleep(RetryInterval)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
retryApps = thisRound
|
|
||||||
}
|
|
||||||
pn.apps.Range(func(_, i interface{}) bool {
|
|
||||||
app := i.(*p2pApp)
|
|
||||||
if app.isActive() {
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
gLog.Printf(LevelINFO, "detect app %s%d disconnect,last hb %s reconnecting...", app.config.Protocol, app.config.SrcPort, app.hbTime)
|
|
||||||
config := app.config
|
|
||||||
// clear peerinfo
|
|
||||||
config.peerConeNatPort = 0
|
|
||||||
config.peerIP = ""
|
|
||||||
config.peerNatType = 0
|
|
||||||
config.peerToken = 0
|
|
||||||
pn.DeleteApp(config)
|
|
||||||
retryApps = append(retryApps, config)
|
|
||||||
return true
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
gLog.Println(LevelINFO, "autoReconnectApp 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, error) {
|
||||||
@@ -198,21 +204,17 @@ func (pn *P2PNetwork) addRelayTunnel(config AppConfig, appid uint64, appkey uint
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
||||||
gLog.Printf(LevelINFO, "addApp %s%d to %s:%s:%d start", config.Protocol, config.SrcPort, 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%d to %s:%s:%d end", config.Protocol, config.SrcPort, 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)
|
||||||
if !pn.online {
|
if !pn.online {
|
||||||
return errors.New("P2PNetwork offline")
|
return errors.New("P2PNetwork offline")
|
||||||
}
|
}
|
||||||
// check if app already exist?
|
// check if app already exist?
|
||||||
appExist := false
|
appExist := false
|
||||||
pn.apps.Range(func(_, i interface{}) bool {
|
_, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
app := i.(*p2pApp)
|
if ok {
|
||||||
if app.config.Protocol == config.Protocol && app.config.SrcPort == config.SrcPort {
|
appExist = true
|
||||||
appExist = true
|
}
|
||||||
return false
|
|
||||||
}
|
|
||||||
return true
|
|
||||||
})
|
|
||||||
if appExist {
|
if appExist {
|
||||||
return errors.New("P2PApp already exist")
|
return errors.New("P2PApp already exist")
|
||||||
}
|
}
|
||||||
@@ -221,7 +223,7 @@ func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
|||||||
t, err := pn.addDirectTunnel(config, 0)
|
t, err := pn.addDirectTunnel(config, 0)
|
||||||
var rtid uint64
|
var rtid uint64
|
||||||
relayNode := ""
|
relayNode := ""
|
||||||
peerNatType := 100
|
peerNatType := NATUnknown
|
||||||
peerIP := ""
|
peerIP := ""
|
||||||
errMsg := ""
|
errMsg := ""
|
||||||
if err != nil && err == ErrorHandshake {
|
if err != nil && err == ErrorHandshake {
|
||||||
@@ -257,13 +259,14 @@ func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
|||||||
pn.write(MsgReport, MsgReportConnect, &req)
|
pn.write(MsgReport, MsgReportConnect, &req)
|
||||||
|
|
||||||
app := p2pApp{
|
app := p2pApp{
|
||||||
id: appID,
|
id: appID,
|
||||||
key: appKey,
|
key: appKey,
|
||||||
tunnel: t,
|
tunnel: t,
|
||||||
config: config,
|
config: config,
|
||||||
rtid: rtid,
|
rtid: rtid,
|
||||||
hbTime: time.Now()}
|
relayNode: relayNode,
|
||||||
pn.apps.Store(appID, &app)
|
hbTime: time.Now()}
|
||||||
|
pn.apps.Store(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort), &app)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
go app.listen()
|
go app.listen()
|
||||||
}
|
}
|
||||||
@@ -274,22 +277,18 @@ func (pn *P2PNetwork) DeleteApp(config AppConfig) {
|
|||||||
gLog.Printf(LevelINFO, "DeleteApp %s%d start", config.Protocol, config.SrcPort)
|
gLog.Printf(LevelINFO, "DeleteApp %s%d start", config.Protocol, config.SrcPort)
|
||||||
defer gLog.Printf(LevelINFO, "DeleteApp %s%d end", config.Protocol, config.SrcPort)
|
defer gLog.Printf(LevelINFO, "DeleteApp %s%d end", config.Protocol, config.SrcPort)
|
||||||
// close the apps of this config
|
// close the apps of this config
|
||||||
pn.apps.Range(func(_, i interface{}) bool {
|
i, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
|
if ok {
|
||||||
app := i.(*p2pApp)
|
app := i.(*p2pApp)
|
||||||
if app.config.Protocol == config.Protocol && app.config.SrcPort == config.SrcPort {
|
gLog.Printf(LevelINFO, "app %s exist, delete it", fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
gLog.Printf(LevelINFO, "app %s exist, delete it", fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
app.close()
|
||||||
app := i.(*p2pApp)
|
pn.apps.Delete(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
app.close()
|
}
|
||||||
pn.apps.Delete(app.id)
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
return true
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) addDirectTunnel(config AppConfig, tid uint64) (*P2PTunnel, error) {
|
func (pn *P2PNetwork) addDirectTunnel(config AppConfig, tid uint64) (*P2PTunnel, error) {
|
||||||
gLog.Printf(LevelINFO, "addDirectTunnel %s%d to %s:%s:%d start", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
gLog.Printf(LevelDEBUG, "addDirectTunnel %s%d to %s:%s:%d start", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
||||||
defer gLog.Printf(LevelINFO, "addDirectTunnel %s%d to %s:%s:%d end", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
defer gLog.Printf(LevelDEBUG, "addDirectTunnel %s%d to %s:%s:%d end", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
||||||
isClient := false
|
isClient := false
|
||||||
// client side tid=0, assign random uint64
|
// client side tid=0, assign random uint64
|
||||||
if tid == 0 {
|
if tid == 0 {
|
||||||
@@ -377,10 +376,10 @@ func (pn *P2PNetwork) init() error {
|
|||||||
pn.config.natType = NATSymmetric
|
pn.config.natType = NATSymmetric
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelINFO, "detect NAT type error:", err)
|
gLog.Println(LevelDEBUG, "detect NAT type error:", err)
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
gLog.Println(LevelINFO, "detect NAT type:", pn.config.natType, " publicIP:", pn.config.publicIP)
|
gLog.Println(LevelDEBUG, "detect NAT type:", pn.config.natType, " publicIP:", pn.config.publicIP)
|
||||||
gatewayURL := fmt.Sprintf("%s:%d", pn.config.ServerHost, pn.config.ServerPort)
|
gatewayURL := fmt.Sprintf("%s:%d", pn.config.ServerHost, pn.config.ServerPort)
|
||||||
forwardPath := "/openp2p/v1/login"
|
forwardPath := "/openp2p/v1/login"
|
||||||
config := tls.Config{InsecureSkipVerify: true} // let's encrypt root cert "DST Root CA X3" expired at 2021/09/29. many old system(windows server 2008 etc) will not trust our cert
|
config := tls.Config{InsecureSkipVerify: true} // let's encrypt root cert "DST Root CA X3" expired at 2021/09/29. many old system(windows server 2008 etc) will not trust our cert
|
||||||
@@ -392,12 +391,7 @@ func (pn *P2PNetwork) init() error {
|
|||||||
q.Add("password", pn.config.Password)
|
q.Add("password", pn.config.Password)
|
||||||
q.Add("version", OpenP2PVersion)
|
q.Add("version", OpenP2PVersion)
|
||||||
q.Add("nattype", fmt.Sprintf("%d", pn.config.natType))
|
q.Add("nattype", fmt.Sprintf("%d", pn.config.natType))
|
||||||
|
q.Add("sharebandwidth", fmt.Sprintf("%d", pn.config.ShareBandwidth))
|
||||||
noShareStr := "false"
|
|
||||||
if pn.config.NoShare {
|
|
||||||
noShareStr = "true"
|
|
||||||
}
|
|
||||||
q.Add("noshare", noShareStr)
|
|
||||||
u.RawQuery = q.Encode()
|
u.RawQuery = q.Encode()
|
||||||
var ws *websocket.Conn
|
var ws *websocket.Conn
|
||||||
ws, _, err = websocket.DefaultDialer.Dial(u.String(), nil)
|
ws, _, err = websocket.DefaultDialer.Dial(u.String(), nil)
|
||||||
@@ -425,7 +419,7 @@ func (pn *P2PNetwork) init() error {
|
|||||||
Version: OpenP2PVersion,
|
Version: OpenP2PVersion,
|
||||||
}
|
}
|
||||||
rsp := netInfo()
|
rsp := netInfo()
|
||||||
gLog.Println(LevelINFO, rsp)
|
gLog.Println(LevelDEBUG, "netinfo:", rsp)
|
||||||
if rsp != nil && rsp.Country != "" {
|
if rsp != nil && rsp.Country != "" {
|
||||||
if len(rsp.IP) == net.IPv6len {
|
if len(rsp.IP) == net.IPv6len {
|
||||||
pn.config.ipv6 = rsp.IP.String()
|
pn.config.ipv6 = rsp.IP.String()
|
||||||
@@ -434,7 +428,7 @@ func (pn *P2PNetwork) init() error {
|
|||||||
req.NetInfo = *rsp
|
req.NetInfo = *rsp
|
||||||
}
|
}
|
||||||
pn.write(MsgReport, MsgReportBasic, &req)
|
pn.write(MsgReport, MsgReportBasic, &req)
|
||||||
gLog.Println(LevelINFO, "P2PNetwork init ok")
|
gLog.Println(LevelDEBUG, "P2PNetwork init ok")
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -472,7 +466,7 @@ func (pn *P2PNetwork) handleMessage(t int, msg []byte) {
|
|||||||
case MsgHeartbeat:
|
case MsgHeartbeat:
|
||||||
gLog.Printf(LevelDEBUG, "P2PNetwork heartbeat ok")
|
gLog.Printf(LevelDEBUG, "P2PNetwork heartbeat ok")
|
||||||
case MsgPush:
|
case MsgPush:
|
||||||
pn.handlePush(head.SubType, msg)
|
handlePush(pn, head.SubType, msg)
|
||||||
default:
|
default:
|
||||||
pn.msgMapMtx.Lock()
|
pn.msgMapMtx.Lock()
|
||||||
ch := pn.msgMap[0]
|
ch := pn.msgMap[0]
|
||||||
@@ -483,7 +477,7 @@ func (pn *P2PNetwork) handleMessage(t int, msg []byte) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) readLoop() {
|
func (pn *P2PNetwork) readLoop() {
|
||||||
gLog.Printf(LevelINFO, "P2PNetwork readLoop start")
|
gLog.Printf(LevelDEBUG, "P2PNetwork readLoop start")
|
||||||
pn.wg.Add(1)
|
pn.wg.Add(1)
|
||||||
defer pn.wg.Done()
|
defer pn.wg.Done()
|
||||||
for pn.running {
|
for pn.running {
|
||||||
@@ -497,7 +491,7 @@ func (pn *P2PNetwork) readLoop() {
|
|||||||
}
|
}
|
||||||
pn.handleMessage(t, msg)
|
pn.handleMessage(t, msg)
|
||||||
}
|
}
|
||||||
gLog.Printf(LevelINFO, "P2PNetwork readLoop end")
|
gLog.Printf(LevelDEBUG, "P2PNetwork readLoop end")
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) write(mainType uint16, subType uint16, packet interface{}) error {
|
func (pn *P2PNetwork) write(mainType uint16, subType uint16, packet interface{}) error {
|
||||||
@@ -592,106 +586,12 @@ func (pn *P2PNetwork) read(node string, mainType uint16, subType uint16, timeout
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) handlePush(subType uint16, msg []byte) error {
|
|
||||||
pushHead := PushHeader{}
|
|
||||||
err := binary.Read(bytes.NewReader(msg[openP2PHeaderSize:openP2PHeaderSize+PushHeaderSize]), binary.LittleEndian, &pushHead)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
gLog.Printf(LevelDEBUG, "handle push msg type:%d, push header:%+v", subType, pushHead)
|
|
||||||
switch subType {
|
|
||||||
case MsgPushConnectReq:
|
|
||||||
req := PushConnectReq{}
|
|
||||||
err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Printf(LevelERROR, "wrong MsgPushConnectReq:%s", err)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
gLog.Printf(LevelINFO, "%s is connecting...", req.From)
|
|
||||||
gLog.Println(LevelDEBUG, "push connect response to ", req.From)
|
|
||||||
// verify token or name&password
|
|
||||||
if VerifyTOTP(req.Token, pn.config.User, pn.config.Password, time.Now().Unix()+(pn.serverTs-pn.localTs)) || // localTs may behind, auto adjust ts
|
|
||||||
VerifyTOTP(req.Token, pn.config.User, pn.config.Password, time.Now().Unix()) ||
|
|
||||||
(req.User == pn.config.User && req.Password == pn.config.Password) {
|
|
||||||
gLog.Printf(LevelINFO, "Access Granted\n")
|
|
||||||
config := AppConfig{}
|
|
||||||
config.peerNatType = req.NatType
|
|
||||||
config.peerConeNatPort = req.ConeNatPort
|
|
||||||
config.peerIP = req.FromIP
|
|
||||||
config.PeerNode = req.From
|
|
||||||
// share relay node will limit bandwidth
|
|
||||||
if req.User != pn.config.User || req.Password != pn.config.Password {
|
|
||||||
gLog.Printf(LevelINFO, "set share bandwidth %d mbps", pn.config.ShareBandwidth)
|
|
||||||
config.shareBandwidth = pn.config.ShareBandwidth
|
|
||||||
}
|
|
||||||
// go pn.AddTunnel(config, req.ID)
|
|
||||||
go pn.addDirectTunnel(config, req.ID)
|
|
||||||
break
|
|
||||||
}
|
|
||||||
gLog.Println(LevelERROR, "Access Denied:", req.From)
|
|
||||||
rsp := PushConnectRsp{
|
|
||||||
Error: 1,
|
|
||||||
Detail: fmt.Sprintf("connect to %s error: Access Denied", pn.config.Node),
|
|
||||||
To: req.From,
|
|
||||||
From: pn.config.Node,
|
|
||||||
}
|
|
||||||
pn.push(req.From, MsgPushConnectRsp, rsp)
|
|
||||||
case MsgPushRsp:
|
|
||||||
rsp := PushRsp{}
|
|
||||||
err := json.Unmarshal(msg[openP2PHeaderSize:], &rsp)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Printf(LevelERROR, "wrong pushRsp:%s", err)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if rsp.Error == 0 {
|
|
||||||
gLog.Printf(LevelDEBUG, "push ok, detail:%s", rsp.Detail)
|
|
||||||
} else {
|
|
||||||
gLog.Printf(LevelERROR, "push error:%d, detail:%s", rsp.Error, rsp.Detail)
|
|
||||||
}
|
|
||||||
case MsgPushAddRelayTunnelReq:
|
|
||||||
req := AddRelayTunnelReq{}
|
|
||||||
err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Printf(LevelERROR, "wrong RelayNodeRsp:%s", err)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
config := AppConfig{}
|
|
||||||
config.PeerNode = req.RelayName
|
|
||||||
config.peerToken = req.RelayToken
|
|
||||||
// set user password, maybe the relay node is your private node
|
|
||||||
config.PeerUser = pn.config.User
|
|
||||||
config.PeerPassword = pn.config.Password
|
|
||||||
go func(r AddRelayTunnelReq) {
|
|
||||||
t, errDt := pn.addDirectTunnel(config, 0)
|
|
||||||
if errDt == nil {
|
|
||||||
// notify peer relay ready
|
|
||||||
msg := TunnelMsg{ID: t.id}
|
|
||||||
pn.push(r.From, MsgPushAddRelayTunnelRsp, msg)
|
|
||||||
SaveKey(req.AppID, req.AppKey)
|
|
||||||
}
|
|
||||||
|
|
||||||
}(req)
|
|
||||||
case MsgPushUpdate:
|
|
||||||
update()
|
|
||||||
if gConf.daemonMode {
|
|
||||||
os.Exit(0)
|
|
||||||
}
|
|
||||||
default:
|
|
||||||
pn.msgMapMtx.Lock()
|
|
||||||
ch := pn.msgMap[pushHead.From]
|
|
||||||
pn.msgMapMtx.Unlock()
|
|
||||||
ch <- msg
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (pn *P2PNetwork) updateAppHeartbeat(appID uint64) {
|
func (pn *P2PNetwork) updateAppHeartbeat(appID uint64) {
|
||||||
pn.apps.Range(func(id, i interface{}) bool {
|
pn.apps.Range(func(id, i interface{}) bool {
|
||||||
key := id.(uint64)
|
app := i.(*p2pApp)
|
||||||
if key != appID {
|
if app.id != appID {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
app := i.(*p2pApp)
|
|
||||||
app.updateHeartbeat()
|
app.updateHeartbeat()
|
||||||
return false
|
return false
|
||||||
})
|
})
|
||||||
|
|||||||
+12
-12
@@ -52,7 +52,7 @@ func (t *P2PTunnel) init() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) connect() error {
|
func (t *P2PTunnel) connect() error {
|
||||||
gLog.Printf(LevelINFO, "start p2pTunnel to %s ", t.config.PeerNode)
|
gLog.Printf(LevelDEBUG, "start p2pTunnel to %s ", t.config.PeerNode)
|
||||||
t.isServer = false
|
t.isServer = false
|
||||||
req := PushConnectReq{
|
req := PushConnectReq{
|
||||||
User: t.config.PeerUser,
|
User: t.config.PeerUser,
|
||||||
@@ -144,7 +144,7 @@ func (t *P2PTunnel) handshake() error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
gLog.Println(LevelINFO, "handshake to ", t.config.PeerNode)
|
gLog.Println(LevelDEBUG, "handshake to ", t.config.PeerNode)
|
||||||
var err error
|
var err error
|
||||||
// TODO: handle NATNone, nodes with public ip has no punching
|
// TODO: handle NATNone, nodes with public ip has no punching
|
||||||
if (t.pn.config.natType == NATCone && t.config.peerNatType == NATCone) || (t.pn.config.natType == NATNone || t.config.peerNatType == NATNone) {
|
if (t.pn.config.natType == NATCone && t.config.peerNatType == NATCone) || (t.pn.config.natType == NATNone || t.config.peerNatType == NATNone) {
|
||||||
@@ -163,7 +163,7 @@ func (t *P2PTunnel) handshake() error {
|
|||||||
gLog.Println(LevelERROR, "punch handshake error:", err)
|
gLog.Println(LevelERROR, "punch handshake error:", err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
gLog.Printf(LevelINFO, "handshake to %s ok", t.config.PeerNode)
|
gLog.Printf(LevelDEBUG, "handshake to %s ok", t.config.PeerNode)
|
||||||
err = t.run()
|
err = t.run()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, err)
|
gLog.Println(LevelERROR, err)
|
||||||
@@ -198,7 +198,7 @@ func (t *P2PTunnel) run() error {
|
|||||||
gLog.Println(LevelDEBUG, string(buff))
|
gLog.Println(LevelDEBUG, string(buff))
|
||||||
}
|
}
|
||||||
qConn.WriteBytes(MsgP2P, MsgTunnelHandshakeAck, []byte("OpenP2P,hello2"))
|
qConn.WriteBytes(MsgP2P, MsgTunnelHandshakeAck, []byte("OpenP2P,hello2"))
|
||||||
gLog.Println(LevelINFO, "quic connection ok")
|
gLog.Println(LevelDEBUG, "quic connection ok")
|
||||||
t.conn = qConn
|
t.conn = qConn
|
||||||
t.setRun(true)
|
t.setRun(true)
|
||||||
go t.readLoop()
|
go t.readLoop()
|
||||||
@@ -216,7 +216,7 @@ func (t *P2PTunnel) run() error {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
t.pn.read(t.config.PeerNode, MsgPush, MsgPushQuicConnect, time.Second*5)
|
t.pn.read(t.config.PeerNode, MsgPush, MsgPushQuicConnect, time.Second*5)
|
||||||
gLog.Println(LevelINFO, "quic dial to ", t.ra.String())
|
gLog.Println(LevelDEBUG, "quic dial to ", t.ra.String())
|
||||||
qConn, e := dialQuic(conn, t.ra, TunnelIdleTimeout)
|
qConn, e := dialQuic(conn, t.ra, TunnelIdleTimeout)
|
||||||
if e != nil {
|
if e != nil {
|
||||||
return fmt.Errorf("quic dial to %s error:%s", t.ra.String(), e)
|
return fmt.Errorf("quic dial to %s error:%s", t.ra.String(), e)
|
||||||
@@ -233,7 +233,7 @@ func (t *P2PTunnel) run() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
gLog.Println(LevelINFO, "rtt=", time.Since(handshakeBegin))
|
gLog.Println(LevelINFO, "rtt=", time.Since(handshakeBegin))
|
||||||
gLog.Println(LevelINFO, "quic connection ok")
|
gLog.Println(LevelDEBUG, "quic connection ok")
|
||||||
t.conn = qConn
|
t.conn = qConn
|
||||||
t.setRun(true)
|
t.setRun(true)
|
||||||
go t.readLoop()
|
go t.readLoop()
|
||||||
@@ -243,7 +243,7 @@ func (t *P2PTunnel) run() error {
|
|||||||
|
|
||||||
func (t *P2PTunnel) readLoop() {
|
func (t *P2PTunnel) readLoop() {
|
||||||
decryptData := make([]byte, ReadBuffLen+PaddingSize) // 16 bytes for padding
|
decryptData := make([]byte, ReadBuffLen+PaddingSize) // 16 bytes for padding
|
||||||
gLog.Printf(LevelINFO, "%d tunnel readloop start", t.id)
|
gLog.Printf(LevelDEBUG, "%d tunnel readloop start", t.id)
|
||||||
for t.isRuning() {
|
for t.isRuning() {
|
||||||
t.conn.SetReadDeadline(time.Now().Add(TunnelIdleTimeout))
|
t.conn.SetReadDeadline(time.Now().Add(TunnelIdleTimeout))
|
||||||
head, body, err := t.conn.ReadMessage()
|
head, body, err := t.conn.ReadMessage()
|
||||||
@@ -333,7 +333,7 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
overlayID := req.ID
|
overlayID := req.ID
|
||||||
gLog.Printf(LevelINFO, "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" {
|
if req.Protocol == "tcp" {
|
||||||
conn, err := net.DialTimeout("tcp", fmt.Sprintf("%s:%d", req.DstIP, req.DstPort), time.Second*5)
|
conn, err := net.DialTimeout("tcp", fmt.Sprintf("%s:%d", req.DstIP, req.DstPort), time.Second*5)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -368,7 +368,7 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
overlayID := req.ID
|
overlayID := req.ID
|
||||||
gLog.Printf(LevelINFO, "%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)
|
otcp := i.(*overlayTCP)
|
||||||
@@ -379,13 +379,13 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
}
|
}
|
||||||
t.setRun(false)
|
t.setRun(false)
|
||||||
t.conn.Close()
|
t.conn.Close()
|
||||||
gLog.Printf(LevelINFO, "%d tunnel readloop end", t.id)
|
gLog.Printf(LevelDEBUG, "%d tunnel readloop end", t.id)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) writeLoop() {
|
func (t *P2PTunnel) writeLoop() {
|
||||||
tc := time.NewTicker(TunnelHeartbeatTime)
|
tc := time.NewTicker(TunnelHeartbeatTime)
|
||||||
defer tc.Stop()
|
defer tc.Stop()
|
||||||
defer gLog.Printf(LevelINFO, "%d tunnel writeloop end", t.id)
|
defer gLog.Printf(LevelDEBUG, "%d tunnel writeloop end", t.id)
|
||||||
for t.isRuning() {
|
for t.isRuning() {
|
||||||
select {
|
select {
|
||||||
case <-tc.C:
|
case <-tc.C:
|
||||||
@@ -402,7 +402,7 @@ func (t *P2PTunnel) writeLoop() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) listen() error {
|
func (t *P2PTunnel) listen() error {
|
||||||
gLog.Printf(LevelINFO, "p2ptunnel wait for connecting")
|
gLog.Printf(LevelDEBUG, "p2ptunnel wait for connecting")
|
||||||
t.isServer = true
|
t.isServer = true
|
||||||
return t.handshake()
|
return t.handshake()
|
||||||
}
|
}
|
||||||
|
|||||||
+30
-1
@@ -10,7 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
const OpenP2PVersion = "0.97.1"
|
const OpenP2PVersion = "0.98.0"
|
||||||
const ProducnName string = "openp2p"
|
const ProducnName string = "openp2p"
|
||||||
|
|
||||||
type openP2PHeader struct {
|
type openP2PHeader struct {
|
||||||
@@ -79,6 +79,7 @@ const (
|
|||||||
MsgPushUpdate = 6
|
MsgPushUpdate = 6
|
||||||
MsgPushReportApps = 7
|
MsgPushReportApps = 7
|
||||||
MsgPushQuicConnect = 8
|
MsgPushQuicConnect = 8
|
||||||
|
MsgPushEditApp = 9
|
||||||
)
|
)
|
||||||
|
|
||||||
// MsgP2P sub type message
|
// MsgP2P sub type message
|
||||||
@@ -109,6 +110,7 @@ const (
|
|||||||
MsgReportBasic = iota
|
MsgReportBasic = iota
|
||||||
MsgReportQuery
|
MsgReportQuery
|
||||||
MsgReportConnect
|
MsgReportConnect
|
||||||
|
MsgReportApps
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -128,6 +130,7 @@ const (
|
|||||||
RetryInterval = time.Second * 30
|
RetryInterval = time.Second * 30
|
||||||
PublicIPEchoTimeout = time.Second * 3
|
PublicIPEchoTimeout = time.Second * 3
|
||||||
NatTestTimeout = time.Second * 10
|
NatTestTimeout = time.Second * 10
|
||||||
|
ClientAPITimeout = time.Second * 10
|
||||||
)
|
)
|
||||||
|
|
||||||
// NATNone has public ip
|
// NATNone has public ip
|
||||||
@@ -135,6 +138,7 @@ const (
|
|||||||
NATNone = 0
|
NATNone = 0
|
||||||
NATCone = 1
|
NATCone = 1
|
||||||
NATSymmetric = 2
|
NATSymmetric = 2
|
||||||
|
NATUnknown = 314
|
||||||
)
|
)
|
||||||
|
|
||||||
func newMessage(mainType uint16, subType uint16, packet interface{}) ([]byte, error) {
|
func newMessage(mainType uint16, subType uint16, packet interface{}) ([]byte, error) {
|
||||||
@@ -271,6 +275,31 @@ type ReportConnect struct {
|
|||||||
Version string `json:"version,omitempty"`
|
Version string `json:"version,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type AppInfo struct {
|
||||||
|
AppName string `json:"appName,omitempty"`
|
||||||
|
Error string `json:"error,omitempty"`
|
||||||
|
Protocol string `json:"protocol,omitempty"`
|
||||||
|
SrcPort int `json:"srcPort,omitempty"`
|
||||||
|
Protocol0 string `json:"protocol0,omitempty"`
|
||||||
|
SrcPort0 int `json:"srcPort0,omitempty"`
|
||||||
|
NatType int `json:"natType,omitempty"`
|
||||||
|
PeerNode string `json:"peerNode,omitempty"`
|
||||||
|
DstPort int `json:"dstPort,omitempty"`
|
||||||
|
DstHost string `json:"dstHost,omitempty"`
|
||||||
|
PeerUser string `json:"peerUser,omitempty"`
|
||||||
|
PeerNatType int `json:"peerNatType,omitempty"`
|
||||||
|
PeerIP string `json:"peerIP,omitempty"`
|
||||||
|
ShareBandwidth int `json:"shareBandWidth,omitempty"`
|
||||||
|
RelayNode string `json:"relayNode,omitempty"`
|
||||||
|
Version string `json:"version,omitempty"`
|
||||||
|
RetryTime string `json:"retryTime,omitempty"`
|
||||||
|
IsActive int `json:"isActive,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type ReportApps struct {
|
||||||
|
Apps []AppInfo
|
||||||
|
}
|
||||||
|
|
||||||
type UpdateInfo struct {
|
type UpdateInfo struct {
|
||||||
Error int `json:"error,omitempty"`
|
Error int `json:"error,omitempty"`
|
||||||
ErrorDetail string `json:"errorDetail,omitempty"`
|
ErrorDetail string `json:"errorDetail,omitempty"`
|
||||||
|
|||||||
Reference in New Issue
Block a user