Compare commits

...
1 Commits
Author SHA1 Message Date
TenderIronh 029d69869f refactor autorunApp and add api for web 2021-12-30 11:18:05 +08:00
14 changed files with 399 additions and 316 deletions
+1 -1
View File
@@ -45,7 +45,7 @@ P2P直连可以让你的设备跑满带宽。不论你的设备在任何网络
![image](/doc/images/officelisten.png) ![image](/doc/images/officelisten.png)
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: **切记将标记大写的参数改成自己的**
+1 -1
View File
@@ -56,7 +56,7 @@ Under the outbreak of covid-19 pandemic, surely remote work becomes a fundamenta
3. Download OpenP2P on your home deviceunzip and execute below command line. 3. Download OpenP2P on your home deviceunzip 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
View File
@@ -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远程桌面3389Linux ssh 22 >* -dstport: 目标服务端口,常见的如windows远程桌面3389Linux ssh 22
+3 -3
View File
@@ -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 peoples 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 peoples 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
+25 -2
View File
@@ -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
+9 -8
View File
@@ -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
View File
@@ -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
}
+3 -3
View File
@@ -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
View File
@@ -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
View File
@@ -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]
+12 -11
View File
@@ -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
View File
@@ -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
View File
@@ -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
View File
@@ -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"`