Compare commits

..
4 Commits
Author SHA1 Message Date
TenderIronh f9b5073e0d fix some bug 2022-05-29 15:48:56 +08:00
TenderIronh 9ea467c7b3 rm rename 2022-05-26 23:43:06 +08:00
TenderIronh c0bad61eb6 tcp punch 2022-05-26 23:18:22 +08:00
TenderIronh bb32133038 config template 2022-05-17 15:53:29 +08:00
23 changed files with 521 additions and 555 deletions
+5
View File
@@ -91,4 +91,9 @@ firewall-cmd --state
## 卸载 ## 卸载
``` ```
./openp2p uninstall ./openp2p uninstall
# 已安装时
# windows
C:\Program Files\OpenP2P\openp2p.exe uninstall
# linux,macos
sudo /usr/local/openp2p/openp2p uninstall
``` ```
+5
View File
@@ -93,4 +93,9 @@ firewall-cmd --state
## Uninstall ## Uninstall
``` ```
./openp2p uninstall ./openp2p uninstall
# when already installed
# windows
C:\Program Files\OpenP2P\openp2p.exe uninstall
# linux,macos
sudo /usr/local/openp2p/openp2p uninstall
``` ```
+9
View File
@@ -188,6 +188,15 @@ func compareVersion(v1, v2 string) int {
return LESS return LESS
} }
func parseMajorVer(ver string) int {
v1Arr := strings.Split(ver, ".")
if len(v1Arr) > 0 {
n, _ := strconv.ParseInt(v1Arr[0], 10, 32)
return int(n)
}
return 0
}
func IsIPv6(address string) bool { func IsIPv6(address string) bool {
return strings.Count(address, ":") >= 2 return strings.Count(address, ":") >= 2
} }
+25
View File
@@ -49,6 +49,11 @@ func assertCompareVersion(t *testing.T, v1 string, v2 string, result int) {
t.Errorf("compare version %s %s fail\n", v1, v2) t.Errorf("compare version %s %s fail\n", v1, v2)
} }
} }
func assertParseMajorVer(t *testing.T, v string, result int) {
if parseMajorVer(v) != result {
t.Errorf("ParseMajorVer %s fail\n", v)
}
}
func TestCompareVersion(t *testing.T) { func TestCompareVersion(t *testing.T) {
// test = // test =
assertCompareVersion(t, "0.98.0", "0.98.0", EQUAL) assertCompareVersion(t, "0.98.0", "0.98.0", EQUAL)
@@ -69,3 +74,23 @@ func TestCompareVersion(t *testing.T) {
assertCompareVersion(t, "", "1.5.0", LESS) assertCompareVersion(t, "", "1.5.0", LESS)
} }
func TestParseMajorVer(t *testing.T) {
assertParseMajorVer(t, "0.98.0", 0)
assertParseMajorVer(t, "0.98", 0)
assertParseMajorVer(t, "1.4.0", 1)
assertParseMajorVer(t, "1.5.0", 1)
assertParseMajorVer(t, "0.98.0.22345", 0)
assertParseMajorVer(t, "1.98.0.12345", 1)
assertParseMajorVer(t, "10.98.0.12345", 10)
assertParseMajorVer(t, "1.4.0", 1)
assertParseMajorVer(t, "1.4", 1)
assertParseMajorVer(t, "1", 1)
assertParseMajorVer(t, "2", 2)
assertParseMajorVer(t, "3", 3)
assertParseMajorVer(t, "2.1.0", 2)
assertParseMajorVer(t, "3.0.0", 3)
}
+30 -28
View File
@@ -11,8 +11,6 @@ import (
var gConf Config var gConf Config
const IntValueNotSet int = -99999999
type AppConfig struct { type AppConfig struct {
// required // required
AppName string AppName string
@@ -28,7 +26,7 @@ type AppConfig struct {
peerToken uint64 peerToken uint64
peerNatType int peerNatType int
hasIPv4 int hasIPv4 int
IPv6 string peerIPv6 string
hasUPNPorNATPMP int hasUPNPorNATPMP int
peerIP string peerIP string
peerConeNatPort int peerConeNatPort int
@@ -39,6 +37,8 @@ type AppConfig struct {
errMsg string errMsg string
connectTime time.Time connectTime time.Time
fromToken uint64 fromToken uint64
linkMode string
isUnderlayServer int // TODO: bool?
} }
// TODO: add loglevel, maxlogfilesize // TODO: add loglevel, maxlogfilesize
@@ -105,10 +105,16 @@ func (c *Config) save() {
} }
} }
func init() {
gConf.LogLevel = 1
gConf.Network.ShareBandwidth = 10
gConf.Network.ServerHost = "api.openp2p.cn"
gConf.Network.ServerPort = WsPort
}
func (c *Config) load() error { func (c *Config) load() error {
c.mtx.Lock() c.mtx.Lock()
c.LogLevel = IntValueNotSet
c.Network.ShareBandwidth = IntValueNotSet
defer c.mtx.Unlock() defer c.mtx.Unlock()
data, err := ioutil.ReadFile("config.json") data, err := ioutil.ReadFile("config.json")
if err != nil { if err != nil {
@@ -154,7 +160,7 @@ type NetworkConfig struct {
publicIP string publicIP string
natType int natType int
hasIPv4 int hasIPv4 int
IPv6 string publicIPv6 string // must lowwer-case not save json
hasUPNPorNATPMP int hasUPNPorNATPMP int
ShareBandwidth int ShareBandwidth int
// server info // server info
@@ -162,11 +168,13 @@ type NetworkConfig struct {
ServerPort int ServerPort int
UDPPort1 int UDPPort1 int
UDPPort2 int UDPPort2 int
TCPPort int
} }
func parseParams(subCommand string) { func parseParams(subCommand string) {
fset := flag.NewFlagSet(subCommand, flag.ExitOnError) fset := flag.NewFlagSet(subCommand, flag.ExitOnError)
serverHost := fset.String("serverhost", "api.openp2p.cn", "server host ") serverHost := fset.String("serverhost", "api.openp2p.cn", "server host ")
serverPort := fset.Int("serverport", WsPort, "server port ")
// serverHost := flag.String("serverhost", "127.0.0.1", "server host ") // for debug // serverHost := flag.String("serverhost", "127.0.0.1", "server host ") // for debug
token := fset.Uint64("token", 0, "token") token := fset.Uint64("token", 0, "token")
node := fset.String("node", "", "node name. 8-31 characters. if not set, it will be hostname") node := fset.String("node", "", "node name. 8-31 characters. if not set, it will be hostname")
@@ -174,13 +182,14 @@ func parseParams(subCommand string) {
dstIP := fset.String("dstip", "127.0.0.1", "destination ip ") dstIP := fset.String("dstip", "127.0.0.1", "destination ip ")
dstPort := fset.Int("dstport", 0, "destination port ") dstPort := fset.Int("dstport", 0, "destination port ")
srcPort := fset.Int("srcport", 0, "source port ") srcPort := fset.Int("srcport", 0, "source port ")
tcpPort := fset.Int("tcpport", 0, "tcp port for upnp or publicip")
protocol := fset.String("protocol", "tcp", "tcp or udp") protocol := fset.String("protocol", "tcp", "tcp or udp")
appName := fset.String("appname", "", "app name") appName := fset.String("appname", "", "app name")
shareBandwidth := fset.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private network no limit") shareBandwidth := fset.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private network no limit")
daemonMode := fset.Bool("d", false, "daemonMode") daemonMode := fset.Bool("d", false, "daemonMode")
notVerbose := fset.Bool("nv", false, "not log console") notVerbose := fset.Bool("nv", false, "not log console")
newconfig := fset.Bool("newconfig", false, "not load existing config.json") newconfig := fset.Bool("newconfig", false, "not load existing config.json")
logLevel := fset.Int("loglevel", 1, "0:debug 1:info 2:warn 3:error") logLevel := fset.Int("loglevel", 0, "0:info 1:warn 2:error 3:debug")
if subCommand == "" { // no subcommand if subCommand == "" { // no subcommand
fset.Parse(os.Args[1:]) fset.Parse(os.Args[1:])
} else { } else {
@@ -197,7 +206,6 @@ func parseParams(subCommand string) {
if !*newconfig { if !*newconfig {
gConf.load() // load old config. otherwise will clear all apps gConf.load() // load old config. otherwise will clear all apps
} }
gConf.LogLevel = *logLevel
if config.SrcPort != 0 { if config.SrcPort != 0 {
gConf.add(config, true) gConf.add(config, true)
} }
@@ -217,14 +225,17 @@ func parseParams(subCommand string) {
if f.Name == "loglevel" { if f.Name == "loglevel" {
gConf.LogLevel = *logLevel gConf.LogLevel = *logLevel
} }
if f.Name == "tcpport" {
gConf.Network.TCPPort = *tcpPort
}
if f.Name == "token" {
gConf.Network.Token = *token
}
}) })
if gConf.Network.ServerHost == "" { if gConf.Network.ServerHost == "" {
gConf.Network.ServerHost = *serverHost gConf.Network.ServerHost = *serverHost
} }
if *token != 0 {
gConf.Network.Token = *token
}
if *node != "" { if *node != "" {
if len(*node) < 8 { if len(*node) < 8 {
gLog.Println(LvERROR, ErrNodeTooShort) gLog.Println(LvERROR, ErrNodeTooShort)
@@ -236,16 +247,17 @@ func parseParams(subCommand string) {
gConf.Network.Node = defaultNodeName() gConf.Network.Node = defaultNodeName()
} }
} }
if gConf.LogLevel == IntValueNotSet { if gConf.Network.TCPPort == 0 {
gConf.LogLevel = *logLevel if *tcpPort == 0 {
p := int(nodeNameToID(gConf.Network.Node)%15000 + 50000)
tcpPort = &p
} }
if gConf.Network.ShareBandwidth == IntValueNotSet { gConf.Network.TCPPort = *tcpPort
gConf.Network.ShareBandwidth = *shareBandwidth
} }
gConf.Network.ServerPort = 27183 gConf.Network.ServerPort = *serverPort
gConf.Network.UDPPort1 = 27182 gConf.Network.UDPPort1 = UDPPort1
gConf.Network.UDPPort2 = 27183 gConf.Network.UDPPort2 = UDPPort2
gLog.setLevel(LogLevel(gConf.LogLevel)) gLog.setLevel(LogLevel(gConf.LogLevel))
if *notVerbose { if *notVerbose {
gLog.setMode(LogFile) gLog.setMode(LogFile)
@@ -253,13 +265,3 @@ func parseParams(subCommand string) {
// gConf.mtx.Unlock() // gConf.mtx.Unlock()
gConf.save() gConf.save()
} }
func (conf *AppConfig) isSupportTCP(pnConf NetworkConfig) bool {
if conf.peerVersion == "" || compareVersion(conf.peerVersion, LeastSupportTCPVersion) == LESS {
return false
}
if pnConf.hasIPv4 == 1 || pnConf.hasUPNPorNATPMP == 1 || conf.hasIPv4 == 1 || conf.hasUPNPorNATPMP == 1 || (IsIPv6(pnConf.IPv6) && IsIPv6(conf.IPv6)) {
return true
}
return false
}
+3 -4
View File
@@ -1,10 +1,9 @@
{ {
"network": { "network": {
"Node": "YOUR_NODE_NAME", "Node": "YOUR_NODE_NAME",
"User": "YOUR_USER_NAME", "Token": "YOUR_TOKEN",
"Password": "YOUR_PASSWORD", "ServerHost": "api.openp2p.cn",
"ServerHost": "openp2p.cn", "ServerPort": 27183,
"ServerPort": 27182,
"UDPPort1": 27182, "UDPPort1": 27182,
"UDPPort2": 27183 "UDPPort2": 27183
}, },
+3
View File
@@ -13,4 +13,7 @@ var (
ErrorNewUser = errors.New("new user") ErrorNewUser = errors.New("new user")
ErrorLogin = errors.New("user or password not correct") ErrorLogin = errors.New("user or password not correct")
ErrNodeTooShort = errors.New("node name too short, it must >=8 charaters") ErrNodeTooShort = errors.New("node name too short, it must >=8 charaters")
ErrPeerOffline = errors.New("peer offline")
ErrMsgFormat = errors.New("message format wrong")
ErrVersionNotCompatible = errors.New("version not compatible")
) )
+2 -1
View File
@@ -6,7 +6,8 @@ require (
github.com/gorilla/websocket v1.4.2 github.com/gorilla/websocket v1.4.2
github.com/kardianos/service v1.2.0 github.com/kardianos/service v1.2.0
github.com/lucas-clemente/quic-go v0.27.0 github.com/lucas-clemente/quic-go v0.27.0
golang.org/x/sys v0.0.0-20210906170528-6f6e22806c34 github.com/openp2p-cn/go-reuseport v0.3.2
golang.org/x/sys v0.0.0-20220422013727-9388b58f7150
) )
require ( require (
+62 -4
View File
@@ -19,7 +19,7 @@ func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
} }
gLog.Printf(LvDEBUG, "handle push msg type:%d, push header:%+v", subType, pushHead) gLog.Printf(LvDEBUG, "handle push msg type:%d, push header:%+v", subType, pushHead)
switch subType { switch subType {
case MsgPushConnectReq: case MsgPushConnectReq: // TODO: handle a msg move to a new function
req := PushConnectReq{} req := PushConnectReq{}
err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req) err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req)
if err != nil { if err != nil {
@@ -28,6 +28,17 @@ func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
} }
gLog.Printf(LvINFO, "%s is connecting...", req.From) gLog.Printf(LvINFO, "%s is connecting...", req.From)
gLog.Println(LvDEBUG, "push connect response to ", req.From) gLog.Println(LvDEBUG, "push connect response to ", req.From)
if compareVersion(req.Version, LeastSupportVersion) == LESS {
gLog.Println(LvERROR, ErrVersionNotCompatible.Error(), ":", req.From)
rsp := PushConnectRsp{
Error: 10,
Detail: ErrVersionNotCompatible.Error(),
To: req.From,
From: pn.config.Node,
}
pn.push(req.From, MsgPushConnectRsp, rsp)
return ErrVersionNotCompatible
}
// verify totp token or token // verify totp token or token
if VerifyTOTP(req.Token, pn.config.Token, time.Now().Unix()+(pn.serverTs-pn.localTs)) || // localTs may behind, auto adjust ts if VerifyTOTP(req.Token, pn.config.Token, time.Now().Unix()+(pn.serverTs-pn.localTs)) || // localTs may behind, auto adjust ts
VerifyTOTP(req.Token, pn.config.Token, time.Now().Unix()) { VerifyTOTP(req.Token, pn.config.Token, time.Now().Unix()) {
@@ -39,7 +50,11 @@ func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
config.PeerNode = req.From config.PeerNode = req.From
config.peerVersion = req.Version config.peerVersion = req.Version
config.fromToken = req.Token config.fromToken = req.Token
config.IPv6 = req.IPv6 config.peerIPv6 = req.IPv6
config.hasIPv4 = req.HasIPv4
config.hasUPNPorNATPMP = req.HasUPNPorNATPMP
config.linkMode = req.LinkMode
config.isUnderlayServer = req.IsUnderlayServer
// share relay node will limit bandwidth // share relay node will limit bandwidth
if req.Token != pn.config.Token { if req.Token != pn.config.Token {
gLog.Printf(LvINFO, "set share bandwidth %d mbps", pn.config.ShareBandwidth) gLog.Printf(LvINFO, "set share bandwidth %d mbps", pn.config.ShareBandwidth)
@@ -98,7 +113,7 @@ func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
SaveKey(req.AppID, req.AppKey) SaveKey(req.AppID, req.AppKey)
case MsgPushUpdate: case MsgPushUpdate:
gLog.Println(LvINFO, "MsgPushUpdate") gLog.Println(LvINFO, "MsgPushUpdate")
update() // download new version first, then exec ./openp2p update update(pn.config.ServerHost, pn.config.ServerPort) // download new version first, then exec ./openp2p update
targetPath := filepath.Join(defaultInstallPath, defaultBinName) targetPath := filepath.Join(defaultInstallPath, defaultBinName)
args := []string{"update"} args := []string{"update"}
env := os.Environ() env := os.Environ()
@@ -134,7 +149,7 @@ func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
} }
relayNode = app.relayNode relayNode = app.relayNode
relayMode = app.relayMode relayMode = app.relayMode
linkMode = app.tunnel.linkMode linkMode = app.tunnel.linkModeWeb
} }
appInfo := AppInfo{ appInfo := AppInfo{
AppName: config.AppName, AppName: config.AppName,
@@ -158,6 +173,49 @@ func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
req.Apps = append(req.Apps, appInfo) req.Apps = append(req.Apps, appInfo)
} }
pn.write(MsgReport, MsgReportApps, &req) pn.write(MsgReport, MsgReportApps, &req)
case MsgPushReportLog:
gLog.Println(LvINFO, "MsgPushReportLog")
req := ReportLogReq{}
err := json.Unmarshal(msg[openP2PHeaderSize:], &req)
if err != nil {
gLog.Printf(LvERROR, "wrong MsgPushReportLog:%s %s", err, string(msg[openP2PHeaderSize:]))
return err
}
if req.FileName == "" {
req.FileName = "openp2p.log"
}
f, err := os.Open(filepath.Join("log", req.FileName))
if err != nil {
gLog.Println(LvERROR, "read log file error:", err)
break
}
fi, err := f.Stat()
if err != nil {
break
}
if req.Offset == 0 && fi.Size() > 4096 {
req.Offset = fi.Size() - 4096
}
if req.Len <= 0 {
req.Len = 4096
}
f.Seek(req.Offset, 0)
if req.Len > 1024*1024 { // too large
break
}
buff := make([]byte, req.Len)
readLength, err := f.Read(buff)
f.Close()
if err != nil {
gLog.Println(LvERROR, "read log content error:", err)
break
}
rsp := ReportLogRsp{}
rsp.Content = string(buff[:readLength])
rsp.FileName = req.FileName
rsp.Total = fi.Size()
rsp.Len = req.Len
pn.write(MsgReport, MsgPushReportLog, &rsp)
case MsgPushEditApp: case MsgPushEditApp:
gLog.Println(LvINFO, "MsgPushEditApp") gLog.Println(LvINFO, "MsgPushEditApp")
newApp := AppInfo{} newApp := AppInfo{}
+6 -2
View File
@@ -12,7 +12,7 @@ import (
func handshakeC2C(t *P2PTunnel) (err error) { func handshakeC2C(t *P2PTunnel) (err error) {
gLog.Printf(LvDEBUG, "handshakeC2C %s:%d:%d to %s:%d", t.pn.config.Node, t.coneLocalPort, t.coneNatPort, t.config.peerIP, t.config.peerConeNatPort) gLog.Printf(LvDEBUG, "handshakeC2C %s:%d:%d to %s:%d", t.pn.config.Node, t.coneLocalPort, t.coneNatPort, t.config.peerIP, t.config.peerConeNatPort)
defer gLog.Printf(LvDEBUG, "handshakeC2C ok") defer gLog.Printf(LvDEBUG, "handshakeC2C end")
conn, err := net.ListenUDP("udp", t.la) conn, err := net.ListenUDP("udp", t.la)
if err != nil { if err != nil {
return err return err
@@ -45,6 +45,7 @@ func handshakeC2C(t *P2PTunnel) (err error) {
} }
if head.MainType == MsgP2P && head.SubType == MsgPunchHandshakeAck { if head.MainType == MsgP2P && head.SubType == MsgPunchHandshakeAck {
gLog.Printf(LvDEBUG, "read %d handshake ack ", t.id) gLog.Printf(LvDEBUG, "read %d handshake ack ", t.id)
gLog.Printf(LvINFO, "handshakeC2C ok")
return nil return nil
} }
} }
@@ -54,9 +55,10 @@ func handshakeC2C(t *P2PTunnel) (err error) {
_, err = UDPWrite(conn, t.ra, MsgP2P, MsgPunchHandshakeAck, P2PHandshakeReq{ID: t.id}) _, err = UDPWrite(conn, t.ra, MsgP2P, MsgPunchHandshakeAck, P2PHandshakeReq{ID: t.id})
if err != nil { if err != nil {
gLog.Println(LvDEBUG, "handshakeC2C write MsgPunchHandshakeAck error", err) gLog.Println(LvDEBUG, "handshakeC2C write MsgPunchHandshakeAck error", err)
}
return err return err
} }
}
gLog.Printf(LvINFO, "handshakeC2C ok")
return nil return nil
} }
@@ -115,6 +117,7 @@ func handshakeC2S(t *P2PTunnel) error {
_, err = UDPWrite(conn, dst, MsgP2P, MsgPunchHandshakeAck, P2PHandshakeReq{ID: t.id}) _, err = UDPWrite(conn, dst, MsgP2P, MsgPunchHandshakeAck, P2PHandshakeReq{ID: t.id})
return err return err
} }
gLog.Printf(LvINFO, "handshakeC2S ok")
return nil return nil
} }
@@ -175,6 +178,7 @@ func handshakeS2C(t *P2PTunnel) error {
return fmt.Errorf("wait handshake failed") return fmt.Errorf("wait handshake failed")
case la := <-gotCh: case la := <-gotCh:
gLog.Println(LvDEBUG, "symmetric handshake ok", la) gLog.Println(LvDEBUG, "symmetric handshake ok", la)
gLog.Printf(LvINFO, "handshakeS2C ok")
} }
return nil return nil
} }
+1 -1
View File
@@ -17,7 +17,7 @@ import (
// ./openp2p install -node hhd1207-222 -token YOUR-TOKEN -sharebandwidth 0 -peernode hhdhome-n1 -dstip 127.0.0.1 -dstport 50022 -protocol tcp -srcport 22 // ./openp2p install -node hhd1207-222 -token YOUR-TOKEN -sharebandwidth 0 -peernode hhdhome-n1 -dstip 127.0.0.1 -dstport 50022 -protocol tcp -srcport 22
func install() { func install() {
gLog.Println(LvINFO, "openp2p start. version: ", OpenP2PVersion) gLog.Println(LvINFO, "openp2p start. version: ", OpenP2PVersion)
gLog.Println(LvINFO, "Contact: QQ Group: 16947733, Email: [email protected]") gLog.Println(LvINFO, "Contact: QQ: 477503927, Email: [email protected]")
gLog.Println(LvINFO, "install start") gLog.Println(LvINFO, "install start")
defer gLog.Println(LvINFO, "install end") defer gLog.Println(LvINFO, "install end")
// auto uninstall // auto uninstall
+63 -9
View File
@@ -6,12 +6,55 @@ import (
"log" "log"
"math/rand" "math/rand"
"net" "net"
"strconv"
"strings"
"time" "time"
reuse "github.com/openp2p-cn/go-reuseport"
) )
var echoConn *net.UDPConn
func natTCP(serverHost string, serverPort int, localPort int) (publicIP string, publicPort int) {
// dialer := &net.Dialer{
// LocalAddr: &net.TCPAddr{
// IP: net.ParseIP("0.0.0.0"),
// Port: localPort,
// },
// }
conn, err := reuse.DialTimeout("tcp4", fmt.Sprintf("%s:%d", "0.0.0.0", localPort), fmt.Sprintf("%s:%d", serverHost, serverPort), time.Second*5)
// conn, err := net.Dial("tcp4", fmt.Sprintf("%s:%d", serverHost, serverPort))
if err != nil {
fmt.Printf("Dial tcp4 %s:%d error:%s", serverHost, serverPort, err)
return
}
defer conn.Close()
_, wrerr := conn.Write([]byte("1"))
if wrerr != nil {
fmt.Printf("Write error: %s\n", wrerr)
return
}
b := make([]byte, 1000)
conn.SetReadDeadline(time.Now().Add(time.Second * 5))
n, rderr := conn.Read(b)
if rderr != nil {
fmt.Printf("Read error: %s\n", rderr)
return
}
arr := strings.Split(string(b[:n]), ":")
if len(arr) < 2 {
return
}
publicIP = arr[0]
port, _ := strconv.ParseInt(arr[1], 10, 32)
publicPort = int(port)
return
}
func natTest(serverHost string, serverPort int, localPort int, echoPort int) (publicIP string, hasPublicIP int, hasUPNPorNATPMP int, publicPort int, err error) { func natTest(serverHost string, serverPort int, localPort int, echoPort int) (publicIP string, hasPublicIP int, hasUPNPorNATPMP int, publicPort int, err error) {
conn, err := net.ListenPacket("udp", fmt.Sprintf(":%d", localPort)) conn, err := net.ListenPacket("udp", fmt.Sprintf(":%d", localPort))
if err != nil { if err != nil {
gLog.Println(LvERROR, "natTest listen udp error:", err)
return "", 0, 0, 0, err return "", 0, 0, 0, err
} }
defer conn.Close() defer conn.Close()
@@ -48,7 +91,7 @@ func natTest(serverHost string, serverPort int, localPort int, echoPort int) (pu
if i == 1 { if i == 1 {
// test upnp or nat-pmp // test upnp or nat-pmp
nat, err := Discover() nat, err := Discover()
if err != nil { if err != nil || nat == nil {
gLog.Println(LvDEBUG, "could not perform UPNP discover:", err) gLog.Println(LvDEBUG, "could not perform UPNP discover:", err)
break break
} }
@@ -59,10 +102,12 @@ func natTest(serverHost string, serverPort int, localPort int, echoPort int) (pu
} }
log.Println("PublicIP:", ext) log.Println("PublicIP:", ext)
externalPort, err := nat.AddPortMapping("udp", echoPort, echoPort, "openp2p", 60) externalPort, err := nat.AddPortMapping("udp", echoPort, echoPort, "openp2p", 30)
if err != nil { if err != nil {
gLog.Println(LvDEBUG, "could not add udp UPNP port mapping", externalPort) gLog.Println(LvDEBUG, "could not add udp UPNP port mapping", externalPort)
break break
} else {
nat.AddPortMapping("tcp", echoPort, echoPort, "openp2p", 604800)
} }
} }
gLog.Printf(LvDEBUG, "public ip test start %s:%d", natRsp.IP, echoPort) gLog.Printf(LvDEBUG, "public ip test start %s:%d", natRsp.IP, echoPort)
@@ -98,14 +143,22 @@ func natTest(serverHost string, serverPort int, localPort int, echoPort int) (pu
func getNATType(host string, udp1 int, udp2 int) (publicIP string, NATType int, hasIPvr int, hasUPNPorNATPMP int, err error) { func getNATType(host string, udp1 int, udp2 int) (publicIP string, NATType int, hasIPvr int, hasUPNPorNATPMP int, err error) {
// the random local port may be used by other. // the random local port may be used by other.
localPort := int(rand.Uint32()%10000 + 50000) localPort := int(rand.Uint32()%15000 + 50000)
echoPort := int(rand.Uint32()%10000 + 50000) echoPort := P2PNetworkInstance(nil).config.TCPPort
go echo(echoPort) go echo(echoPort)
// _, natPort := natTCP(host, 27181, localPort)
// gLog.Println(LvINFO, "nattcp:", natPort)
// _, natPort = natTCP(host, 27180, localPort)
// gLog.Println(LvINFO, "nattcp:", natPort)
ip1, hasIPv4, hasUPNPorNATPMP, port1, err := natTest(host, udp1, localPort, echoPort) ip1, hasIPv4, hasUPNPorNATPMP, port1, err := natTest(host, udp1, localPort, echoPort)
gLog.Printf(LvDEBUG, "local port:%d nat port:%d", localPort, port1) gLog.Printf(LvDEBUG, "local port:%d nat port:%d", localPort, port1)
if err != nil { if err != nil {
return "", 0, hasIPv4, hasUPNPorNATPMP, err return "", 0, hasIPv4, hasUPNPorNATPMP, err
} }
if echoConn != nil {
echoConn.Close()
echoConn = nil
}
// if hasPublicIP == 1 || hasUPNPorNATPMP == 1 { // if hasPublicIP == 1 || hasUPNPorNATPMP == 1 {
// return ip1, NATNone, hasUPNPorNATPMP, nil // return ip1, NATNone, hasUPNPorNATPMP, nil
// } // }
@@ -122,18 +175,19 @@ func getNATType(host string, udp1 int, udp2 int) (publicIP string, NATType int,
} }
func echo(echoPort int) { func echo(echoPort int) {
conn, err := net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4zero, Port: echoPort}) var err error
echoConn, err = net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4zero, Port: echoPort})
if err != nil { if err != nil {
gLog.Println(LvERROR, "echo server listen error:", err) gLog.Println(LvERROR, "echo server listen error:", err)
return return
} }
buf := make([]byte, 1600) buf := make([]byte, 1600)
defer conn.Close() // close outside for breaking the ReadFromUDP
// wait 5s for echo testing // wait 5s for echo testing
conn.SetReadDeadline(time.Now().Add(time.Second * 30)) echoConn.SetReadDeadline(time.Now().Add(time.Second * 30))
n, addr, err := conn.ReadFromUDP(buf) n, addr, err := echoConn.ReadFromUDP(buf)
if err != nil { if err != nil {
return return
} }
conn.WriteToUDP(buf[0:n], addr) echoConn.WriteToUDP(buf[0:n], addr)
} }
+1 -1
View File
@@ -42,7 +42,7 @@ func main() {
} }
parseParams("") parseParams("")
gLog.Println(LvINFO, "openp2p start. version: ", OpenP2PVersion) gLog.Println(LvINFO, "openp2p start. version: ", OpenP2PVersion)
gLog.Println(LvINFO, "Contact: QQ Group: 16947733, Email: [email protected]") gLog.Println(LvINFO, "Contact: QQ: 477503927, Email: [email protected]")
if gConf.daemonMode { if gConf.daemonMode {
d := daemon{} d := daemon{}
+75 -17
View File
@@ -160,6 +160,7 @@ 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.PeerNode)
defer gLog.Printf(LvINFO, "addRelayTunnel to %s end", config.PeerNode) defer gLog.Printf(LvINFO, "addRelayTunnel to %s end", config.PeerNode)
// request a relay node or specify manually(TODO)
pn.write(MsgRelay, MsgRelayNodeReq, &RelayNodeReq{config.PeerNode}) pn.write(MsgRelay, MsgRelayNodeReq, &RelayNodeReq{config.PeerNode})
head, body := pn.read("", MsgRelay, MsgRelayNodeRsp, time.Second*10) head, body := pn.read("", MsgRelay, MsgRelayNodeRsp, time.Second*10)
if head == nil { if head == nil {
@@ -179,6 +180,7 @@ func (pn *P2PNetwork) addRelayTunnel(config AppConfig) (*P2PTunnel, uint64, stri
relayConfig := config relayConfig := config
relayConfig.PeerNode = rsp.RelayName relayConfig.PeerNode = rsp.RelayName
relayConfig.peerToken = rsp.RelayToken relayConfig.peerToken = rsp.RelayToken
///
t, err := pn.addDirectTunnel(relayConfig, 0) t, err := pn.addDirectTunnel(relayConfig, 0)
if err != nil { if err != nil {
gLog.Println(LvERROR, "direct connect error:", err) gLog.Println(LvERROR, "direct connect error:", err)
@@ -238,13 +240,7 @@ func (pn *P2PNetwork) AddApp(config AppConfig) error {
peerIP = t.config.peerIP peerIP = t.config.peerIP
} }
// TODO: if tcp failed, should try udp punching, nattype should refactor also, when NATNONE and failed we don't know the peerNatType // TODO: if tcp failed, should try udp punching, nattype should refactor also, when NATNONE and failed we don't know the peerNatType
// if err != nil && err == ErrorHandshake && t.isSupportTCP() {
// t, err = pn.addDirectTunnel(config, 0)
// if t != nil {
// peerNatType = t.config.peerNatType
// peerIP = t.config.peerIP
// }
// }
if err != nil && err == ErrorHandshake { if err != nil && err == ErrorHandshake {
gLog.Println(LvERROR, "direct connect failed, try to relay") gLog.Println(LvERROR, "direct connect failed, try to relay")
t, rtid, relayMode, err = pn.addRelayTunnel(config) t, rtid, relayMode, err = pn.addRelayTunnel(config)
@@ -349,8 +345,10 @@ func (pn *P2PNetwork) addDirectTunnel(config AppConfig, tid uint64) (*P2PTunnel,
} }
return true return true
}) })
if exist {
return t, nil
}
// create tunnel if not exist // create tunnel if not exist
if !exist {
t = &P2PTunnel{pn: pn, t = &P2PTunnel{pn: pn,
config: config, config: config,
id: tid, id: tid,
@@ -358,25 +356,85 @@ func (pn *P2PNetwork) addDirectTunnel(config AppConfig, tid uint64) (*P2PTunnel,
pn.msgMapMtx.Lock() pn.msgMapMtx.Lock()
pn.msgMap[nodeNameToID(config.PeerNode)] = make(chan []byte, 50) pn.msgMap[nodeNameToID(config.PeerNode)] = make(chan []byte, 50)
pn.msgMapMtx.Unlock() pn.msgMapMtx.Unlock()
t.init() // server side
if !isClient {
err := pn.newTunnel(t, tid, isClient)
return t, err // always return
}
// client side
// peer info
initErr := t.requestPeerInfo()
if initErr != nil {
gLog.Println(LvERROR, "init error:", initErr)
return nil, initErr
}
err := ErrorHandshake
// try TCP6
if IsIPv6(t.config.peerIPv6) && IsIPv6(t.pn.config.publicIPv6) {
gLog.Println(LvINFO, "try TCP6")
t.config.linkMode = LinkModeTCP6
t.config.isUnderlayServer = 0
if err = pn.newTunnel(t, tid, isClient); err == nil {
return t, nil
}
}
// TODO: try UDP6
// try TCP4
if t.config.hasIPv4 == 1 || t.pn.config.hasIPv4 == 1 || t.config.hasUPNPorNATPMP == 1 || t.pn.config.hasUPNPorNATPMP == 1 {
gLog.Println(LvINFO, "try TCP4")
t.config.linkMode = LinkModeTCP4
if t.config.hasIPv4 == 1 || t.config.hasUPNPorNATPMP == 1 {
t.config.isUnderlayServer = 0
} else {
t.config.isUnderlayServer = 1
}
if err = pn.newTunnel(t, tid, isClient); err == nil {
return t, nil
}
}
// TODO: try UDP4
// try TCPPunch
if t.config.peerNatType == NATCone && t.pn.config.natType == NATCone { // TODO: support c2s
gLog.Println(LvINFO, "try TCP4 Punch")
t.config.linkMode = LinkModeTCPPunch
t.config.isUnderlayServer = 0
if err = pn.newTunnel(t, tid, isClient); err == nil {
return t, nil
}
}
// try UDPPunch
if t.config.peerNatType == NATCone || t.pn.config.natType == NATCone {
gLog.Println(LvINFO, "try UDP4 Punch")
t.config.linkMode = LinkModeUDPPunch
t.config.isUnderlayServer = 0
if err = pn.newTunnel(t, tid, isClient); err == nil {
return t, nil
}
}
return nil, err
}
func (pn *P2PNetwork) newTunnel(t *P2PTunnel, tid uint64, isClient bool) error {
t.initPort()
if isClient { if isClient {
if err := t.connect(); err != nil { if err := t.connect(); err != nil {
gLog.Println(LvERROR, "p2pTunnel connect error:", err) gLog.Println(LvERROR, "p2pTunnel connect error:", err)
return t, err return err
} }
} else { } else {
if err := t.listen(); err != nil { if err := t.listen(); err != nil {
gLog.Println(LvERROR, "p2pTunnel listen error:", err) gLog.Println(LvERROR, "p2pTunnel listen error:", err)
return t, err return err
}
} }
} }
// store it when success // store it when success
gLog.Printf(LvDEBUG, "store tunnel %d", tid) gLog.Printf(LvDEBUG, "store tunnel %d", tid)
pn.allTunnels.Store(tid, t) pn.allTunnels.Store(tid, t)
return t, nil return nil
} }
func (pn *P2PNetwork) init() error { func (pn *P2PNetwork) init() error {
gLog.Println(LvINFO, "init start") gLog.Println(LvINFO, "init start")
var err error var err error
@@ -443,7 +501,7 @@ func (pn *P2PNetwork) init() error {
gLog.Println(LvDEBUG, "netinfo:", rsp) gLog.Println(LvDEBUG, "netinfo:", rsp)
if rsp != nil && rsp.Country != "" { if rsp != nil && rsp.Country != "" {
if IsIPv6(rsp.IP.String()) { if IsIPv6(rsp.IP.String()) {
pn.config.IPv6 = rsp.IP.String() pn.config.publicIPv6 = rsp.IP.String()
req.IPv6 = rsp.IP.String() req.IPv6 = rsp.IP.String()
} }
req.NetInfo = *rsp req.NetInfo = *rsp
@@ -627,7 +685,7 @@ func (pn *P2PNetwork) updateAppHeartbeat(appID uint64) {
} }
func (pn *P2PNetwork) refreshIPv6() { func (pn *P2PNetwork) refreshIPv6() {
if !IsIPv6(pn.config.IPv6) { // not support ipv6, not refresh if !IsIPv6(pn.config.publicIPv6) { // not support ipv6, not refresh
return return
} }
client := &http.Client{Timeout: time.Second * 10} client := &http.Client{Timeout: time.Second * 10}
@@ -643,5 +701,5 @@ func (pn *P2PNetwork) refreshIPv6() {
gLog.Println(LvINFO, "netInfo error:", err, n) gLog.Println(LvINFO, "netInfo error:", err, n)
return return
} }
pn.config.IPv6 = string(buf[:n]) pn.config.publicIPv6 = string(buf[:n])
} }
+80 -66
View File
@@ -28,45 +28,70 @@ type P2PTunnel struct {
tunnelServer bool // different from underlayServer tunnelServer bool // different from underlayServer
coneLocalPort int coneLocalPort int
coneNatPort int coneNatPort int
linkMode string linkModeWeb string // use config.linkmode
} }
func (t *P2PTunnel) init() { func (t *P2PTunnel) requestPeerInfo() error {
// request peer info
t.pn.write(MsgQuery, MsgQueryPeerInfoReq, &QueryPeerInfoReq{t.pn.config.Token, t.config.PeerNode})
head, body := t.pn.read("", MsgQuery, MsgQueryPeerInfoRsp, time.Second*10)
if head == nil {
return ErrPeerOffline
}
rsp := QueryPeerInfoRsp{}
err := json.Unmarshal(body, &rsp)
if err != nil {
gLog.Printf(LvERROR, "wrong QueryPeerInfoRsp:%s", err)
return ErrMsgFormat
}
if rsp.Online == 0 {
return ErrPeerOffline
}
if compareVersion(rsp.Version, LeastSupportVersion) == LESS {
return ErrVersionNotCompatible
}
t.config.peerVersion = rsp.Version
t.config.hasIPv4 = rsp.HasIPv4
t.config.peerIP = rsp.IPv4
t.config.peerIPv6 = rsp.IPv6
t.config.hasUPNPorNATPMP = rsp.HasUPNPorNATPMP
t.config.peerNatType = rsp.NatType
///
return nil
}
func (t *P2PTunnel) initPort() {
t.running = true t.running = true
t.hbMtx.Lock() t.hbMtx.Lock()
t.hbTime = time.Now() t.hbTime = time.Now()
t.hbMtx.Unlock() t.hbMtx.Unlock()
t.hbTimeRelay = time.Now().Add(time.Second * 600) // TODO: test fake time t.hbTimeRelay = time.Now().Add(time.Second * 600) // TODO: test fake time
localPort := int(rand.Uint32()%10000 + 50000) localPort := int(rand.Uint32()%15000 + 50000) // if the process has bug, will add many upnp port. use specify p2p port by param
if t.pn.config.natType == NATCone { if t.config.linkMode == LinkModeTCP6 {
t.pn.refreshIPv6()
}
if t.config.linkMode == LinkModeTCP6 || t.config.linkMode == LinkModeTCP4 {
t.coneLocalPort = t.pn.config.TCPPort
t.coneNatPort = t.pn.config.TCPPort // symmetric doesn't need coneNatPort
}
if t.config.linkMode == LinkModeUDPPunch {
// prepare one random cone hole // prepare one random cone hole
_, _, _, port1, _ := natTest(t.pn.config.ServerHost, t.pn.config.UDPPort1, localPort, 0) _, _, _, natPort, _ := natTest(t.pn.config.ServerHost, t.pn.config.UDPPort1, localPort, 0)
t.coneLocalPort = localPort t.coneLocalPort = localPort
t.coneNatPort = port1 t.coneNatPort = natPort
t.la = &net.UDPAddr{IP: net.ParseIP(t.pn.config.localIP), Port: t.coneLocalPort} }
} else { if t.config.linkMode == LinkModeTCPPunch {
// prepare one random cone hole
_, natPort := natTCP(t.pn.config.ServerHost, IfconfigPort1, localPort)
t.coneLocalPort = localPort t.coneLocalPort = localPort
t.coneNatPort = localPort // symmetric doesn't need coneNatPort t.coneNatPort = natPort
if t.pn.config.hasUPNPorNATPMP == 1 {
nat, err := Discover()
if err != nil {
gLog.Println(LvDEBUG, "could not perform UPNP discover:", err)
} else {
externalPort, err := nat.AddPortMapping("tcp", localPort, localPort, "openp2p", 30) // timeout the connection still alive, make the timeout short
if err != nil {
gLog.Println(LvDEBUG, "could not add udp UPNP port mapping", externalPort)
}
}
} }
t.la = &net.UDPAddr{IP: net.ParseIP(t.pn.config.localIP), Port: t.coneLocalPort} t.la = &net.UDPAddr{IP: net.ParseIP(t.pn.config.localIP), Port: t.coneLocalPort}
}
gLog.Printf(LvDEBUG, "prepare punching port %d:%d", t.coneLocalPort, t.coneNatPort) gLog.Printf(LvDEBUG, "prepare punching port %d:%d", t.coneLocalPort, t.coneNatPort)
} }
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.PeerNode)
t.tunnelServer = false t.tunnelServer = false
t.pn.refreshIPv6()
appKey := uint64(0) appKey := uint64(0)
req := PushConnectReq{ req := PushConnectReq{
Token: t.config.peerToken, Token: t.config.peerToken,
@@ -75,11 +100,13 @@ func (t *P2PTunnel) connect() error {
ConeNatPort: t.coneNatPort, ConeNatPort: t.coneNatPort,
NatType: t.pn.config.natType, NatType: t.pn.config.natType,
HasIPv4: t.pn.config.hasIPv4, HasIPv4: t.pn.config.hasIPv4,
IPv6: t.pn.config.IPv6, IPv6: t.pn.config.publicIPv6,
HasUPNPorNATPMP: t.pn.config.hasUPNPorNATPMP, HasUPNPorNATPMP: t.pn.config.hasUPNPorNATPMP,
ID: t.id, ID: t.id,
AppKey: appKey, AppKey: appKey,
Version: OpenP2PVersion, Version: OpenP2PVersion,
LinkMode: t.config.linkMode,
IsUnderlayServer: t.config.isUnderlayServer ^ 1,
} }
if req.Token == 0 { // no relay token if req.Token == 0 { // no relay token
req.Token = t.pn.config.Token req.Token = t.pn.config.Token
@@ -101,7 +128,7 @@ func (t *P2PTunnel) connect() error {
} }
t.config.peerNatType = rsp.NatType t.config.peerNatType = rsp.NatType
t.config.hasIPv4 = rsp.HasIPv4 t.config.hasIPv4 = rsp.HasIPv4
t.config.IPv6 = rsp.IPv6 t.config.peerIPv6 = rsp.IPv6
t.config.hasUPNPorNATPMP = rsp.HasUPNPorNATPMP t.config.hasUPNPorNATPMP = rsp.HasUPNPorNATPMP
t.config.peerVersion = rsp.Version t.config.peerVersion = rsp.Version
t.config.peerConeNatPort = rsp.ConeNatPort t.config.peerConeNatPort = rsp.ConeNatPort
@@ -162,7 +189,7 @@ func (t *P2PTunnel) close() {
} }
func (t *P2PTunnel) start() error { func (t *P2PTunnel) start() error {
if !t.config.isSupportTCP(t.pn.config) { if t.config.linkMode == LinkModeUDPPunch {
if err := t.handshake(); err != nil { if err := t.handshake(); err != nil {
return err return err
} }
@@ -176,7 +203,7 @@ func (t *P2PTunnel) start() error {
} }
func (t *P2PTunnel) handshake() error { func (t *P2PTunnel) handshake() error {
if t.config.peerConeNatPort > 0 { if t.config.peerConeNatPort > 0 { // only peer is cone should prepare t.ra
var err error var err error
t.ra, err = net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", t.config.peerIP, t.config.peerConeNatPort)) t.ra, err = net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", t.config.peerIP, t.config.peerConeNatPort))
if err != nil { if err != nil {
@@ -207,23 +234,22 @@ func (t *P2PTunnel) handshake() error {
} }
func (t *P2PTunnel) connectUnderlay() (err error) { func (t *P2PTunnel) connectUnderlay() (err error) {
if !t.config.isSupportTCP(t.pn.config) { switch t.config.linkMode {
t.conn, err = t.connectUnderlayQuic() case LinkModeTCP6:
if err != nil {
return err
}
} else {
if IsIPv6(t.pn.config.IPv6) && IsIPv6(t.config.IPv6) { // both have ipv6
t.conn, err = t.connectUnderlayTCP6() t.conn, err = t.connectUnderlayTCP6()
if err != nil { case LinkModeTCP4:
return err
}
} else { // hasipv4 or upnp
t.conn, err = t.connectUnderlayTCP() t.conn, err = t.connectUnderlayTCP()
case LinkModeTCPPunch:
t.conn, err = t.connectUnderlayTCP()
case LinkModeUDPPunch:
t.conn, err = t.connectUnderlayQuic()
}
if err != nil { if err != nil {
return err return err
} }
} if t.conn == nil {
return errors.New("connect underlay error")
} }
t.setRun(true) t.setRun(true)
go t.readLoop() go t.readLoop()
@@ -235,7 +261,7 @@ func (t *P2PTunnel) connectUnderlayQuic() (c underlay, err error) {
gLog.Println(LvINFO, "connectUnderlayQuic start") gLog.Println(LvINFO, "connectUnderlayQuic start")
defer gLog.Println(LvINFO, "connectUnderlayQuic end") defer gLog.Println(LvINFO, "connectUnderlayQuic end")
var qConn *underlayQUIC var qConn *underlayQUIC
if t.isUnderlayServer() { if t.config.isUnderlayServer == 1 {
time.Sleep(time.Millisecond * 10) // punching udp port will need some times in some env time.Sleep(time.Millisecond * 10) // punching udp port will need some times in some env
qConn, err = listenQuic(t.la.String(), TunnelIdleTimeout) qConn, err = listenQuic(t.la.String(), TunnelIdleTimeout)
if err != nil { if err != nil {
@@ -288,7 +314,7 @@ func (t *P2PTunnel) connectUnderlayQuic() (c underlay, err error) {
gLog.Println(LvINFO, "rtt=", time.Since(handshakeBegin)) gLog.Println(LvINFO, "rtt=", time.Since(handshakeBegin))
gLog.Println(LvDEBUG, "quic connection ok") gLog.Println(LvDEBUG, "quic connection ok")
t.linkMode = LinkModeUDPPunch t.linkModeWeb = LinkModeUDPPunch
return qConn, nil return qConn, nil
} }
@@ -297,37 +323,36 @@ func (t *P2PTunnel) connectUnderlayTCP() (c underlay, err error) {
gLog.Println(LvINFO, "connectUnderlayTCP start") gLog.Println(LvINFO, "connectUnderlayTCP start")
defer gLog.Println(LvINFO, "connectUnderlayTCP end") defer gLog.Println(LvINFO, "connectUnderlayTCP end")
var qConn *underlayTCP var qConn *underlayTCP
if t.isUnderlayServer() { if t.config.isUnderlayServer == 1 {
t.pn.push(t.config.PeerNode, MsgPushUnderlayConnect, nil) t.pn.push(t.config.PeerNode, MsgPushUnderlayConnect, nil)
qConn, err = listenTCP(t.coneNatPort, TunnelIdleTimeout) qConn, err = listenTCP(t.config.peerIP, t.config.peerConeNatPort, t.coneLocalPort, t.config.linkMode)
if err != nil { if err != nil {
return nil, fmt.Errorf("listen TCP error:%s", err) return nil, fmt.Errorf("listen TCP error:%s", err)
} }
_, buff, err := qConn.ReadBuffer() _, buff, err := qConn.ReadBuffer()
if err != nil { if err != nil {
qConn.listener.Close()
return nil, fmt.Errorf("read start msg error:%s", err) return nil, fmt.Errorf("read start msg error:%s", err)
} }
if buff != nil { if buff != nil {
gLog.Println(LvDEBUG, string(buff)) gLog.Println(LvDEBUG, string(buff))
} }
qConn.WriteBytes(MsgP2P, MsgTunnelHandshakeAck, []byte("OpenP2P,hello2")) qConn.WriteBytes(MsgP2P, MsgTunnelHandshakeAck, []byte("OpenP2P,hello2"))
gLog.Println(LvDEBUG, "TCP connection ok") gLog.Println(LvINFO, "TCP connection ok")
return qConn, nil return qConn, nil
} }
//else //else
t.pn.read(t.config.PeerNode, MsgPush, MsgPushUnderlayConnect, time.Second*5) t.pn.read(t.config.PeerNode, MsgPush, MsgPushUnderlayConnect, time.Second*5)
gLog.Println(LvDEBUG, "TCP dial to ", t.ra.String()) gLog.Println(LvDEBUG, "TCP dial to ", t.config.peerIP, ":", t.config.peerConeNatPort)
qConn, err = dialTCP(t.config.peerIP, t.config.peerConeNatPort) qConn, err = dialTCP(t.config.peerIP, t.config.peerConeNatPort, t.coneLocalPort, t.config.linkMode)
if err != nil { if err != nil {
return nil, fmt.Errorf("TCP dial to %s error:%s", t.ra.String(), err) return nil, fmt.Errorf("TCP dial to %s:%d error:%s", t.config.peerIP, t.config.peerConeNatPort, err)
} }
handshakeBegin := time.Now() handshakeBegin := time.Now()
qConn.WriteBytes(MsgP2P, MsgTunnelHandshake, []byte("OpenP2P,hello")) qConn.WriteBytes(MsgP2P, MsgTunnelHandshake, []byte("OpenP2P,hello"))
_, buff, err := qConn.ReadBuffer() _, buff, err := qConn.ReadBuffer()
if err != nil { if err != nil {
qConn.listener.Close()
return nil, fmt.Errorf("read MsgTunnelHandshake error:%s", err) return nil, fmt.Errorf("read MsgTunnelHandshake error:%s", err)
} }
if buff != nil { if buff != nil {
@@ -335,8 +360,8 @@ func (t *P2PTunnel) connectUnderlayTCP() (c underlay, err error) {
} }
gLog.Println(LvINFO, "rtt=", time.Since(handshakeBegin)) gLog.Println(LvINFO, "rtt=", time.Since(handshakeBegin))
gLog.Println(LvDEBUG, "TCP connection ok") gLog.Println(LvINFO, "TCP connection ok")
t.linkMode = LinkModeIPv4 t.linkModeWeb = LinkModeIPv4
return qConn, nil return qConn, nil
} }
@@ -344,7 +369,7 @@ func (t *P2PTunnel) connectUnderlayTCP6() (c underlay, err error) {
gLog.Println(LvINFO, "connectUnderlayTCP6 start") gLog.Println(LvINFO, "connectUnderlayTCP6 start")
defer gLog.Println(LvINFO, "connectUnderlayTCP6 end") defer gLog.Println(LvINFO, "connectUnderlayTCP6 end")
var qConn *underlayTCP6 var qConn *underlayTCP6
if t.isUnderlayServer() { if t.config.isUnderlayServer == 1 {
t.pn.push(t.config.PeerNode, MsgPushUnderlayConnect, nil) t.pn.push(t.config.PeerNode, MsgPushUnderlayConnect, nil)
qConn, err = listenTCP6(t.coneNatPort, TunnelIdleTimeout) qConn, err = listenTCP6(t.coneNatPort, TunnelIdleTimeout)
if err != nil { if err != nil {
@@ -365,10 +390,10 @@ func (t *P2PTunnel) connectUnderlayTCP6() (c underlay, err error) {
//else //else
t.pn.read(t.config.PeerNode, MsgPush, MsgPushUnderlayConnect, time.Second*5) t.pn.read(t.config.PeerNode, MsgPush, MsgPushUnderlayConnect, time.Second*5)
gLog.Println(LvDEBUG, "TCP6 dial to ", t.config.IPv6) gLog.Println(LvDEBUG, "TCP6 dial to ", t.config.peerIPv6)
qConn, err = dialTCP6(t.config.IPv6, t.config.peerConeNatPort) qConn, err = dialTCP6(t.config.peerIPv6, t.config.peerConeNatPort)
if err != nil { if err != nil {
return nil, fmt.Errorf("TCP6 dial to %s:%d error:%s", t.config.IPv6, 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()
qConn.WriteBytes(MsgP2P, MsgTunnelHandshake, []byte("OpenP2P,hello")) qConn.WriteBytes(MsgP2P, MsgTunnelHandshake, []byte("OpenP2P,hello"))
@@ -383,7 +408,7 @@ func (t *P2PTunnel) connectUnderlayTCP6() (c underlay, err error) {
gLog.Println(LvINFO, "rtt=", time.Since(handshakeBegin)) gLog.Println(LvINFO, "rtt=", time.Since(handshakeBegin))
gLog.Println(LvDEBUG, "TCP6 connection ok") gLog.Println(LvDEBUG, "TCP6 connection ok")
t.linkMode = LinkModeIPv6 t.linkModeWeb = LinkModeIPv6
return qConn, nil return qConn, nil
} }
@@ -569,7 +594,7 @@ func (t *P2PTunnel) listen() error {
// only private node set ipv6 // only private node set ipv6
if t.config.fromToken == t.pn.config.Token { if t.config.fromToken == t.pn.config.Token {
t.pn.refreshIPv6() t.pn.refreshIPv6()
rsp.IPv6 = t.pn.config.IPv6 rsp.IPv6 = t.pn.config.publicIPv6
} }
t.pn.push(t.config.PeerNode, MsgPushConnectRsp, rsp) t.pn.push(t.config.PeerNode, MsgPushConnectRsp, rsp)
@@ -594,14 +619,3 @@ func (t *P2PTunnel) closeOverlayConns(appID uint64) {
return true return true
}) })
} }
func (t *P2PTunnel) isUnderlayServer() bool {
if (t.pn.config.hasIPv4 == 1 || t.pn.config.hasUPNPorNATPMP == 1) && (t.config.hasIPv4 != 1 || t.config.hasUPNPorNATPMP != 1) {
return true
}
if (t.pn.config.hasIPv4 != 1 || t.pn.config.hasUPNPorNATPMP != 1) && (t.config.hasIPv4 == 1 || t.config.hasUPNPorNATPMP == 1) {
return false
}
// NAT or both has public IP
return t.tunnelServer
}
+53 -4
View File
@@ -10,9 +10,17 @@ import (
"time" "time"
) )
const OpenP2PVersion = "2.0.1" const OpenP2PVersion = "3.1.0"
const ProducnName string = "openp2p" const ProducnName string = "openp2p"
const LeastSupportTCPVersion = "1.5.0" const LeastSupportVersion = "3.0.0"
const (
IfconfigPort1 = 27180
IfconfigPort2 = 27181
WsPort = 27183
UDPPort1 = 27182
UDPPort2 = 27183
)
type openP2PHeader struct { type openP2PHeader struct {
DataLen uint32 DataLen uint32
@@ -68,6 +76,7 @@ const (
MsgP2P = 4 MsgP2P = 4
MsgRelay = 5 MsgRelay = 5
MsgReport = 6 MsgReport = 6
MsgQuery = 7
) )
// TODO: seperate node push and web push. // TODO: seperate node push and web push.
@@ -86,6 +95,7 @@ const (
MsgPushRestart = 11 MsgPushRestart = 11
MsgPushEditNode = 12 MsgPushEditNode = 12
MsgPushAPPKey = 13 MsgPushAPPKey = 13
MsgPushReportLog = 14
) )
// MsgP2P sub type message // MsgP2P sub type message
@@ -117,6 +127,7 @@ const (
MsgReportQuery MsgReportQuery
MsgReportConnect MsgReportConnect
MsgReportApps MsgReportApps
MsgReportLog
) )
const ( const (
@@ -158,8 +169,18 @@ const (
// linkmode // linkmode
const ( const (
LinkModeUDPPunch = "udppunch" LinkModeUDPPunch = "udppunch"
LinkModeIPv4 = "ipv4" LinkModeTCPPunch = "tcppunch"
LinkModeIPv6 = "ipv6" LinkModeIPv4 = "ipv4" // for web
LinkModeIPv6 = "ipv6" // for web
LinkModeTCP6 = "tcp6"
LinkModeTCP4 = "tcp4"
LinkModeUDP6 = "udp6"
LinkModeUDP4 = "udp4"
)
const (
MsgQueryPeerInfoReq = iota
MsgQueryPeerInfoRsp
) )
func newMessage(mainType uint16, subType uint16, packet interface{}) ([]byte, error) { func newMessage(mainType uint16, subType uint16, packet interface{}) ([]byte, error) {
@@ -199,6 +220,8 @@ type PushConnectReq struct {
FromIP string `json:"fromIP,omitempty"` FromIP string `json:"fromIP,omitempty"`
ID uint64 `json:"id,omitempty"` ID uint64 `json:"id,omitempty"`
AppKey uint64 `json:"appKey,omitempty"` // for underlay tcp AppKey uint64 `json:"appKey,omitempty"` // for underlay tcp
LinkMode string `json:"linkMode,omitempty"`
IsUnderlayServer int `json:"isServer,omitempty"` // Requset spec peer is server
} }
type PushConnectRsp struct { type PushConnectRsp struct {
Error int `json:"error,omitempty"` Error int `json:"error,omitempty"`
@@ -342,6 +365,18 @@ type ReportApps struct {
Apps []AppInfo Apps []AppInfo
} }
type ReportLogReq struct {
FileName string `json:"fileName,omitempty"`
Offset int64 `json:"offset,omitempty"`
Len int64 `json:"len,omitempty"`
}
type ReportLogRsp struct {
FileName string `json:"fileName,omitempty"`
Content string `json:"content,omitempty"`
Len int64 `json:"len,omitempty"`
Total int64 `json:"total,omitempty"`
}
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"`
@@ -380,3 +415,17 @@ type EditNode struct {
NewName string `json:"newName,omitempty"` NewName string `json:"newName,omitempty"`
Bandwidth int `json:"bandwidth,omitempty"` Bandwidth int `json:"bandwidth,omitempty"`
} }
type QueryPeerInfoReq struct {
Token uint64 `json:"token,omitempty"` // if public totp token
PeerNode string `json:"peerNode,omitempty"`
}
type QueryPeerInfoRsp struct {
Online int `json:"online,omitempty"`
Version string `json:"version,omitempty"`
NatType int `json:"natType,omitempty"`
IPv4 string `json:"IPv4,omitempty"`
HasIPv4 int `json:"hasIPv4,omitempty"` // has public ipv4
IPv6 string `json:"IPv6,omitempty"` // if public relay node, ipv6 not set
HasUPNPorNATPMP int `json:"hasUPNPorNATPMP,omitempty"`
}
-151
View File
@@ -1,151 +0,0 @@
package main
import (
"context"
"crypto/rand"
"crypto/rsa"
"crypto/tls"
"crypto/x509"
"encoding/json"
"encoding/pem"
"fmt"
"io"
"math/big"
"net"
"sync"
"time"
"github.com/lucas-clemente/quic-go"
)
//quic.DialContext do not support version 44,disable it
var quicVersion []quic.VersionNumber
type quicConn struct {
listener quic.Listener
writeMtx *sync.Mutex
quic.Stream
quic.Session
}
func (conn *quicConn) ReadMessage() (*openP2PHeader, []byte, error) {
headBuf := make([]byte, openP2PHeaderSize)
_, err := io.ReadFull(conn, headBuf)
if err != nil {
return nil, nil, err
}
head, err := decodeHeader(headBuf)
if err != nil {
return nil, nil, err
}
dataBuf := make([]byte, head.DataLen)
_, err = io.ReadFull(conn, dataBuf)
return head, dataBuf, err
}
func (conn *quicConn) WriteBytes(mainType uint16, subType uint16, data []byte) error {
writeBytes := append(encodeHeader(mainType, subType, uint32(len(data))), data...)
conn.writeMtx.Lock()
_, err := conn.Write(writeBytes)
conn.writeMtx.Unlock()
return err
}
func (conn *quicConn) WriteBuffer(data []byte) error {
conn.writeMtx.Lock()
_, err := conn.Write(data)
conn.writeMtx.Unlock()
return err
}
func (conn *quicConn) WriteMessage(mainType uint16, subType uint16, packet interface{}) error {
// TODO: call newMessage
data, err := json.Marshal(packet)
if err != nil {
return err
}
writeBytes := append(encodeHeader(mainType, subType, uint32(len(data))), data...)
conn.writeMtx.Lock()
_, err = conn.Write(writeBytes)
conn.writeMtx.Unlock()
return err
}
func (conn *quicConn) Close() error {
conn.Stream.CancelRead(1)
conn.Session.CloseWithError(0, "")
return nil
}
func (conn *quicConn) CloseListener() {
if conn.listener != nil {
conn.listener.Close()
}
}
func (conn *quicConn) Accept() error {
ctx, cancel := context.WithTimeout(context.Background(), time.Second*10)
defer cancel()
sess, err := conn.listener.Accept(ctx)
if err != nil {
return err
}
stream, err := sess.AcceptStream(context.Background())
if err != nil {
return err
}
conn.Stream = stream
conn.Session = sess
return nil
}
func listenQuic(addr string, idleTimeout time.Duration) (*quicConn, error) {
gLog.Println(LevelDEBUG, "quic listen on ", addr)
listener, err := quic.ListenAddr(addr, generateTLSConfig(),
&quic.Config{Versions: quicVersion, MaxIdleTimeout: idleTimeout, DisablePathMTUDiscovery: true})
if err != nil {
return nil, fmt.Errorf("quic.ListenAddr error:%s", err)
}
return &quicConn{listener: listener, writeMtx: &sync.Mutex{}}, nil
}
func dialQuic(conn *net.UDPConn, remoteAddr *net.UDPAddr, idleTimeout time.Duration) (*quicConn, error) {
tlsConf := &tls.Config{
InsecureSkipVerify: true,
NextProtos: []string{"openp2pv1"},
}
session, err := quic.DialContext(context.Background(), conn, remoteAddr, conn.LocalAddr().String(), tlsConf,
&quic.Config{Versions: quicVersion, MaxIdleTimeout: idleTimeout, DisablePathMTUDiscovery: true})
if err != nil {
return nil, fmt.Errorf("quic.DialContext error:%s", err)
}
stream, err := session.OpenStreamSync(context.Background())
if err != nil {
return nil, fmt.Errorf("OpenStreamSync error:%s", err)
}
qConn := &quicConn{nil, &sync.Mutex{}, stream, session}
return qConn, nil
}
// Setup a bare-bones TLS config for the server
func generateTLSConfig() *tls.Config {
key, err := rsa.GenerateKey(rand.Reader, 1024)
if err != nil {
panic(err)
}
template := x509.Certificate{SerialNumber: big.NewInt(1)}
certDER, err := x509.CreateCertificate(rand.Reader, &template, &template, &key.PublicKey, key)
if err != nil {
panic(err)
}
keyPEM := pem.EncodeToMemory(&pem.Block{Type: "RSA PRIVATE KEY", Bytes: x509.MarshalPKCS1PrivateKey(key)})
certPEM := pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: certDER})
tlsCert, err := tls.X509KeyPair(certPEM, keyPEM)
if err != nil {
panic(err)
}
return &tls.Config{
Certificates: []tls.Certificate{tlsCert},
NextProtos: []string{"openp2pv1"},
}
}
-35
View File
@@ -1,35 +0,0 @@
//go:build darwin
// +build darwin
package main
import (
"strings"
"syscall"
)
const (
defaultInstallPath = "/usr/local/openp2p"
defaultBinName = "openp2p"
)
func getOsName() (osName string) {
output := execOutput("sw_vers", "-productVersion")
osName = "Mac OS X " + strings.TrimSpace(output)
return
}
func setRLimit() error {
var limit syscall.Rlimit
if err := syscall.Getrlimit(syscall.RLIMIT_NOFILE, &limit); err != nil {
return err
}
limit.Cur = 10240
if err := syscall.Setrlimit(syscall.RLIMIT_NOFILE, &limit); err != nil {
return err
}
return nil
}
func setFirewall() {
}
-75
View File
@@ -1,75 +0,0 @@
//go:build linux
// +build linux
package main
import (
"bufio"
"bytes"
"io/ioutil"
"os"
"strings"
"syscall"
)
const (
defaultInstallPath = "/usr/local/openp2p"
defaultBinName = "openp2p"
)
func getOsName() (osName string) {
var sysnamePath string
sysnamePath = "/etc/redhat-release"
_, err := os.Stat(sysnamePath)
if err != nil && os.IsNotExist(err) {
str := "PRETTY_NAME="
f, err := os.Open("/etc/os-release")
if err != nil && os.IsNotExist(err) {
str = "DISTRIB_ID="
f, err = os.Open("/etc/openwrt_release")
}
if err == nil {
buf := bufio.NewReader(f)
for {
line, err := buf.ReadString('\n')
if err == nil {
line = strings.TrimSpace(line)
pos := strings.Count(line, str)
if pos > 0 {
len1 := len([]rune(str)) + 1
rs := []rune(line)
osName = string(rs[len1 : (len(rs))-1])
break
}
} else {
break
}
}
}
} else {
buff, err := ioutil.ReadFile(sysnamePath)
if err == nil {
osName = string(bytes.TrimSpace(buff))
}
}
if osName == "" {
osName = "Linux"
}
return
}
func setRLimit() error {
var limit syscall.Rlimit
if err := syscall.Getrlimit(syscall.RLIMIT_NOFILE, &limit); err != nil {
return err
}
limit.Max = 1024 * 1024
limit.Cur = limit.Max
if err := syscall.Setrlimit(syscall.RLIMIT_NOFILE, &limit); err != nil {
return err
}
return nil
}
func setFirewall() {
}
-56
View File
@@ -1,56 +0,0 @@
//go:build windows
// +build windows
package main
import (
"fmt"
"os"
"os/exec"
"path/filepath"
"strings"
"golang.org/x/sys/windows/registry"
)
const (
defaultInstallPath = "C:\\Program Files\\OpenP2P"
defaultBinName = "openp2p.exe"
)
func getOsName() (osName string) {
k, err := registry.OpenKey(registry.LOCAL_MACHINE, `SOFTWARE\Microsoft\Windows NT\CurrentVersion`, registry.QUERY_VALUE|registry.WOW64_64KEY)
if err != nil {
return
}
defer k.Close()
pn, _, err := k.GetStringValue("ProductName")
if err == nil {
osName = pn
}
return
}
func setRLimit() error {
return nil
}
func setFirewall() {
fullPath, err := filepath.Abs(os.Args[0])
if err != nil {
gLog.Println(LevelERROR, "add firewall error:", err)
return
}
isXP := false
osName := getOsName()
if strings.Contains(osName, "XP") || strings.Contains(osName, "2003") {
isXP = true
}
if isXP {
exec.Command("cmd.exe", `/c`, fmt.Sprintf(`netsh firewall del allowedprogram "%s"`, fullPath)).Run()
exec.Command("cmd.exe", `/c`, fmt.Sprintf(`netsh firewall add allowedprogram "%s" "%s" ENABLE`, ProducnName, fullPath)).Run()
} else { // win7 or later
exec.Command("cmd.exe", `/c`, fmt.Sprintf(`netsh advfirewall firewall del rule name="%s"`, ProducnName)).Run()
exec.Command("cmd.exe", `/c`, fmt.Sprintf(`netsh advfirewall firewall add rule name="%s" dir=in action=allow program="%s" enable=yes`, ProducnName, fullPath)).Run()
}
}
+23 -7
View File
@@ -7,10 +7,11 @@ import (
"net" "net"
"sync" "sync"
"time" "time"
reuse "github.com/openp2p-cn/go-reuseport"
) )
type underlayTCP struct { type underlayTCP struct {
listener net.Listener
writeMtx *sync.Mutex writeMtx *sync.Mutex
net.Conn net.Conn
} }
@@ -66,13 +67,20 @@ func (conn *underlayTCP) Close() error {
return conn.Conn.Close() return conn.Conn.Close()
} }
func listenTCP(port int, idleTimeout time.Duration) (*underlayTCP, error) { func listenTCP(host string, port int, localPort int, mode string) (*underlayTCP, error) {
addr, _ := net.ResolveTCPAddr("tcp", fmt.Sprintf("0.0.0.0:%d", port)) if mode == LinkModeTCPPunch {
c, err := reuse.DialTimeout("tcp", fmt.Sprintf("0.0.0.0:%d", localPort), fmt.Sprintf("%s:%d", host, port), SymmetricHandshakeAckTimeout) // TODO: timeout
if err != nil {
gLog.Println(LvDEBUG, "send tcp punch: ", err)
return nil, err
}
return &underlayTCP{writeMtx: &sync.Mutex{}, Conn: c}, nil
}
addr, _ := net.ResolveTCPAddr("tcp", fmt.Sprintf("0.0.0.0:%d", localPort))
l, err := net.ListenTCP("tcp", addr) l, err := net.ListenTCP("tcp", addr)
if err != nil { if err != nil {
return nil, err return nil, err
} }
defer l.Close()
l.SetDeadline(time.Now().Add(SymmetricHandshakeAckTimeout)) l.SetDeadline(time.Now().Add(SymmetricHandshakeAckTimeout))
c, err := l.Accept() c, err := l.Accept()
defer l.Close() defer l.Close()
@@ -82,11 +90,19 @@ func listenTCP(port int, idleTimeout time.Duration) (*underlayTCP, error) {
return &underlayTCP{writeMtx: &sync.Mutex{}, Conn: c}, nil return &underlayTCP{writeMtx: &sync.Mutex{}, Conn: c}, nil
} }
func dialTCP(host string, port int) (*underlayTCP, error) { func dialTCP(host string, port int, localPort int, mode string) (*underlayTCP, error) {
c, err := net.DialTimeout("tcp", fmt.Sprintf("%s:%d", host, port), SymmetricHandshakeAckTimeout) var c net.Conn
var err error
if mode == LinkModeTCPPunch {
c, err = reuse.DialTimeout("tcp", fmt.Sprintf("0.0.0.0:%d", localPort), fmt.Sprintf("%s:%d", host, port), SymmetricHandshakeAckTimeout)
} else {
c, err = net.DialTimeout("tcp", fmt.Sprintf("%s:%d", host, port), SymmetricHandshakeAckTimeout)
}
if err != nil { if err != nil {
fmt.Printf("Dial %s:%d error:%s", host, port, err) gLog.Printf(LvERROR, "Dial %s:%d error:%s", host, port, err)
return nil, err return nil, err
} }
gLog.Printf(LvDEBUG, "Dial %s:%d OK", host, port)
return &underlayTCP{writeMtx: &sync.Mutex{}, Conn: c}, nil return &underlayTCP{writeMtx: &sync.Mutex{}, Conn: c}, nil
} }
+2 -2
View File
@@ -16,7 +16,7 @@ import (
"time" "time"
) )
func update() { func update(host string, port int) {
gLog.Println(LvINFO, "update start") gLog.Println(LvINFO, "update start")
defer gLog.Println(LvINFO, "update end") defer gLog.Println(LvINFO, "update end")
c := http.Client{ c := http.Client{
@@ -27,7 +27,7 @@ func update() {
} }
goos := runtime.GOOS goos := runtime.GOOS
goarch := runtime.GOARCH goarch := runtime.GOARCH
rsp, err := c.Get(fmt.Sprintf("https://api.openp2p.cn:27183/api/v1/update?fromver=%s&os=%s&arch=%s", OpenP2PVersion, goos, goarch)) rsp, err := c.Get(fmt.Sprintf("https://%s:%d/api/v1/update?fromver=%s&os=%s&arch=%s", host, port, OpenP2PVersion, goos, goarch))
if err != nil { if err != nil {
gLog.Println(LvERROR, "update:query update list failed:", err) gLog.Println(LvERROR, "update:query update list failed:", err)
return return
+11 -30
View File
@@ -1,12 +1,8 @@
/* /*
Taken from taipei-torrent Taken from taipei-torrent And fix some bugs
Just enough UPnP to be able to forward ports
*/ */
package main package main
// BUG(jae): TODO: use syscalls to get actual ourIP. http://pastebin.com/9exZG4rh
import ( import (
"bytes" "bytes"
"encoding/xml" "encoding/xml"
@@ -34,11 +30,12 @@ type NAT interface {
} }
func Discover() (nat NAT, err error) { func Discover() (nat NAT, err error) {
localIP := localIPv4()
ssdp, err := net.ResolveUDPAddr("udp4", "239.255.255.250:1900") ssdp, err := net.ResolveUDPAddr("udp4", "239.255.255.250:1900")
if err != nil { if err != nil {
return return
} }
conn, err := net.ListenPacket("udp4", ":0") conn, err := net.ListenPacket("udp4", fmt.Sprintf("%s:0", localIP))
if err != nil { if err != nil {
return return
} }
@@ -67,6 +64,7 @@ func Discover() (nat NAT, err error) {
var n int var n int
_, _, err = socket.ReadFromUDP(answerBytes) _, _, err = socket.ReadFromUDP(answerBytes)
if err != nil { if err != nil {
gLog.Println(LvERROR, "UPNP discover error:", err)
return return
} }
@@ -98,12 +96,10 @@ func Discover() (nat NAT, err error) {
if err != nil { if err != nil {
return return
} }
var ourIP net.IP
ourIP, err = localIPv4()
if err != nil { if err != nil {
return return
} }
nat = &upnpNAT{serviceURL: serviceURL, ourIP: ourIP.String(), urnDomain: urnDomain} nat = &upnpNAT{serviceURL: serviceURL, ourIP: localIP, urnDomain: urnDomain}
return return
} }
} }
@@ -174,29 +170,14 @@ func getChildService(d *Device, serviceType string) *UPNPService {
return nil return nil
} }
func localIPv4() (net.IP, error) { func localIPv4() string { // TODO: multi nic will wrong
tt, err := net.Interfaces() conn, err := net.Dial("udp", "8.8.8.8:80")
if err != nil { if err != nil {
return nil, err return ""
} }
for _, t := range tt { defer conn.Close()
aa, err := t.Addrs() localAddr := conn.LocalAddr().(*net.UDPAddr)
if err != nil { return localAddr.IP.String()
return nil, err
}
for _, a := range aa {
ipnet, ok := a.(*net.IPNet)
if !ok {
continue
}
v4 := ipnet.IP.To4()
if v4 == nil || v4[0] == 127 { // loopback address
continue
}
return v4, nil
}
}
return nil, errors.New("cannot find local IP address")
} }
func getServiceURL(rootURL string) (url, urnDomain string, err error) { func getServiceURL(rootURL string) (url, urnDomain string, err error) {