Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f015b828fc | ||
|
|
df1e16e708 | ||
|
|
c68094cc12 | ||
|
|
a0df0b1e95 | ||
|
|
9c3d557f5d | ||
|
|
2dea3a718d |
@@ -21,3 +21,4 @@ wintun.dll
|
|||||||
app/.idea/
|
app/.idea/
|
||||||
*_debug_bin*
|
*_debug_bin*
|
||||||
cmd/openp2p
|
cmd/openp2p
|
||||||
|
vendor/
|
||||||
+2
-2
@@ -31,7 +31,7 @@ P2P直连可以让你的设备跑满带宽。不论你的设备在任何网络
|
|||||||
## 快速入门
|
## 快速入门
|
||||||
仅需简单4步就能用起来。
|
仅需简单4步就能用起来。
|
||||||
下面是一个远程办公例子:在家里连入办公室Windows电脑。
|
下面是一个远程办公例子:在家里连入办公室Windows电脑。
|
||||||
(另外一个快速入门视频 https://www.bilibili.com/video/BV1Et4y1P7bF/)
|
(另外一个快速入门视频 <https://www.bilibili.com/video/BV1Et4y1P7bF/>)
|
||||||
### 1.注册
|
### 1.注册
|
||||||
前往<https://console.openp2p.cn> 注册新用户,暂无需任何认证
|
前往<https://console.openp2p.cn> 注册新用户,暂无需任何认证
|
||||||
|
|
||||||
@@ -96,7 +96,7 @@ Windows默认会阻止没有花钱买它家证书签名过的程序,选择“
|
|||||||
服务端有个调度模型,根据带宽、ping值、稳定性、服务时长,尽可能地使共享节点均匀地提供服务。连接共享节点使用TOTP密码,hmac-sha256算法校验,它是一次性密码,和我们平时使用的手机验证码或银行密码器一样的原理。
|
服务端有个调度模型,根据带宽、ping值、稳定性、服务时长,尽可能地使共享节点均匀地提供服务。连接共享节点使用TOTP密码,hmac-sha256算法校验,它是一次性密码,和我们平时使用的手机验证码或银行密码器一样的原理。
|
||||||
|
|
||||||
## 编译
|
## 编译
|
||||||
go version go1.18.1+
|
go version 1.20 only (支持win7)
|
||||||
cd到代码根目录,执行
|
cd到代码根目录,执行
|
||||||
```
|
```
|
||||||
make
|
make
|
||||||
|
|||||||
@@ -103,7 +103,7 @@ That's right, the relay node is naturally an man-in-middle, so AES encryption is
|
|||||||
The server side has a scheduling model, which calculate bandwith, ping value,stability and service duration to provide a well-proportioned service to every share node. It uses TOTP(Time-based One-time Password) with hmac-sha256 algorithem, its theory as same as the cellphone validation code or bank cipher coder.
|
The server side has a scheduling model, which calculate bandwith, ping value,stability and service duration to provide a well-proportioned service to every share node. It uses TOTP(Time-based One-time Password) with hmac-sha256 algorithem, its theory as same as the cellphone validation code or bank cipher coder.
|
||||||
|
|
||||||
## Build
|
## Build
|
||||||
go version go1.18.1+
|
go version 1.20 only (support win7)
|
||||||
cd root directory of the socure code and execute
|
cd root directory of the socure code and execute
|
||||||
```
|
```
|
||||||
make
|
make
|
||||||
|
|||||||
+3
-2
@@ -2,9 +2,10 @@
|
|||||||
depends on openjdk 11, gradle 8.1.3, ndk 21
|
depends on openjdk 11, gradle 8.1.3, ndk 21
|
||||||
```
|
```
|
||||||
|
|
||||||
go install golang.org/x/mobile/cmd/gomobile@latest
|
# latest version not support go1.20
|
||||||
|
go install golang.org/x/mobile/cmd/gomobile@7c4916698cc93475ebfea76748ee0faba2deb2a5
|
||||||
gomobile init
|
gomobile init
|
||||||
go get -v golang.org/x/mobile/bind
|
go get -v golang.org/x/mobile/bind@7c4916698cc93475ebfea76748ee0faba2deb2a5
|
||||||
cd core
|
cd core
|
||||||
gomobile bind -target android -v
|
gomobile bind -target android -v
|
||||||
if [[ $? -ne 0 ]]; then
|
if [[ $? -ne 0 ]]; then
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
package openp2p
|
package openp2p
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
"testing"
|
"testing"
|
||||||
)
|
)
|
||||||
@@ -114,3 +115,15 @@ func TestIsIPv6(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestNodeID(t *testing.T) {
|
||||||
|
node1 := "n1-stable"
|
||||||
|
node2 := "tony-stable"
|
||||||
|
nodeID1 := NodeNameToID(node1)
|
||||||
|
nodeID2 := NodeNameToID(node2)
|
||||||
|
if nodeID1 < nodeID2 {
|
||||||
|
fmt.Printf("%s < %s\n", node1, node2)
|
||||||
|
} else {
|
||||||
|
fmt.Printf("%s >= %s\n", node1, node2)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+39
-22
@@ -3,6 +3,7 @@ package openp2p
|
|||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"flag"
|
"flag"
|
||||||
|
"fmt"
|
||||||
"os"
|
"os"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
@@ -28,6 +29,7 @@ type AppConfig struct {
|
|||||||
ForceRelay int // default:0 disable;1 enable
|
ForceRelay int // default:0 disable;1 enable
|
||||||
Enabled int // default:1
|
Enabled int // default:1
|
||||||
// runtime info
|
// runtime info
|
||||||
|
relayMode string // private|public
|
||||||
peerVersion string
|
peerVersion string
|
||||||
peerToken uint64
|
peerToken uint64
|
||||||
peerNatType int
|
peerNatType int
|
||||||
@@ -64,17 +66,25 @@ func (c *AppConfig) ID() uint64 {
|
|||||||
return uint64(c.SrcPort)*10 + 1
|
return uint64(c.SrcPort)*10 + 1
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c *AppConfig) LogPeerNode() string {
|
||||||
|
if c.relayMode == "public" { // memapp
|
||||||
|
return fmt.Sprintf("%d", NodeNameToID(c.PeerNode))
|
||||||
|
}
|
||||||
|
return c.PeerNode
|
||||||
|
}
|
||||||
|
|
||||||
type Config struct {
|
type Config struct {
|
||||||
Network NetworkConfig `json:"network"`
|
Network NetworkConfig `json:"network"`
|
||||||
Apps []*AppConfig `json:"apps"`
|
Apps []*AppConfig `json:"apps"`
|
||||||
|
|
||||||
LogLevel int
|
LogLevel int
|
||||||
|
MaxLogSize int
|
||||||
daemonMode bool
|
daemonMode bool
|
||||||
mtx sync.Mutex
|
mtx sync.Mutex
|
||||||
sdwanMtx sync.Mutex
|
sdwanMtx sync.Mutex
|
||||||
sdwan SDWANInfo
|
sdwan SDWANInfo
|
||||||
delNodes []SDWANNode
|
delNodes []*SDWANNode
|
||||||
addNodes []SDWANNode
|
addNodes []*SDWANNode
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) getSDWAN() SDWANInfo {
|
func (c *Config) getSDWAN() SDWANInfo {
|
||||||
@@ -83,23 +93,30 @@ func (c *Config) getSDWAN() SDWANInfo {
|
|||||||
return c.sdwan
|
return c.sdwan
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) getDelNodes() []SDWANNode {
|
func (c *Config) getDelNodes() []*SDWANNode {
|
||||||
c.sdwanMtx.Lock()
|
c.sdwanMtx.Lock()
|
||||||
defer c.sdwanMtx.Unlock()
|
defer c.sdwanMtx.Unlock()
|
||||||
return c.delNodes
|
return c.delNodes
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) getAddNodes() []SDWANNode {
|
func (c *Config) getAddNodes() []*SDWANNode {
|
||||||
c.sdwanMtx.Lock()
|
c.sdwanMtx.Lock()
|
||||||
defer c.sdwanMtx.Unlock()
|
defer c.sdwanMtx.Unlock()
|
||||||
return c.addNodes
|
return c.addNodes
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c *Config) resetSDWAN() {
|
||||||
|
c.sdwanMtx.Lock()
|
||||||
|
defer c.sdwanMtx.Unlock()
|
||||||
|
c.delNodes = []*SDWANNode{}
|
||||||
|
c.addNodes = []*SDWANNode{}
|
||||||
|
c.sdwan = SDWANInfo{}
|
||||||
|
}
|
||||||
func (c *Config) setSDWAN(s SDWANInfo) {
|
func (c *Config) setSDWAN(s SDWANInfo) {
|
||||||
c.sdwanMtx.Lock()
|
c.sdwanMtx.Lock()
|
||||||
defer c.sdwanMtx.Unlock()
|
defer c.sdwanMtx.Unlock()
|
||||||
// get old-new
|
// get old-new
|
||||||
c.delNodes = []SDWANNode{}
|
c.delNodes = []*SDWANNode{}
|
||||||
for _, oldNode := range c.sdwan.Nodes {
|
for _, oldNode := range c.sdwan.Nodes {
|
||||||
isDeleted := true
|
isDeleted := true
|
||||||
for _, newNode := range s.Nodes {
|
for _, newNode := range s.Nodes {
|
||||||
@@ -113,7 +130,7 @@ func (c *Config) setSDWAN(s SDWANInfo) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
// get new-old
|
// get new-old
|
||||||
c.addNodes = []SDWANNode{}
|
c.addNodes = []*SDWANNode{}
|
||||||
for _, newNode := range s.Nodes {
|
for _, newNode := range s.Nodes {
|
||||||
isNew := true
|
isNew := true
|
||||||
for _, oldNode := range c.sdwan.Nodes {
|
for _, oldNode := range c.sdwan.Nodes {
|
||||||
@@ -147,7 +164,7 @@ func (c *Config) retryApp(peerNode string) {
|
|||||||
GNetwork.apps.Range(func(id, i interface{}) bool {
|
GNetwork.apps.Range(func(id, i interface{}) bool {
|
||||||
app := i.(*p2pApp)
|
app := i.(*p2pApp)
|
||||||
if app.config.PeerNode == peerNode {
|
if app.config.PeerNode == peerNode {
|
||||||
gLog.Println(LvDEBUG, "retry app ", peerNode)
|
gLog.Println(LvDEBUG, "retry app ", app.config.LogPeerNode())
|
||||||
app.config.retryNum = 0
|
app.config.retryNum = 0
|
||||||
app.config.nextRetryTime = time.Now()
|
app.config.nextRetryTime = time.Now()
|
||||||
app.retryRelayNum = 0
|
app.retryRelayNum = 0
|
||||||
@@ -157,7 +174,7 @@ func (c *Config) retryApp(peerNode string) {
|
|||||||
app.hbMtx.Unlock()
|
app.hbMtx.Unlock()
|
||||||
}
|
}
|
||||||
if app.config.RelayNode == peerNode {
|
if app.config.RelayNode == peerNode {
|
||||||
gLog.Println(LvDEBUG, "retry app ", peerNode)
|
gLog.Println(LvDEBUG, "retry app ", app.config.LogPeerNode())
|
||||||
app.retryRelayNum = 0
|
app.retryRelayNum = 0
|
||||||
app.nextRetryRelayTime = time.Now()
|
app.nextRetryRelayTime = time.Now()
|
||||||
app.hbMtx.Lock()
|
app.hbMtx.Lock()
|
||||||
@@ -171,7 +188,7 @@ func (c *Config) retryApp(peerNode string) {
|
|||||||
func (c *Config) retryAllApp() {
|
func (c *Config) retryAllApp() {
|
||||||
GNetwork.apps.Range(func(id, i interface{}) bool {
|
GNetwork.apps.Range(func(id, i interface{}) bool {
|
||||||
app := i.(*p2pApp)
|
app := i.(*p2pApp)
|
||||||
gLog.Println(LvDEBUG, "retry app ", app.config.PeerNode)
|
gLog.Println(LvDEBUG, "retry app ", app.config.LogPeerNode())
|
||||||
app.config.retryNum = 0
|
app.config.retryNum = 0
|
||||||
app.config.nextRetryTime = time.Now()
|
app.config.nextRetryTime = time.Now()
|
||||||
app.retryRelayNum = 0
|
app.retryRelayNum = 0
|
||||||
@@ -189,7 +206,7 @@ func (c *Config) retryAllMemApp() {
|
|||||||
if app.config.SrcPort != 0 {
|
if app.config.SrcPort != 0 {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
gLog.Println(LvDEBUG, "retry app ", app.config.PeerNode)
|
gLog.Println(LvDEBUG, "retry app ", app.config.LogPeerNode())
|
||||||
app.config.retryNum = 0
|
app.config.retryNum = 0
|
||||||
app.config.nextRetryTime = time.Now()
|
app.config.nextRetryTime = time.Now()
|
||||||
app.retryRelayNum = 0
|
app.retryRelayNum = 0
|
||||||
@@ -221,17 +238,8 @@ func (c *Config) delete(app AppConfig) {
|
|||||||
defer c.mtx.Unlock()
|
defer c.mtx.Unlock()
|
||||||
defer c.save()
|
defer c.save()
|
||||||
for i := 0; i < len(c.Apps); i++ {
|
for i := 0; i < len(c.Apps); i++ {
|
||||||
got := false
|
if (app.SrcPort != 0 && c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort) || // normal app
|
||||||
if app.SrcPort != 0 { // normal p2papp
|
(app.SrcPort == 0 && c.Apps[i].PeerNode == app.PeerNode) { // memapp
|
||||||
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
|
||||||
got = true
|
|
||||||
}
|
|
||||||
} else { // memapp
|
|
||||||
if c.Apps[i].PeerNode == app.PeerNode {
|
|
||||||
got = true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if got {
|
|
||||||
if i == len(c.Apps)-1 {
|
if i == len(c.Apps)-1 {
|
||||||
c.Apps = c.Apps[:i]
|
c.Apps = c.Apps[:i]
|
||||||
} else {
|
} else {
|
||||||
@@ -240,12 +248,14 @@ func (c *Config) delete(app AppConfig) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) save() {
|
func (c *Config) save() {
|
||||||
// c.mtx.Lock()
|
// c.mtx.Lock()
|
||||||
// defer c.mtx.Unlock() // internal call
|
// defer c.mtx.Unlock() // internal call
|
||||||
|
if c.Network.Token == 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
data, _ := json.MarshalIndent(c, "", " ")
|
data, _ := json.MarshalIndent(c, "", " ")
|
||||||
err := os.WriteFile("config.json", data, 0644)
|
err := os.WriteFile("config.json", data, 0644)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -256,6 +266,9 @@ func (c *Config) save() {
|
|||||||
func (c *Config) saveCache() {
|
func (c *Config) saveCache() {
|
||||||
// c.mtx.Lock()
|
// c.mtx.Lock()
|
||||||
// defer c.mtx.Unlock() // internal call
|
// defer c.mtx.Unlock() // internal call
|
||||||
|
if c.Network.Token == 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
data, _ := json.MarshalIndent(c, "", " ")
|
data, _ := json.MarshalIndent(c, "", " ")
|
||||||
err := os.WriteFile("config.json0", data, 0644)
|
err := os.WriteFile("config.json0", data, 0644)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -265,6 +278,7 @@ func (c *Config) saveCache() {
|
|||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
gConf.LogLevel = int(LvINFO)
|
gConf.LogLevel = int(LvINFO)
|
||||||
|
gConf.MaxLogSize = 1024 * 1024
|
||||||
gConf.Network.ShareBandwidth = 10
|
gConf.Network.ShareBandwidth = 10
|
||||||
gConf.Network.ServerHost = "api.openp2p.cn"
|
gConf.Network.ServerHost = "api.openp2p.cn"
|
||||||
gConf.Network.ServerPort = WsPort
|
gConf.Network.ServerPort = WsPort
|
||||||
@@ -448,6 +462,9 @@ func parseParams(subCommand string, cmd string) {
|
|||||||
if f.Name == "loglevel" {
|
if f.Name == "loglevel" {
|
||||||
gConf.LogLevel = *logLevel
|
gConf.LogLevel = *logLevel
|
||||||
}
|
}
|
||||||
|
if f.Name == "maxlogsize" {
|
||||||
|
gConf.MaxLogSize = *maxLogSize
|
||||||
|
}
|
||||||
if f.Name == "tcpport" {
|
if f.Name == "tcpport" {
|
||||||
gConf.Network.TCPPort = *tcpPort
|
gConf.Network.TCPPort = *tcpPort
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -29,4 +29,5 @@ var (
|
|||||||
ErrPeerConnectRelay = errors.New("peer connect relayNode error")
|
ErrPeerConnectRelay = errors.New("peer connect relayNode error")
|
||||||
ErrBuildTunnelBusy = errors.New("build tunnel busy")
|
ErrBuildTunnelBusy = errors.New("build tunnel busy")
|
||||||
ErrMemAppTunnelNotFound = errors.New("memapp tunnel not found")
|
ErrMemAppTunnelNotFound = errors.New("memapp tunnel not found")
|
||||||
|
ErrRemoteServiceUnable = errors.New("remote service unable")
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ import (
|
|||||||
"encoding/binary"
|
"encoding/binary"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"net"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"reflect"
|
"reflect"
|
||||||
@@ -44,6 +45,7 @@ func handlePush(subType uint16, msg []byte) error {
|
|||||||
config := AppConfig{}
|
config := AppConfig{}
|
||||||
config.PeerNode = req.RelayName
|
config.PeerNode = req.RelayName
|
||||||
config.peerToken = req.RelayToken
|
config.peerToken = req.RelayToken
|
||||||
|
config.relayMode = req.RelayMode
|
||||||
go func(r AddRelayTunnelReq) {
|
go func(r AddRelayTunnelReq) {
|
||||||
t, errDt := GNetwork.addDirectTunnel(config, 0)
|
t, errDt := GNetwork.addDirectTunnel(config, 0)
|
||||||
if errDt == nil {
|
if errDt == nil {
|
||||||
@@ -142,6 +144,8 @@ func handlePush(subType uint16, msg []byte) error {
|
|||||||
err = handleLog(msg)
|
err = handleLog(msg)
|
||||||
case MsgPushReportGoroutine:
|
case MsgPushReportGoroutine:
|
||||||
err = handleReportGoroutine()
|
err = handleReportGoroutine()
|
||||||
|
case MsgPushCheckRemoteService:
|
||||||
|
err = handleCheckRemoteService(msg)
|
||||||
case MsgPushEditApp:
|
case MsgPushEditApp:
|
||||||
err = handleEditApp(msg)
|
err = handleEditApp(msg)
|
||||||
case MsgPushEditNode:
|
case MsgPushEditNode:
|
||||||
@@ -458,3 +462,21 @@ func handleReportGoroutine() (err error) {
|
|||||||
stackLen := runtime.Stack(buf, true)
|
stackLen := runtime.Stack(buf, true)
|
||||||
return GNetwork.write(MsgReport, MsgPushReportLog, string(buf[:stackLen]))
|
return GNetwork.write(MsgReport, MsgPushReportLog, string(buf[:stackLen]))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func handleCheckRemoteService(msg []byte) (err error) {
|
||||||
|
gLog.Println(LvDEBUG, "handleCheckRemoteService")
|
||||||
|
req := CheckRemoteService{}
|
||||||
|
if err = json.Unmarshal(msg[openP2PHeaderSize:], &req); err != nil {
|
||||||
|
gLog.Printf(LvERROR, "wrong %v:%s %s", reflect.TypeOf(req), err, string(msg[openP2PHeaderSize:]))
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
rsp := PushRsp{Error: 0}
|
||||||
|
conn, err := net.DialTimeout("tcp", fmt.Sprintf("%s:%d", req.Host, req.Port), time.Second*3)
|
||||||
|
if err != nil {
|
||||||
|
rsp.Error = 1
|
||||||
|
rsp.Detail = ErrRemoteServiceUnable.Error()
|
||||||
|
} else {
|
||||||
|
conn.Close()
|
||||||
|
}
|
||||||
|
return GNetwork.write(MsgReport, MsgReportResponse, rsp)
|
||||||
|
}
|
||||||
|
|||||||
@@ -52,7 +52,7 @@ func addRoute(dst, gw, ifname string) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func delRoute(dst, gw string) error {
|
func delRoute(dst, gw string) error {
|
||||||
err := exec.Command("route", "delete", dst, gw).Run()
|
err := exec.Command("route", "delete", dst, "-gateway", gw).Run()
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
func delRoutesByGateway(gateway string) error {
|
func delRoutesByGateway(gateway string) error {
|
||||||
@@ -68,13 +68,14 @@ func delRoutesByGateway(gateway string) error {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
fields := strings.Fields(line)
|
fields := strings.Fields(line)
|
||||||
if len(fields) >= 7 && fields[0] == "default" && fields[len(fields)-1] == gateway {
|
if len(fields) >= 2 {
|
||||||
delCmd := exec.Command("route", "delete", "default", gateway)
|
cmd := exec.Command("route", "delete", fields[0], gateway)
|
||||||
err := delCmd.Run()
|
err := cmd.Run()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
gLog.Printf(LvERROR, "Delete route %s error:%s", fields[0], err)
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
fmt.Printf("Delete route ok: %s %s\n", "default", gateway)
|
gLog.Printf(LvINFO, "Delete route ok: %s %s\n", fields[0], gateway)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
+3
-2
@@ -124,9 +124,10 @@ func delRoutesByGateway(gateway string) error {
|
|||||||
delCmd := exec.Command("route", "del", "-net", fields[0], "gw", gateway)
|
delCmd := exec.Command("route", "del", "-net", fields[0], "gw", gateway)
|
||||||
err := delCmd.Run()
|
err := delCmd.Run()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
gLog.Printf(LvERROR, "Delete route %s error:%s", fields[0], err)
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
fmt.Printf("Delete route ok: %s %s %s\n", fields[0], fields[1], gateway)
|
gLog.Printf(LvINFO, "Delete route ok: %s %s %s\n", fields[0], fields[1], gateway)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -133,9 +133,10 @@ func delRoutesByGateway(gateway string) error {
|
|||||||
cmd := exec.Command("route", "delete", fields[0], "mask", fields[1], gateway)
|
cmd := exec.Command("route", "delete", fields[0], "mask", fields[1], gateway)
|
||||||
err := cmd.Run()
|
err := cmd.Run()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fmt.Println("Delete route error:", err)
|
gLog.Printf(LvERROR, "Delete route %s error:%s", fields[0], err)
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
fmt.Printf("Delete route ok: %s %s %s\n", fields[0], fields[1], gateway)
|
gLog.Printf(LvINFO, "Delete route ok: %s %s %s\n", fields[0], fields[1], gateway)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
+14
-14
@@ -138,7 +138,7 @@ func (app *p2pApp) checkDirectTunnel() error {
|
|||||||
app.config.retryNum = 1
|
app.config.retryNum = 1
|
||||||
}
|
}
|
||||||
if app.config.retryNum > 0 { // first time not show reconnect log
|
if app.config.retryNum > 0 { // first time not show reconnect log
|
||||||
gLog.Printf(LvINFO, "detect app %s appid:%d disconnect, reconnecting the %d times...", app.config.PeerNode, app.id, app.config.retryNum)
|
gLog.Printf(LvINFO, "detect app %s appid:%d disconnect, reconnecting the %d times...", app.config.LogPeerNode(), app.id, app.config.retryNum)
|
||||||
}
|
}
|
||||||
app.config.retryNum++
|
app.config.retryNum++
|
||||||
app.config.retryTime = time.Now()
|
app.config.retryTime = time.Now()
|
||||||
@@ -149,7 +149,7 @@ func (app *p2pApp) checkDirectTunnel() error {
|
|||||||
app.config.errMsg = err.Error()
|
app.config.errMsg = err.Error()
|
||||||
if err == ErrPeerOffline && app.config.retryNum > 2 { // stop retry, waiting for online
|
if err == ErrPeerOffline && app.config.retryNum > 2 { // stop retry, waiting for online
|
||||||
app.config.retryNum = retryLimit
|
app.config.retryNum = retryLimit
|
||||||
gLog.Printf(LvINFO, " %s offline, it will auto reconnect when peer node online", app.config.PeerNode)
|
gLog.Printf(LvINFO, " %s offline, it will auto reconnect when peer node online", app.config.LogPeerNode())
|
||||||
}
|
}
|
||||||
if err == ErrBuildTunnelBusy {
|
if err == ErrBuildTunnelBusy {
|
||||||
app.config.retryNum--
|
app.config.retryNum--
|
||||||
@@ -174,7 +174,7 @@ func (app *p2pApp) buildDirectTunnel() error {
|
|||||||
pn := GNetwork
|
pn := GNetwork
|
||||||
initErr := pn.requestPeerInfo(&app.config)
|
initErr := pn.requestPeerInfo(&app.config)
|
||||||
if initErr != nil {
|
if initErr != nil {
|
||||||
gLog.Printf(LvERROR, "%s init error:%s", app.config.PeerNode, initErr)
|
gLog.Printf(LvERROR, "%s init error:%s", app.config.LogPeerNode(), initErr)
|
||||||
return initErr
|
return initErr
|
||||||
}
|
}
|
||||||
t, err = pn.addDirectTunnel(app.config, 0)
|
t, err = pn.addDirectTunnel(app.config, 0)
|
||||||
@@ -212,7 +212,7 @@ func (app *p2pApp) buildDirectTunnel() error {
|
|||||||
AppID: app.id,
|
AppID: app.id,
|
||||||
AppKey: app.key,
|
AppKey: app.key,
|
||||||
}
|
}
|
||||||
gLog.Printf(LvDEBUG, "sync appkey direct to %s", app.config.PeerNode)
|
gLog.Printf(LvDEBUG, "sync appkey direct to %s", app.config.LogPeerNode())
|
||||||
pn.push(app.config.PeerNode, MsgPushAPPKey, &syncKeyReq)
|
pn.push(app.config.PeerNode, MsgPushAPPKey, &syncKeyReq)
|
||||||
app.setDirectTunnel(t)
|
app.setDirectTunnel(t)
|
||||||
|
|
||||||
@@ -220,7 +220,7 @@ func (app *p2pApp) buildDirectTunnel() error {
|
|||||||
if app.config.SrcPort == 0 {
|
if app.config.SrcPort == 0 {
|
||||||
req := ServerSideSaveMemApp{From: gConf.Network.Node, Node: gConf.Network.Node, TunnelID: t.id, RelayTunnelID: 0, AppID: app.id}
|
req := ServerSideSaveMemApp{From: gConf.Network.Node, Node: gConf.Network.Node, TunnelID: t.id, RelayTunnelID: 0, AppID: app.id}
|
||||||
pn.push(app.config.PeerNode, MsgPushServerSideSaveMemApp, &req)
|
pn.push(app.config.PeerNode, MsgPushServerSideSaveMemApp, &req)
|
||||||
gLog.Printf(LvDEBUG, "push %s ServerSideSaveMemApp: %s", app.config.PeerNode, prettyJson(req))
|
gLog.Printf(LvDEBUG, "push %s ServerSideSaveMemApp: %s", app.config.LogPeerNode(), prettyJson(req))
|
||||||
}
|
}
|
||||||
gLog.Printf(LvDEBUG, "%s use tunnel %d", app.config.AppName, t.id)
|
gLog.Printf(LvDEBUG, "%s use tunnel %d", app.config.AppName, t.id)
|
||||||
return nil
|
return nil
|
||||||
@@ -244,7 +244,7 @@ func (app *p2pApp) checkRelayTunnel() error {
|
|||||||
app.retryRelayNum = 1
|
app.retryRelayNum = 1
|
||||||
}
|
}
|
||||||
if app.retryRelayNum > 0 { // first time not show reconnect log
|
if app.retryRelayNum > 0 { // first time not show reconnect log
|
||||||
gLog.Printf(LvINFO, "detect app %s appid:%d relay disconnect, reconnecting the %d times...", app.config.PeerNode, app.id, app.retryRelayNum)
|
gLog.Printf(LvINFO, "detect app %s appid:%d relay disconnect, reconnecting the %d times...", app.config.LogPeerNode(), app.id, app.retryRelayNum)
|
||||||
}
|
}
|
||||||
app.setRelayTunnel(nil) // reset relayTunnel
|
app.setRelayTunnel(nil) // reset relayTunnel
|
||||||
app.retryRelayNum++
|
app.retryRelayNum++
|
||||||
@@ -256,7 +256,7 @@ func (app *p2pApp) checkRelayTunnel() error {
|
|||||||
app.errMsg = err.Error()
|
app.errMsg = err.Error()
|
||||||
if err == ErrPeerOffline && app.retryRelayNum > 2 { // stop retry, waiting for online
|
if err == ErrPeerOffline && app.retryRelayNum > 2 { // stop retry, waiting for online
|
||||||
app.retryRelayNum = retryLimit
|
app.retryRelayNum = retryLimit
|
||||||
gLog.Printf(LvINFO, " %s offline, it will auto reconnect when peer node online", app.config.PeerNode)
|
gLog.Printf(LvINFO, " %s offline, it will auto reconnect when peer node online", app.config.LogPeerNode())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if app.Tunnel() != nil {
|
if app.Tunnel() != nil {
|
||||||
@@ -282,7 +282,7 @@ func (app *p2pApp) buildRelayTunnel() error {
|
|||||||
config := app.config
|
config := app.config
|
||||||
initErr := pn.requestPeerInfo(&config)
|
initErr := pn.requestPeerInfo(&config)
|
||||||
if initErr != nil {
|
if initErr != nil {
|
||||||
gLog.Printf(LvERROR, "%s init error:%s", config.PeerNode, initErr)
|
gLog.Printf(LvERROR, "%s init error:%s", config.LogPeerNode(), initErr)
|
||||||
return initErr
|
return initErr
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -318,7 +318,7 @@ func (app *p2pApp) buildRelayTunnel() error {
|
|||||||
AppID: app.id,
|
AppID: app.id,
|
||||||
AppKey: app.key,
|
AppKey: app.key,
|
||||||
}
|
}
|
||||||
gLog.Printf(LvDEBUG, "sync appkey relay to %s", config.PeerNode)
|
gLog.Printf(LvDEBUG, "sync appkey relay to %s", config.LogPeerNode())
|
||||||
pn.push(config.PeerNode, MsgPushAPPKey, &syncKeyReq)
|
pn.push(config.PeerNode, MsgPushAPPKey, &syncKeyReq)
|
||||||
app.setRelayTunnelID(rtid)
|
app.setRelayTunnelID(rtid)
|
||||||
app.setRelayTunnel(t)
|
app.setRelayTunnel(t)
|
||||||
@@ -330,7 +330,7 @@ func (app *p2pApp) buildRelayTunnel() error {
|
|||||||
if config.SrcPort == 0 {
|
if config.SrcPort == 0 {
|
||||||
req := ServerSideSaveMemApp{From: gConf.Network.Node, Node: relayNode, TunnelID: rtid, RelayTunnelID: t.id, AppID: app.id, RelayMode: relayMode}
|
req := ServerSideSaveMemApp{From: gConf.Network.Node, Node: relayNode, TunnelID: rtid, RelayTunnelID: t.id, AppID: app.id, RelayMode: relayMode}
|
||||||
pn.push(config.PeerNode, MsgPushServerSideSaveMemApp, &req)
|
pn.push(config.PeerNode, MsgPushServerSideSaveMemApp, &req)
|
||||||
gLog.Printf(LvDEBUG, "push %s relay ServerSideSaveMemApp: %s", config.PeerNode, prettyJson(req))
|
gLog.Printf(LvDEBUG, "push %s relay ServerSideSaveMemApp: %s", config.LogPeerNode(), prettyJson(req))
|
||||||
}
|
}
|
||||||
gLog.Printf(LvDEBUG, "%s use tunnel %d", app.config.AppName, t.id)
|
gLog.Printf(LvDEBUG, "%s use tunnel %d", app.config.AppName, t.id)
|
||||||
return nil
|
return nil
|
||||||
@@ -594,8 +594,8 @@ func (app *p2pApp) close() {
|
|||||||
func (app *p2pApp) relayHeartbeatLoop() {
|
func (app *p2pApp) relayHeartbeatLoop() {
|
||||||
app.wg.Add(1)
|
app.wg.Add(1)
|
||||||
defer app.wg.Done()
|
defer app.wg.Done()
|
||||||
gLog.Printf(LvDEBUG, "%s appid:%d relayHeartbeat to rtid:%d start", app.config.PeerNode, app.id, app.rtid)
|
gLog.Printf(LvDEBUG, "%s appid:%d relayHeartbeat to rtid:%d start", app.config.LogPeerNode(), app.id, app.rtid)
|
||||||
defer gLog.Printf(LvDEBUG, "%s appid:%d relayHeartbeat to rtid%d end", app.config.PeerNode, app.id, app.rtid)
|
defer gLog.Printf(LvDEBUG, "%s appid:%d relayHeartbeat to rtid%d end", app.config.LogPeerNode(), app.id, app.rtid)
|
||||||
|
|
||||||
for app.running {
|
for app.running {
|
||||||
if app.RelayTunnel() == nil || !app.RelayTunnel().isRuning() {
|
if app.RelayTunnel() == nil || !app.RelayTunnel().isRuning() {
|
||||||
@@ -606,11 +606,11 @@ func (app *p2pApp) relayHeartbeatLoop() {
|
|||||||
AppID: app.id}
|
AppID: app.id}
|
||||||
err := app.RelayTunnel().WriteMessage(app.rtid, MsgP2P, MsgRelayHeartbeat, &req)
|
err := app.RelayTunnel().WriteMessage(app.rtid, MsgP2P, MsgRelayHeartbeat, &req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Printf(LvERROR, "%s appid:%d rtid:%d write relay tunnel heartbeat error %s", app.config.PeerNode, app.id, app.rtid, err)
|
gLog.Printf(LvERROR, "%s appid:%d rtid:%d write relay tunnel heartbeat error %s", app.config.LogPeerNode(), app.id, app.rtid, err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// TODO: debug relay heartbeat
|
// TODO: debug relay heartbeat
|
||||||
gLog.Printf(LvDEBUG, "%s appid:%d rtid:%d write relay tunnel heartbeat ok", app.config.PeerNode, app.id, app.rtid)
|
gLog.Printf(LvDEBUG, "%s appid:%d rtid:%d write relay tunnel heartbeat ok", app.config.LogPeerNode(), app.id, app.rtid)
|
||||||
time.Sleep(TunnelHeartbeatTime)
|
time.Sleep(TunnelHeartbeatTime)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+22
-23
@@ -115,6 +115,7 @@ func (pn *P2PNetwork) run() {
|
|||||||
pn.write(MsgHeartbeat, 0, "")
|
pn.write(MsgHeartbeat, 0, "")
|
||||||
case <-pn.restartCh:
|
case <-pn.restartCh:
|
||||||
gLog.Printf(LvDEBUG, "got restart channel")
|
gLog.Printf(LvDEBUG, "got restart channel")
|
||||||
|
GNetwork.sdwan.reset()
|
||||||
pn.online = false
|
pn.online = false
|
||||||
pn.wgReconnect.Wait() // wait read/autorunapp goroutine end
|
pn.wgReconnect.Wait() // wait read/autorunapp goroutine end
|
||||||
delay := ClientAPITimeout + time.Duration(rand.Int()%pn.loginMaxDelaySeconds)*time.Second
|
delay := ClientAPITimeout + time.Duration(rand.Int()%pn.loginMaxDelaySeconds)*time.Second
|
||||||
@@ -124,8 +125,9 @@ func (pn *P2PNetwork) run() {
|
|||||||
gLog.Println(LvERROR, "P2PNetwork init error:", err)
|
gLog.Println(LvERROR, "P2PNetwork init error:", err)
|
||||||
}
|
}
|
||||||
gConf.retryAllApp()
|
gConf.retryAllApp()
|
||||||
|
|
||||||
case t := <-pn.tunnelCloseCh:
|
case t := <-pn.tunnelCloseCh:
|
||||||
gLog.Printf(LvDEBUG, "got tunnelCloseCh %s", t.config.PeerNode)
|
gLog.Printf(LvDEBUG, "got tunnelCloseCh %s", t.config.LogPeerNode())
|
||||||
pn.apps.Range(func(id, i interface{}) bool {
|
pn.apps.Range(func(id, i interface{}) bool {
|
||||||
app := i.(*p2pApp)
|
app := i.(*p2pApp)
|
||||||
if app.DirectTunnel() == t {
|
if app.DirectTunnel() == t {
|
||||||
@@ -195,12 +197,12 @@ func (pn *P2PNetwork) autorunApp() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) addRelayTunnel(config AppConfig) (*P2PTunnel, uint64, string, error) {
|
func (pn *P2PNetwork) addRelayTunnel(config AppConfig) (*P2PTunnel, uint64, string, error) {
|
||||||
gLog.Printf(LvINFO, "addRelayTunnel to %s start", config.PeerNode)
|
gLog.Printf(LvINFO, "addRelayTunnel to %s start", config.LogPeerNode())
|
||||||
defer gLog.Printf(LvINFO, "addRelayTunnel to %s end", config.PeerNode)
|
defer gLog.Printf(LvINFO, "addRelayTunnel to %s end", config.LogPeerNode())
|
||||||
relayConfig := AppConfig{
|
relayConfig := AppConfig{
|
||||||
PeerNode: config.RelayNode,
|
PeerNode: config.RelayNode,
|
||||||
peerToken: config.peerToken}
|
peerToken: config.peerToken,
|
||||||
relayMode := "private"
|
relayMode: "private"}
|
||||||
if relayConfig.PeerNode == "" {
|
if relayConfig.PeerNode == "" {
|
||||||
// find existing relay tunnel
|
// find existing relay tunnel
|
||||||
pn.apps.Range(func(id, i interface{}) bool {
|
pn.apps.Range(func(id, i interface{}) bool {
|
||||||
@@ -212,7 +214,7 @@ func (pn *P2PNetwork) addRelayTunnel(config AppConfig) (*P2PTunnel, uint64, stri
|
|||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
relayConfig.PeerNode = app.RelayTunnel().config.PeerNode
|
relayConfig.PeerNode = app.RelayTunnel().config.PeerNode
|
||||||
gLog.Printf(LvDEBUG, "found existing relay tunnel %s", relayConfig.PeerNode)
|
gLog.Printf(LvDEBUG, "found existing relay tunnel %s", relayConfig.LogPeerNode())
|
||||||
return false
|
return false
|
||||||
})
|
})
|
||||||
if relayConfig.PeerNode == "" { // request relay node
|
if relayConfig.PeerNode == "" { // request relay node
|
||||||
@@ -231,11 +233,11 @@ func (pn *P2PNetwork) addRelayTunnel(config AppConfig) (*P2PTunnel, uint64, stri
|
|||||||
gLog.Printf(LvERROR, "MsgRelayNodeReq error")
|
gLog.Printf(LvERROR, "MsgRelayNodeReq error")
|
||||||
return nil, 0, "", errors.New("MsgRelayNodeReq error")
|
return nil, 0, "", errors.New("MsgRelayNodeReq error")
|
||||||
}
|
}
|
||||||
gLog.Printf(LvDEBUG, "got relay node:%s", rsp.RelayName)
|
gLog.Printf(LvDEBUG, "got relay node:%s", relayConfig.LogPeerNode())
|
||||||
|
|
||||||
relayConfig.PeerNode = rsp.RelayName
|
relayConfig.PeerNode = rsp.RelayName
|
||||||
relayConfig.peerToken = rsp.RelayToken
|
relayConfig.peerToken = rsp.RelayToken
|
||||||
relayMode = rsp.Mode
|
relayConfig.relayMode = rsp.Mode
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
@@ -250,10 +252,10 @@ func (pn *P2PNetwork) addRelayTunnel(config AppConfig) (*P2PTunnel, uint64, stri
|
|||||||
From: gConf.Network.Node,
|
From: gConf.Network.Node,
|
||||||
RelayName: relayConfig.PeerNode,
|
RelayName: relayConfig.PeerNode,
|
||||||
RelayToken: relayConfig.peerToken,
|
RelayToken: relayConfig.peerToken,
|
||||||
RelayMode: relayMode,
|
RelayMode: relayConfig.relayMode,
|
||||||
RelayTunnelID: t.id,
|
RelayTunnelID: t.id,
|
||||||
}
|
}
|
||||||
gLog.Printf(LvDEBUG, "push %s the relay node(%s)", config.PeerNode, relayConfig.PeerNode)
|
gLog.Printf(LvDEBUG, "push %s the relay node(%s)", config.LogPeerNode(), relayConfig.LogPeerNode())
|
||||||
pn.push(config.PeerNode, MsgPushAddRelayTunnelReq, &req)
|
pn.push(config.PeerNode, MsgPushAddRelayTunnelReq, &req)
|
||||||
|
|
||||||
// wait relay ready
|
// wait relay ready
|
||||||
@@ -267,13 +269,13 @@ func (pn *P2PNetwork) addRelayTunnel(config AppConfig) (*P2PTunnel, uint64, stri
|
|||||||
gLog.Println(LvDEBUG, ErrPeerConnectRelay)
|
gLog.Println(LvDEBUG, ErrPeerConnectRelay)
|
||||||
return nil, 0, "", ErrPeerConnectRelay
|
return nil, 0, "", ErrPeerConnectRelay
|
||||||
}
|
}
|
||||||
return t, rspID.ID, relayMode, err
|
return t, rspID.ID, relayConfig.relayMode, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// use *AppConfig to save status
|
// use *AppConfig to save status
|
||||||
func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
||||||
gLog.Printf(LvINFO, "addApp %s to %s:%s:%d start", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
gLog.Printf(LvINFO, "addApp %s to %s:%s:%d start", config.AppName, config.LogPeerNode(), config.DstHost, config.DstPort)
|
||||||
defer gLog.Printf(LvINFO, "addApp %s to %s:%s:%d end", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
defer gLog.Printf(LvINFO, "addApp %s to %s:%s:%d end", config.AppName, config.LogPeerNode(), config.DstHost, config.DstPort)
|
||||||
if !pn.online {
|
if !pn.online {
|
||||||
return errors.New("P2PNetwork offline")
|
return errors.New("P2PNetwork offline")
|
||||||
}
|
}
|
||||||
@@ -304,8 +306,8 @@ func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) DeleteApp(config AppConfig) {
|
func (pn *P2PNetwork) DeleteApp(config AppConfig) {
|
||||||
gLog.Printf(LvINFO, "DeleteApp %s to %s:%s:%d start", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
gLog.Printf(LvINFO, "DeleteApp %s to %s:%s:%d start", config.AppName, config.LogPeerNode(), config.DstHost, config.DstPort)
|
||||||
defer gLog.Printf(LvINFO, "DeleteApp %s to %s:%s:%d end", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
defer gLog.Printf(LvINFO, "DeleteApp %s to %s:%s:%d end", config.AppName, config.LogPeerNode(), config.DstHost, config.DstPort)
|
||||||
// close the apps of this config
|
// close the apps of this config
|
||||||
i, ok := pn.apps.Load(config.ID())
|
i, ok := pn.apps.Load(config.ID())
|
||||||
if ok {
|
if ok {
|
||||||
@@ -339,8 +341,8 @@ func (pn *P2PNetwork) findTunnel(peerNode string) (t *P2PTunnel) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) addDirectTunnel(config AppConfig, tid uint64) (t *P2PTunnel, err error) {
|
func (pn *P2PNetwork) addDirectTunnel(config AppConfig, tid uint64) (t *P2PTunnel, err error) {
|
||||||
gLog.Printf(LvDEBUG, "addDirectTunnel %s%d to %s:%s:%d tid:%d start", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort, tid)
|
gLog.Printf(LvDEBUG, "addDirectTunnel %s%d to %s:%s:%d tid:%d start", config.Protocol, config.SrcPort, config.LogPeerNode(), config.DstHost, config.DstPort, tid)
|
||||||
defer gLog.Printf(LvDEBUG, "addDirectTunnel %s%d to %s:%s:%d tid:%d end", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort, tid)
|
defer gLog.Printf(LvDEBUG, "addDirectTunnel %s%d to %s:%s:%d tid:%d end", config.Protocol, config.SrcPort, config.LogPeerNode(), config.DstHost, config.DstPort, tid)
|
||||||
isClient := false
|
isClient := false
|
||||||
// client side tid=0, assign random uint64
|
// client side tid=0, assign random uint64
|
||||||
if tid == 0 {
|
if tid == 0 {
|
||||||
@@ -360,12 +362,12 @@ func (pn *P2PNetwork) addDirectTunnel(config AppConfig, tid uint64) (t *P2PTunne
|
|||||||
// peer info
|
// peer info
|
||||||
initErr := pn.requestPeerInfo(&config)
|
initErr := pn.requestPeerInfo(&config)
|
||||||
if initErr != nil {
|
if initErr != nil {
|
||||||
gLog.Printf(LvERROR, "%s init error:%s", config.PeerNode, initErr)
|
gLog.Printf(LvERROR, "%s init error:%s", config.LogPeerNode(), initErr)
|
||||||
|
|
||||||
return nil, initErr
|
return nil, initErr
|
||||||
}
|
}
|
||||||
gLog.Printf(LvDEBUG, "config.peerNode=%s,config.peerVersion=%s,config.peerIP=%s,config.peerLanIP=%s,gConf.Network.publicIP=%s,config.peerIPv6=%s,config.hasIPv4=%d,config.hasUPNPorNATPMP=%d,gConf.Network.hasIPv4=%d,gConf.Network.hasUPNPorNATPMP=%d,config.peerNatType=%d,gConf.Network.natType=%d,",
|
gLog.Printf(LvDEBUG, "config.peerNode=%s,config.peerVersion=%s,config.peerIP=%s,config.peerLanIP=%s,gConf.Network.publicIP=%s,config.peerIPv6=%s,config.hasIPv4=%d,config.hasUPNPorNATPMP=%d,gConf.Network.hasIPv4=%d,gConf.Network.hasUPNPorNATPMP=%d,config.peerNatType=%d,gConf.Network.natType=%d,",
|
||||||
config.PeerNode, config.peerVersion, config.peerIP, config.peerLanIP, gConf.Network.publicIP, config.peerIPv6, config.hasIPv4, config.hasUPNPorNATPMP, gConf.Network.hasIPv4, gConf.Network.hasUPNPorNATPMP, config.peerNatType, gConf.Network.natType)
|
config.LogPeerNode(), config.peerVersion, config.peerIP, config.peerLanIP, gConf.Network.publicIP, config.peerIPv6, config.hasIPv4, config.hasUPNPorNATPMP, gConf.Network.hasIPv4, gConf.Network.hasUPNPorNATPMP, config.peerNatType, gConf.Network.natType)
|
||||||
// try Intranet
|
// try Intranet
|
||||||
if config.peerIP == gConf.Network.publicIP && compareVersion(config.peerVersion, SupportIntranetVersion) >= 0 { // old version client has no peerLanIP
|
if config.peerIP == gConf.Network.publicIP && compareVersion(config.peerVersion, SupportIntranetVersion) >= 0 { // old version client has no peerLanIP
|
||||||
gLog.Println(LvINFO, "try Intranet")
|
gLog.Println(LvINFO, "try Intranet")
|
||||||
@@ -533,7 +535,6 @@ func (pn *P2PNetwork) init() error {
|
|||||||
caCertPool, errCert := x509.SystemCertPool()
|
caCertPool, errCert := x509.SystemCertPool()
|
||||||
if errCert != nil {
|
if errCert != nil {
|
||||||
gLog.Println(LvERROR, "Failed to load system root CAs:", errCert)
|
gLog.Println(LvERROR, "Failed to load system root CAs:", errCert)
|
||||||
} else {
|
|
||||||
caCertPool = x509.NewCertPool()
|
caCertPool = x509.NewCertPool()
|
||||||
}
|
}
|
||||||
caCertPool.AppendCertsFromPEM([]byte(rootCA))
|
caCertPool.AppendCertsFromPEM([]byte(rootCA))
|
||||||
@@ -624,8 +625,6 @@ func (pn *P2PNetwork) handleMessage(msg []byte) {
|
|||||||
gLog.Printf(LvERROR, "login error:%d, detail:%s", rsp.Error, rsp.Detail)
|
gLog.Printf(LvERROR, "login error:%d, detail:%s", rsp.Error, rsp.Detail)
|
||||||
pn.running = false
|
pn.running = false
|
||||||
} else {
|
} else {
|
||||||
gConf.Network.Token = rsp.Token
|
|
||||||
gConf.Network.User = rsp.User
|
|
||||||
gConf.setToken(rsp.Token)
|
gConf.setToken(rsp.Token)
|
||||||
gConf.setUser(rsp.User)
|
gConf.setUser(rsp.User)
|
||||||
if len(rsp.Node) >= MinNodeNameLen {
|
if len(rsp.Node) >= MinNodeNameLen {
|
||||||
@@ -726,7 +725,7 @@ func (pn *P2PNetwork) relay(to uint64, body []byte) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) push(to string, subType uint16, packet interface{}) error {
|
func (pn *P2PNetwork) push(to string, subType uint16, packet interface{}) error {
|
||||||
gLog.Printf(LvDEBUG, "push msgType %d to %s", subType, to)
|
// gLog.Printf(LvDEBUG, "push msgType %d to %s", subType, to)
|
||||||
if !pn.online {
|
if !pn.online {
|
||||||
return errors.New("client offline")
|
return errors.New("client offline")
|
||||||
}
|
}
|
||||||
|
|||||||
+23
-23
@@ -63,7 +63,7 @@ func (t *P2PTunnel) initPort() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) connect() error {
|
func (t *P2PTunnel) connect() error {
|
||||||
gLog.Printf(LvDEBUG, "start p2pTunnel to %s ", t.config.PeerNode)
|
gLog.Printf(LvDEBUG, "start p2pTunnel to %s ", t.config.LogPeerNode())
|
||||||
t.tunnelServer = false
|
t.tunnelServer = false
|
||||||
appKey := uint64(0)
|
appKey := uint64(0)
|
||||||
req := PushConnectReq{
|
req := PushConnectReq{
|
||||||
@@ -170,7 +170,7 @@ func (t *P2PTunnel) close() {
|
|||||||
t.conn.Close()
|
t.conn.Close()
|
||||||
}
|
}
|
||||||
t.pn.allTunnels.Delete(t.id)
|
t.pn.allTunnels.Delete(t.id)
|
||||||
gLog.Printf(LvINFO, "%d p2ptunnel close %s ", t.id, t.config.PeerNode)
|
gLog.Printf(LvINFO, "%d p2ptunnel close %s ", t.id, t.config.LogPeerNode())
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) start() error {
|
func (t *P2PTunnel) start() error {
|
||||||
@@ -205,7 +205,7 @@ func (t *P2PTunnel) handshake() error {
|
|||||||
gLog.Printf(LvDEBUG, "sleep %d ms", ts/time.Millisecond)
|
gLog.Printf(LvDEBUG, "sleep %d ms", ts/time.Millisecond)
|
||||||
time.Sleep(ts)
|
time.Sleep(ts)
|
||||||
}
|
}
|
||||||
gLog.Println(LvDEBUG, "handshake to ", t.config.PeerNode)
|
gLog.Println(LvDEBUG, "handshake to ", t.config.LogPeerNode())
|
||||||
var err error
|
var err error
|
||||||
if gConf.Network.natType == NATCone && t.config.peerNatType == NATCone {
|
if gConf.Network.natType == NATCone && t.config.peerNatType == NATCone {
|
||||||
err = handshakeC2C(t)
|
err = handshakeC2C(t)
|
||||||
@@ -223,7 +223,7 @@ func (t *P2PTunnel) handshake() error {
|
|||||||
gLog.Println(LvERROR, "punch handshake error:", err)
|
gLog.Println(LvERROR, "punch handshake error:", err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
gLog.Printf(LvDEBUG, "handshake to %s ok", t.config.PeerNode)
|
gLog.Printf(LvDEBUG, "handshake to %s ok", t.config.LogPeerNode())
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -258,8 +258,8 @@ func (t *P2PTunnel) connectUnderlay() (err error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) connectUnderlayUDP() (c underlay, err error) {
|
func (t *P2PTunnel) connectUnderlayUDP() (c underlay, err error) {
|
||||||
gLog.Printf(LvDEBUG, "connectUnderlayUDP %s start ", t.config.PeerNode)
|
gLog.Printf(LvDEBUG, "connectUnderlayUDP %s start ", t.config.LogPeerNode())
|
||||||
defer gLog.Printf(LvDEBUG, "connectUnderlayUDP %s end ", t.config.PeerNode)
|
defer gLog.Printf(LvDEBUG, "connectUnderlayUDP %s end ", t.config.LogPeerNode())
|
||||||
var ul underlay
|
var ul underlay
|
||||||
underlayProtocol := t.config.UnderlayProtocol
|
underlayProtocol := t.config.UnderlayProtocol
|
||||||
if underlayProtocol == "" {
|
if underlayProtocol == "" {
|
||||||
@@ -330,8 +330,8 @@ func (t *P2PTunnel) connectUnderlayUDP() (c underlay, err error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) connectUnderlayTCP() (c underlay, err error) {
|
func (t *P2PTunnel) connectUnderlayTCP() (c underlay, err error) {
|
||||||
gLog.Printf(LvDEBUG, "connectUnderlayTCP %s start ", t.config.PeerNode)
|
gLog.Printf(LvDEBUG, "connectUnderlayTCP %s start ", t.config.LogPeerNode())
|
||||||
defer gLog.Printf(LvDEBUG, "connectUnderlayTCP %s end ", t.config.PeerNode)
|
defer gLog.Printf(LvDEBUG, "connectUnderlayTCP %s end ", t.config.LogPeerNode())
|
||||||
var ul *underlayTCP
|
var ul *underlayTCP
|
||||||
peerIP := t.config.peerIP
|
peerIP := t.config.peerIP
|
||||||
if t.config.linkMode == LinkModeIntranet {
|
if t.config.linkMode == LinkModeIntranet {
|
||||||
@@ -392,8 +392,8 @@ func (t *P2PTunnel) connectUnderlayTCP() (c underlay, err error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) connectUnderlayTCPSymmetric() (c underlay, err error) {
|
func (t *P2PTunnel) connectUnderlayTCPSymmetric() (c underlay, err error) {
|
||||||
gLog.Printf(LvDEBUG, "connectUnderlayTCPSymmetric %s start ", t.config.PeerNode)
|
gLog.Printf(LvDEBUG, "connectUnderlayTCPSymmetric %s start ", t.config.LogPeerNode())
|
||||||
defer gLog.Printf(LvDEBUG, "connectUnderlayTCPSymmetric %s end ", t.config.PeerNode)
|
defer gLog.Printf(LvDEBUG, "connectUnderlayTCPSymmetric %s end ", t.config.LogPeerNode())
|
||||||
ts := time.Duration(int64(t.punchTs) + t.pn.dt + t.pn.ddtma*int64(time.Since(t.pn.hbTime)+PunchTsDelay)/int64(NetworkHeartbeatTime) - time.Now().UnixNano())
|
ts := time.Duration(int64(t.punchTs) + t.pn.dt + t.pn.ddtma*int64(time.Since(t.pn.hbTime)+PunchTsDelay)/int64(NetworkHeartbeatTime) - time.Now().UnixNano())
|
||||||
if ts > PunchTsDelay || ts < 0 {
|
if ts > PunchTsDelay || ts < 0 {
|
||||||
ts = PunchTsDelay
|
ts = PunchTsDelay
|
||||||
@@ -426,7 +426,7 @@ func (t *P2PTunnel) connectUnderlayTCPSymmetric() (c underlay, err error) {
|
|||||||
}
|
}
|
||||||
_, buff, err := ul.ReadBuffer()
|
_, buff, err := ul.ReadBuffer()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Printf(LvERROR, "utcp.ReadBuffer error:", err)
|
gLog.Println(LvDEBUG, "c2s ul.ReadBuffer error:", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
req := P2PHandshakeReq{}
|
req := P2PHandshakeReq{}
|
||||||
@@ -455,7 +455,7 @@ func (t *P2PTunnel) connectUnderlayTCPSymmetric() (c underlay, err error) {
|
|||||||
|
|
||||||
_, buff, err := ul.ReadBuffer()
|
_, buff, err := ul.ReadBuffer()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Printf(LvERROR, "utcp.ReadBuffer error:", err)
|
gLog.Println(LvDEBUG, "s2c ul.ReadBuffer error:", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
req := P2PHandshakeReq{}
|
req := P2PHandshakeReq{}
|
||||||
@@ -486,8 +486,8 @@ func (t *P2PTunnel) connectUnderlayTCPSymmetric() (c underlay, err error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) connectUnderlayTCP6() (c underlay, err error) {
|
func (t *P2PTunnel) connectUnderlayTCP6() (c underlay, err error) {
|
||||||
gLog.Printf(LvDEBUG, "connectUnderlayTCP6 %s start ", t.config.PeerNode)
|
gLog.Printf(LvDEBUG, "connectUnderlayTCP6 %s start ", t.config.LogPeerNode())
|
||||||
defer gLog.Printf(LvDEBUG, "connectUnderlayTCP6 %s end ", t.config.PeerNode)
|
defer gLog.Printf(LvDEBUG, "connectUnderlayTCP6 %s end ", t.config.LogPeerNode())
|
||||||
var ul *underlayTCP6
|
var ul *underlayTCP6
|
||||||
if t.config.isUnderlayServer == 1 {
|
if t.config.isUnderlayServer == 1 {
|
||||||
t.pn.push(t.config.PeerNode, MsgPushUnderlayConnect, nil)
|
t.pn.push(t.config.PeerNode, MsgPushUnderlayConnect, nil)
|
||||||
@@ -512,14 +512,14 @@ func (t *P2PTunnel) connectUnderlayTCP6() (c underlay, err error) {
|
|||||||
t.pn.read(t.config.PeerNode, MsgPush, MsgPushUnderlayConnect, ReadMsgTimeout)
|
t.pn.read(t.config.PeerNode, MsgPush, MsgPushUnderlayConnect, ReadMsgTimeout)
|
||||||
gLog.Println(LvDEBUG, "TCP6 dial to ", t.config.peerIPv6)
|
gLog.Println(LvDEBUG, "TCP6 dial to ", t.config.peerIPv6)
|
||||||
ul, err = dialTCP6(t.config.peerIPv6, t.config.peerConeNatPort)
|
ul, err = dialTCP6(t.config.peerIPv6, t.config.peerConeNatPort)
|
||||||
if err != nil {
|
if err != nil || ul == nil {
|
||||||
return nil, fmt.Errorf("TCP6 dial to %s:%d error:%s", t.config.peerIPv6, t.config.peerConeNatPort, err)
|
return nil, fmt.Errorf("TCP6 dial to %s:%d error:%s", t.config.peerIPv6, t.config.peerConeNatPort, err)
|
||||||
}
|
}
|
||||||
handshakeBegin := time.Now()
|
handshakeBegin := time.Now()
|
||||||
ul.WriteBytes(MsgP2P, MsgTunnelHandshake, []byte("OpenP2P,hello"))
|
ul.WriteBytes(MsgP2P, MsgTunnelHandshake, []byte("OpenP2P,hello"))
|
||||||
_, buff, err := ul.ReadBuffer()
|
_, buff, errR := ul.ReadBuffer()
|
||||||
if err != nil {
|
if errR != nil {
|
||||||
return nil, fmt.Errorf("read MsgTunnelHandshake error:%s", err)
|
return nil, fmt.Errorf("read MsgTunnelHandshake error:%s", errR)
|
||||||
}
|
}
|
||||||
if buff != nil {
|
if buff != nil {
|
||||||
gLog.Println(LvDEBUG, string(buff))
|
gLog.Println(LvDEBUG, string(buff))
|
||||||
@@ -597,7 +597,7 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
tunnelID := binary.LittleEndian.Uint64(body[:8])
|
tunnelID := binary.LittleEndian.Uint64(body[:8])
|
||||||
gLog.Printf(LvDev, "relay data to %d, len=%d", tunnelID, head.DataLen-RelayHeaderSize)
|
gLog.Printf(LvDev, "relay data to %d, len=%d", tunnelID, head.DataLen-RelayHeaderSize)
|
||||||
if err := t.pn.relay(tunnelID, body[RelayHeaderSize:]); err != nil {
|
if err := t.pn.relay(tunnelID, body[RelayHeaderSize:]); err != nil {
|
||||||
gLog.Printf(LvERROR, "%s:%d relay to %d len=%d error:%s", t.config.PeerNode, t.id, tunnelID, len(body), ErrRelayTunnelNotFound)
|
gLog.Printf(LvERROR, "%s:%d relay to %d len=%d error:%s", t.config.LogPeerNode(), t.id, tunnelID, len(body), ErrRelayTunnelNotFound)
|
||||||
}
|
}
|
||||||
case MsgRelayHeartbeat:
|
case MsgRelayHeartbeat:
|
||||||
req := RelayHeartbeat{}
|
req := RelayHeartbeat{}
|
||||||
@@ -691,8 +691,8 @@ func (t *P2PTunnel) writeLoop() {
|
|||||||
t.hbMtx.Unlock()
|
t.hbMtx.Unlock()
|
||||||
tc := time.NewTicker(TunnelHeartbeatTime)
|
tc := time.NewTicker(TunnelHeartbeatTime)
|
||||||
defer tc.Stop()
|
defer tc.Stop()
|
||||||
gLog.Printf(LvDEBUG, "%s:%d tunnel writeLoop start", t.config.PeerNode, t.id)
|
gLog.Printf(LvDEBUG, "%s:%d tunnel writeLoop start", t.config.LogPeerNode(), t.id)
|
||||||
defer gLog.Printf(LvDEBUG, "%s:%d tunnel writeLoop end", t.config.PeerNode, t.id)
|
defer gLog.Printf(LvDEBUG, "%s:%d tunnel writeLoop end", t.config.LogPeerNode(), t.id)
|
||||||
for t.isRuning() {
|
for t.isRuning() {
|
||||||
select {
|
select {
|
||||||
case buff := <-t.writeDataSmall:
|
case buff := <-t.writeDataSmall:
|
||||||
@@ -781,13 +781,13 @@ func (t *P2PTunnel) asyncWriteNodeData(mainType, subType uint16, data []byte) {
|
|||||||
case t.writeDataSmall <- writeBytes:
|
case t.writeDataSmall <- writeBytes:
|
||||||
// gLog.Printf(LvWARN, "%s:%d t.writeDataSmall write %d", t.config.PeerNode, t.id, len(t.writeDataSmall))
|
// gLog.Printf(LvWARN, "%s:%d t.writeDataSmall write %d", t.config.PeerNode, t.id, len(t.writeDataSmall))
|
||||||
default:
|
default:
|
||||||
gLog.Printf(LvWARN, "%s:%d t.writeDataSmall is full, drop it", t.config.PeerNode, t.id)
|
gLog.Printf(LvWARN, "%s:%d t.writeDataSmall is full, drop it", t.config.LogPeerNode(), t.id)
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
select {
|
select {
|
||||||
case t.writeData <- writeBytes:
|
case t.writeData <- writeBytes:
|
||||||
default:
|
default:
|
||||||
gLog.Printf(LvWARN, "%s:%d t.writeData is full, drop it", t.config.PeerNode, t.id)
|
gLog.Printf(LvWARN, "%s:%d t.writeData is full, drop it", t.config.LogPeerNode(), t.id)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,87 @@
|
|||||||
|
package openp2p
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"net"
|
||||||
|
"os"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"golang.org/x/net/icmp"
|
||||||
|
"golang.org/x/net/ipv4"
|
||||||
|
)
|
||||||
|
|
||||||
|
// 定义ICMP回显请求和应答的结构
|
||||||
|
type ICMPMessage struct {
|
||||||
|
Type uint8
|
||||||
|
Code uint8
|
||||||
|
Checksum uint16
|
||||||
|
Ident uint16
|
||||||
|
Seq uint16
|
||||||
|
Data []byte
|
||||||
|
}
|
||||||
|
|
||||||
|
// Ping sends an ICMP Echo request to the specified host and returns the response time.
|
||||||
|
func Ping(host string) (time.Duration, error) {
|
||||||
|
// Resolve the IP address of the host
|
||||||
|
ipAddr, err := net.ResolveIPAddr("ip4", host)
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("failed to resolve host: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Create an ICMP listener
|
||||||
|
conn, err := net.ListenPacket("ip4:icmp", "0.0.0.0")
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("failed to create ICMP connection: %v", err)
|
||||||
|
}
|
||||||
|
defer conn.Close()
|
||||||
|
|
||||||
|
// Create an ICMP Echo request message
|
||||||
|
message := icmp.Message{
|
||||||
|
Type: ipv4.ICMPTypeEcho,
|
||||||
|
Code: 0,
|
||||||
|
Body: &icmp.Echo{
|
||||||
|
ID: os.Getpid() & 0xffff,
|
||||||
|
Seq: 1,
|
||||||
|
Data: []byte("HELLO-R-U-THERE"),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
// Marshal the message into binary form
|
||||||
|
messageBytes, err := message.Marshal(nil)
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("failed to marshal ICMP message: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Send the ICMP Echo request
|
||||||
|
start := time.Now()
|
||||||
|
if _, err := conn.WriteTo(messageBytes, ipAddr); err != nil {
|
||||||
|
return 0, fmt.Errorf("failed to send ICMP request: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Set a deadline for the response
|
||||||
|
err = conn.SetReadDeadline(time.Now().Add(3 * time.Second))
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("failed to set read deadline: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Read the ICMP response
|
||||||
|
response := make([]byte, 1500)
|
||||||
|
n, _, err := conn.ReadFrom(response)
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("failed to read ICMP response: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Parse the ICMP response message
|
||||||
|
parsedMessage, err := icmp.ParseMessage(ipv4.ICMPTypeEchoReply.Protocol(), response[:n])
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("failed to parse ICMP response: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check if the response is an Echo reply
|
||||||
|
if parsedMessage.Type == ipv4.ICMPTypeEchoReply {
|
||||||
|
duration := time.Since(start)
|
||||||
|
return duration, nil
|
||||||
|
} else {
|
||||||
|
return 0, fmt.Errorf("unexpected ICMP message: %+v", parsedMessage)
|
||||||
|
}
|
||||||
|
}
|
||||||
+9
-2
@@ -10,7 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
const OpenP2PVersion = "3.18.4"
|
const OpenP2PVersion = "3.21.8"
|
||||||
const ProductName string = "openp2p"
|
const ProductName string = "openp2p"
|
||||||
const LeastSupportVersion = "3.0.0"
|
const LeastSupportVersion = "3.0.0"
|
||||||
const SyncServerTimeVersion = "3.9.0"
|
const SyncServerTimeVersion = "3.9.0"
|
||||||
@@ -108,6 +108,7 @@ const (
|
|||||||
MsgPushReportGoroutine = 16
|
MsgPushReportGoroutine = 16
|
||||||
MsgPushReportMemApps = 17
|
MsgPushReportMemApps = 17
|
||||||
MsgPushServerSideSaveMemApp = 18
|
MsgPushServerSideSaveMemApp = 18
|
||||||
|
MsgPushCheckRemoteService = 19
|
||||||
)
|
)
|
||||||
|
|
||||||
// MsgP2P sub type message
|
// MsgP2P sub type message
|
||||||
@@ -143,6 +144,7 @@ const (
|
|||||||
MsgReportApps
|
MsgReportApps
|
||||||
MsgReportLog
|
MsgReportLog
|
||||||
MsgReportMemApps
|
MsgReportMemApps
|
||||||
|
MsgReportResponse
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -493,7 +495,7 @@ type SDWANInfo struct {
|
|||||||
ForceRelay int32 `json:"forceRelay,omitempty"`
|
ForceRelay int32 `json:"forceRelay,omitempty"`
|
||||||
PunchPriority int32 `json:"punchPriority,omitempty"`
|
PunchPriority int32 `json:"punchPriority,omitempty"`
|
||||||
Enable int32 `json:"enable,omitempty"`
|
Enable int32 `json:"enable,omitempty"`
|
||||||
Nodes []SDWANNode
|
Nodes []*SDWANNode
|
||||||
}
|
}
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -510,6 +512,11 @@ type ServerSideSaveMemApp struct {
|
|||||||
AppID uint64 `json:"appID,omitempty"`
|
AppID uint64 `json:"appID,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type CheckRemoteService struct {
|
||||||
|
Host string `json:"host,omitempty"`
|
||||||
|
Port uint32 `json:"port,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
const rootCA = `-----BEGIN CERTIFICATE-----
|
const rootCA = `-----BEGIN CERTIFICATE-----
|
||||||
MIIDhTCCAm0CFHm0cd8dnGCbUW/OcS56jf0gvRk7MA0GCSqGSIb3DQEBCwUAMH4x
|
MIIDhTCCAm0CFHm0cd8dnGCbUW/OcS56jf0gvRk7MA0GCSqGSIb3DQEBCwUAMH4x
|
||||||
CzAJBgNVBAYTAkNOMQswCQYDVQQIDAJHRDETMBEGA1UECgwKb3BlbnAycC5jbjET
|
CzAJBgNVBAYTAkNOMQswCQYDVQQIDAJHRDETMBEGA1UECgwKb3BlbnAycC5jbjET
|
||||||
|
|||||||
+40
-9
@@ -52,18 +52,36 @@ type p2pSDWAN struct {
|
|||||||
internalRoute *IPTree
|
internalRoute *IPTree
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *p2pSDWAN) reset() {
|
||||||
|
gLog.Println(LvINFO, "reset sdwan when network disconnected")
|
||||||
|
// clear sysroute
|
||||||
|
delRoutesByGateway(s.gateway.String())
|
||||||
|
// clear internel route
|
||||||
|
s.internalRoute = NewIPTree("")
|
||||||
|
// clear p2papp
|
||||||
|
for _, node := range gConf.getAddNodes() {
|
||||||
|
gConf.delete(AppConfig{SrcPort: 0, PeerNode: node.Name})
|
||||||
|
}
|
||||||
|
|
||||||
|
gConf.resetSDWAN()
|
||||||
|
}
|
||||||
func (s *p2pSDWAN) init(name string) error {
|
func (s *p2pSDWAN) init(name string) error {
|
||||||
if gConf.getSDWAN().Gateway == "" {
|
if gConf.getSDWAN().Gateway == "" {
|
||||||
gLog.Println(LvDEBUG, "not in sdwan clear all ")
|
gLog.Println(LvDEBUG, "sdwan init: not in sdwan clear all ")
|
||||||
}
|
}
|
||||||
if s.internalRoute == nil {
|
if s.internalRoute == nil {
|
||||||
s.internalRoute = NewIPTree("")
|
s.internalRoute = NewIPTree("")
|
||||||
}
|
}
|
||||||
|
|
||||||
s.nodeName = name
|
s.nodeName = name
|
||||||
s.gateway, s.subnet, _ = net.ParseCIDR(gConf.getSDWAN().Gateway)
|
if gw, sn, err := net.ParseCIDR(gConf.getSDWAN().Gateway); err == nil { // preserve old gateway
|
||||||
|
s.gateway = gw
|
||||||
|
s.subnet = sn
|
||||||
|
}
|
||||||
|
|
||||||
for _, node := range gConf.getDelNodes() {
|
for _, node := range gConf.getDelNodes() {
|
||||||
gLog.Println(LvDEBUG, "deal deleted node: ", node.Name)
|
gLog.Println(LvDEBUG, "sdwan init: deal deleted node: ", node.Name)
|
||||||
|
gLog.Printf(LvDEBUG, "sdwan init: delRoute: %s, %s ", node.IP, s.gateway.String())
|
||||||
delRoute(node.IP, s.gateway.String())
|
delRoute(node.IP, s.gateway.String())
|
||||||
s.internalRoute.Del(node.IP, node.IP)
|
s.internalRoute.Del(node.IP, node.IP)
|
||||||
ipNum, _ := inetAtoN(node.IP)
|
ipNum, _ := inetAtoN(node.IP)
|
||||||
@@ -88,24 +106,26 @@ func (s *p2pSDWAN) init(name string) error {
|
|||||||
}
|
}
|
||||||
s.internalRoute.Del(minIP.String(), maxIP.String())
|
s.internalRoute.Del(minIP.String(), maxIP.String())
|
||||||
delRoute(ipnet.String(), s.gateway.String())
|
delRoute(ipnet.String(), s.gateway.String())
|
||||||
|
gLog.Printf(LvDEBUG, "sdwan init: resource delRoute: %s, %s ", ipnet.String(), s.gateway.String())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
for _, node := range gConf.getAddNodes() {
|
for _, node := range gConf.getAddNodes() {
|
||||||
gLog.Println(LvDEBUG, "deal add node: ", node.Name)
|
gLog.Println(LvDEBUG, "sdwan init: deal add node: ", node.Name)
|
||||||
ipNet := &net.IPNet{
|
ipNet := &net.IPNet{
|
||||||
IP: net.ParseIP(node.IP),
|
IP: net.ParseIP(node.IP),
|
||||||
Mask: s.subnet.Mask,
|
Mask: s.subnet.Mask,
|
||||||
}
|
}
|
||||||
if node.Name == s.nodeName {
|
if node.Name == s.nodeName {
|
||||||
s.virtualIP = ipNet
|
s.virtualIP = ipNet
|
||||||
gLog.Println(LvINFO, "start tun ", ipNet.String())
|
gLog.Println(LvINFO, "sdwan init: start tun ", ipNet.String())
|
||||||
err := s.StartTun()
|
err := s.StartTun()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LvERROR, "start tun error:", err)
|
gLog.Println(LvERROR, "sdwan init: start tun error:", err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
gLog.Println(LvINFO, "start tun ok")
|
gLog.Println(LvINFO, "sdwan init: start tun ok")
|
||||||
allowTunForward()
|
allowTunForward()
|
||||||
|
gLog.Printf(LvDEBUG, "sdwan init: addRoute %s %s %s", s.subnet.String(), s.gateway.String(), s.tun.tunName)
|
||||||
addRoute(s.subnet.String(), s.gateway.String(), s.tun.tunName)
|
addRoute(s.subnet.String(), s.gateway.String(), s.tun.tunName)
|
||||||
// addRoute("255.255.255.255/32", s.gateway.String(), s.tun.tunName) // for broadcast
|
// addRoute("255.255.255.255/32", s.gateway.String(), s.tun.tunName) // for broadcast
|
||||||
// addRoute("224.0.0.0/4", s.gateway.String(), s.tun.tunName) // for multicast
|
// addRoute("224.0.0.0/4", s.gateway.String(), s.tun.tunName) // for multicast
|
||||||
@@ -124,18 +144,28 @@ func (s *p2pSDWAN) init(name string) error {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if len(node.Resource) > 0 {
|
if len(node.Resource) > 0 {
|
||||||
gLog.Printf(LvINFO, "deal add node: %s resource: %s", node.Name, node.Resource)
|
gLog.Printf(LvINFO, "sdwan init: deal add node: %s resource: %s", node.Name, node.Resource)
|
||||||
arr := strings.Split(node.Resource, ",")
|
arr := strings.Split(node.Resource, ",")
|
||||||
for _, r := range arr {
|
for _, r := range arr {
|
||||||
// add internal route
|
// add internal route
|
||||||
_, ipnet, err := net.ParseCIDR(r)
|
_, ipnet, err := net.ParseCIDR(r)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fmt.Println("Error parsing CIDR:", err)
|
fmt.Println("sdwan init: Error parsing CIDR:", err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if ipnet.Contains(net.ParseIP(gConf.Network.localIP)) { // local ip and resource in the same lan
|
if ipnet.Contains(net.ParseIP(gConf.Network.localIP)) { // local ip and resource in the same lan
|
||||||
|
gLog.Printf(LvDEBUG, "sdwan init: local ip %s in this resource %s, ignore", gConf.Network.localIP, ipnet.IP.String())
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
// local net could access this single ip
|
||||||
|
if ipnet.Mask[0] == 255 && ipnet.Mask[1] == 255 && ipnet.Mask[2] == 255 && ipnet.Mask[3] == 255 {
|
||||||
|
gLog.Printf(LvDEBUG, "sdwan init: ping %s start", ipnet.IP.String())
|
||||||
|
if _, err := Ping(ipnet.IP.String()); err == nil {
|
||||||
|
gLog.Printf(LvDEBUG, "sdwan init: ping %s ok, ignore this resource", ipnet.IP.String())
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
gLog.Printf(LvDEBUG, "sdwan init: ping %s failed", ipnet.IP.String())
|
||||||
|
}
|
||||||
minIP := ipnet.IP
|
minIP := ipnet.IP
|
||||||
maxIP := make(net.IP, len(minIP))
|
maxIP := make(net.IP, len(minIP))
|
||||||
copy(maxIP, minIP)
|
copy(maxIP, minIP)
|
||||||
@@ -144,6 +174,7 @@ func (s *p2pSDWAN) init(name string) error {
|
|||||||
}
|
}
|
||||||
s.internalRoute.Add(minIP.String(), maxIP.String(), &sdwanNode{name: node.Name, id: NodeNameToID(node.Name)})
|
s.internalRoute.Add(minIP.String(), maxIP.String(), &sdwanNode{name: node.Name, id: NodeNameToID(node.Name)})
|
||||||
// add sys route
|
// add sys route
|
||||||
|
gLog.Printf(LvDEBUG, "sdwan init: addRoute %s %s %s", ipnet.String(), s.gateway.String(), s.tun.tunName)
|
||||||
addRoute(ipnet.String(), s.gateway.String(), s.tun.tunName)
|
addRoute(ipnet.String(), s.gateway.String(), s.tun.tunName)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+1
-1
@@ -47,7 +47,7 @@ func (vl *v4Listener) handleConnection(c net.Conn) {
|
|||||||
utcp.SetReadDeadline(time.Now().Add(UnderlayTCPConnectTimeout))
|
utcp.SetReadDeadline(time.Now().Add(UnderlayTCPConnectTimeout))
|
||||||
_, buff, err := utcp.ReadBuffer()
|
_, buff, err := utcp.ReadBuffer()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Printf(LvERROR, "utcp.ReadBuffer error:", err)
|
gLog.Println(LvERROR, "utcp.ReadBuffer error:", err)
|
||||||
}
|
}
|
||||||
utcp.WriteBytes(MsgP2P, MsgTunnelHandshakeAck, buff)
|
utcp.WriteBytes(MsgP2P, MsgTunnelHandshakeAck, buff)
|
||||||
var tid uint64
|
var tid uint64
|
||||||
|
|||||||
@@ -1,5 +1,9 @@
|
|||||||
package main
|
package main
|
||||||
|
|
||||||
|
// On Windows env
|
||||||
|
// cd lib
|
||||||
|
// go build -o openp2p.dll -buildmode=c-shared openp2p.go
|
||||||
|
// caller example see example/dll
|
||||||
import (
|
import (
|
||||||
op "openp2p/core"
|
op "openp2p/core"
|
||||||
)
|
)
|
||||||
|
|||||||
Reference in New Issue
Block a user