Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
133fe046f8 | ||
|
|
95b46f51d0 | ||
|
|
7686af39e0 | ||
|
|
16b937ebd7 | ||
|
|
ac454ec694 | ||
|
|
029d69869f |
@@ -28,38 +28,43 @@ P2P直连可以让你的设备跑满带宽。不论你的设备在任何网络
|
|||||||
基于OpenP2P只需数行代码,就能让原来只能局域网通信的程序,变成任何内网都能通信
|
基于OpenP2P只需数行代码,就能让原来只能局域网通信的程序,变成任何内网都能通信
|
||||||
|
|
||||||
## 快速入门
|
## 快速入门
|
||||||
|
仅需简单4步就能用起来。
|
||||||
|
下面是一个远程办公例子:在家里连入办公室Windows电脑。
|
||||||
|
### 1.注册
|
||||||
|
前往<https://console.openp2p.cn> 注册新用户,暂无需任何认证
|
||||||
|
|
||||||
> :warning: 本文所有命令, Windows环境使用"openp2p.exe", Linux环境使用"./openp2p"
|

|
||||||
|
### 2.安装
|
||||||
|
分别在本地和远程电脑下载后双击运行,一键安装
|
||||||
|
|
||||||
|

|
||||||
|
|
||||||
以一个最常见的例子说明OpenP2P如何使用:远程办公,在家里连入办公室Windows电脑。
|
Windows默认会阻止没有花钱买它家证书签名过的程序,选择“仍要运行”即可。
|
||||||
相信很多人在疫情下远程办公是刚需。
|
|
||||||
1. 先确认办公室电脑已开启远程桌面功能(如何开启参考官方说明https://docs.microsoft.com/zh-cn/windows-server/remote/remote-desktop-services/clients/remote-desktop-allow-access)
|
|
||||||
2. 在办公室下载最新的`OpenP2P`[下载页](https://openp2p.cn/),解压出来,在命令行执行
|
|
||||||
```
|
|
||||||
openp2p.exe install -node OFFICEPC1 -user USERNAME1 -password PASSWORD1
|
|
||||||
```
|
|
||||||
|
|
||||||
> :warning: **切记将标记大写的参数改成自己的,3个参数的长度必须>=8个字符**
|
|
||||||
|
|
||||||

|

|
||||||
3. 在家里下载最新的OpenP2P,解压出来,在命令行执行
|
|
||||||
```
|

|
||||||
openp2p.exe -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 --peernode OFFICEPC1 --dstip 127.0.0.1 --dstport 3389 --srcport 23389 --protocol tcp
|
### 3.新建P2P应用
|
||||||
```
|
|
||||||
> :warning: **切记将标记大写的参数改成自己的**
|

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

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

|
||||||
|
|
||||||
|
### 4.使用P2P应用
|
||||||
|
在“MyHomePC”设备上能看到刚才创建的P2P应用,连接下图显示的“本地监听端口”即可。
|
||||||
|
|
||||||
|

|
||||||
|
|
||||||
|
在家里Windows电脑,按Win+R输入mstsc打开远程桌面,输入127.0.0.1:23389 /admin
|
||||||
|
|
||||||

|
|
||||||

|
|
||||||
`LISTEN ON PORT 23389 START` 看到这行日志表示P2PApp建立成功,监听23389端口。只需连接本机的127.0.0.1:23389就相当于连接公司Windows电脑的3389端口。
|
|
||||||
|
|
||||||
4. 在家里Windows电脑,按Win+R输入mstsc打开远程桌面,输入127.0.0.1:23389 /admin
|
|
||||||

|

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

|

|
||||||
|
|
||||||
## 详细使用说明
|
## 详细使用说明
|
||||||
[这里](/USAGE-ZH.md)详细介绍如何使用和运行参数
|
[这里](/USAGE-ZH.md)介绍如何手动运行
|
||||||
|
|
||||||
## 典型应用场景
|
## 典型应用场景
|
||||||
特别适合大流量的内网访问
|
特别适合大流量的内网访问
|
||||||
|
|||||||
@@ -30,52 +30,45 @@ P2P direct connection lets your devices make good use of bandwidth. Your device
|
|||||||
Your applicaiton can call OpenP2P with a few code to make any internal networks communicate with each other.
|
Your applicaiton can call OpenP2P with a few code to make any internal networks communicate with each other.
|
||||||
|
|
||||||
## Get Started
|
## Get Started
|
||||||
A common scenario to introduce OpenP2P: remote work. At home connects to office's Linux PC .
|
Just 4 simple steps to use.
|
||||||
Under the outbreak of covid-19 pandemic, surely remote work becomes a fundamental demand.
|
Here's an example of remote work: connecting to an office Windows computer at home.
|
||||||
|
|
||||||
|
### 1.Register
|
||||||
|
Go to <https://console.openp2p.cn> register a new user
|
||||||
|
|
||||||
> :warning: all commands in this doc, Windows env uses "openp2p.exe", Linux env uses "./openp2p"
|

|
||||||
|
### 2.Install
|
||||||
|
Download on local and remote computers and double-click to run, one-click installation
|
||||||
|
|
||||||
1. Make sure your office device(Linux) has opened the access of ssh.
|

|
||||||
```
|
|
||||||
netstat -nl | grep 22
|
|
||||||
```
|
|
||||||
Output sample
|
|
||||||

|
|
||||||
|
|
||||||
2. Download the latest version of `OpenP2P` [Download Page](https://openp2p.cn/),unzip the downloaded package, and execute below command line.
|
By default, Windows will block programs that have not been signed by the Microsoft's certificate, and you can select "Run anyway".
|
||||||
```
|
|
||||||
tar xzvf ${PackageName}
|
|
||||||
./openp2p install -node OFFICEPC1 -user USERNAME1 -password PASSWORD1
|
|
||||||
```
|
|
||||||
|
|
||||||
> :warning: **Must change the parameters marked in UPPERCASE to your own. These 3 parameters must >= 8 charaters**
|

|
||||||
|
|
||||||
Output sample
|

|
||||||

|
### 3.New P2PApp
|
||||||
|
|
||||||
3. Download OpenP2P on your home device,unzip and execute below command line.
|

|
||||||
```
|
|
||||||
openp2p.exe -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 --peernode OFFICEPC1 --dstip 127.0.0.1 --dstport 22 --srcport 22022 --protocol tcp
|
|
||||||
```
|
|
||||||
|
|
||||||
> :warning: **Must change the parameters marked in UPPERCASE to your own**
|
|
||||||
|
|
||||||
Output sample
|

|
||||||

|
|
||||||
The log of `LISTEN ON PORT 22022 START` indicates P2PApp runs successfully on your home device, listing port is 22022. Once connects to local ip:port,127.0.0.1:22022, it means the home device has conneccted to the office device's port, 22.
|
|
||||||

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

|
||||||
|
|
||||||
4. Test the connection between office device and home device.In your home deivce, run SSH to login the office device.
|
### 4.Use P2PApp
|
||||||
```
|
You can see the P2P application you just created on the "MyHomePC" device, just connect to the "local listening port" shown in the figure below.
|
||||||
ssh -p22022 [email protected]:22022
|
|
||||||
```
|

|
||||||

|
|
||||||
|
On MyHomePC, press Win+R and enter MSTSC to open the remote desktop, input `127.0.0.1:23389 /admin`
|
||||||
|
|
||||||
|

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

|
||||||
|
|
||||||
|
|
||||||
## Usage
|
## Usage
|
||||||
[Here](/USAGE.md) is a detailed description of how to use and running parameters
|
[Here](/USAGE.md) describes how to run manually
|
||||||
|
|
||||||
## Scenarios
|
## Scenarios
|
||||||
Especially suitable for large traffic intranet access.
|
Especially suitable for large traffic intranet access.
|
||||||
|
|||||||
@@ -1,37 +1,70 @@
|
|||||||
# 详细运行参数说明
|
# 手动运行说明
|
||||||
|
大部分情况通过<https://console.openp2p.cn> 操作即可。有些情况需要手动运行
|
||||||
> :warning: 本文所有命令, Windows环境使用"openp2p.exe", Linux环境使用"./openp2p"
|
> :warning: 本文所有命令, Windows环境使用"openp2p.exe", Linux环境使用"./openp2p"
|
||||||
|
|
||||||
|
|
||||||
## 安装和监听
|
## 安装和监听
|
||||||
```
|
```
|
||||||
./openp2p install -node OFFICEPC1 -user USERNAME1 -password PASSWORD1
|
./openp2p install -node OFFICEPC1 -token TOKEN
|
||||||
或
|
或
|
||||||
./openp2p -d -node OFFICEPC1 -user USERNAME1 -password PASSWORD1
|
./openp2p -d -node OFFICEPC1 -token TOKEN
|
||||||
# 注意Windows系统把“./openp2p” 换成“openp2p.exe”
|
# 注意Windows系统把“./openp2p” 换成“openp2p.exe”
|
||||||
```
|
```
|
||||||
>* install: 安装模式【推荐】,会安装成系统服务,这样它就能随系统自动启动
|
>* install: 安装模式【推荐】,会安装成系统服务,这样它就能随系统自动启动
|
||||||
>* -d: daemon模式。发现worker进程意外退出就会自动启动新的worker进程
|
>* -d: daemon模式。发现worker进程意外退出就会自动启动新的worker进程
|
||||||
>* -node: 独一无二的节点名字,唯一标识
|
>* -node: 独一无二的节点名字,唯一标识
|
||||||
>* -user: 独一无二的用户名字,该节点属于这个user
|
>* -token: 在<console.openp2p.cn>“我的”里面找到
|
||||||
>* -password: 密码
|
>* -sharebandwidth: 作为共享节点时提供带宽,默认10mbps. 如果是光纤大带宽,设置越大效果越好. -1表示不共享,该节点只在私有的P2P网络使用。不加入共享的P2P网络,这样也意味着无法使用别人的共享节点
|
||||||
>* -sharebandwidth: 作为共享节点时提供带宽,默认10mbps. 如果是光纤大带宽,设置越大效果越好
|
|
||||||
>* -loglevel: 需要查看更多调试日志,设置0;默认是1
|
>* -loglevel: 需要查看更多调试日志,设置0;默认是1
|
||||||
>* -noshare: 不共享,该节点只在私有的P2P网络使用。不加入共享的P2P网络,这样也意味着无法使用别人的共享节点
|
|
||||||
|
|
||||||
## 连接
|
## 连接
|
||||||
```
|
```
|
||||||
./openp2p -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389 -protocol tcp
|
./openp2p -d -node HOMEPC123 -token TOKEN -appname OfficeWindowsRemote -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389
|
||||||
使用配置文件,建立多个P2PApp
|
使用配置文件,建立多个P2PApp
|
||||||
./openp2p -d -f
|
./openp2p -d
|
||||||
./openp2p -f
|
|
||||||
```
|
```
|
||||||
|
>* -appname: 这个P2P应用名字
|
||||||
>* -peernode: 目标节点名字
|
>* -peernode: 目标节点名字
|
||||||
>* -dstip: 目标服务地址,默认本机127.0.0.1
|
>* -dstip: 目标服务地址,默认本机127.0.0.1
|
||||||
>* -dstport: 目标服务端口,常见的如windows远程桌面3389,Linux ssh 22
|
>* -dstport: 目标服务端口,常见的如windows远程桌面3389,Linux ssh 22
|
||||||
>* -protocol: 目标服务协议 tcp、udp
|
>* -protocol: 目标服务协议 tcp、udp
|
||||||
>* -peeruser: 目标用户,如果是同一个用户下的节点,则无需设置
|
|
||||||
>* -peerpassword: 目标密码,如果是同一个用户下的节点,则无需设置
|
## 配置文件
|
||||||
>* -f: 配置文件,如果希望配置多个P2PApp参考[config.json](/config.json)
|
一般保存在当前目录,安装模式下会保存到 `C:\Program Files\OpenP2P\config.json` 或 `/usr/local/openp2p/config.json`
|
||||||
|
希望修改参数,或者配置多个P2PApp可手动修改配置文件
|
||||||
|
|
||||||
|
配置实例
|
||||||
|
```
|
||||||
|
{
|
||||||
|
"network": {
|
||||||
|
"Node": "hhd1207-222",
|
||||||
|
"Token": "TOKEN",
|
||||||
|
"ShareBandwidth": -1,
|
||||||
|
"ServerHost": "api.openp2p.cn",
|
||||||
|
"ServerPort": 27183,
|
||||||
|
"UDPPort1": 27182,
|
||||||
|
"UDPPort2": 27183
|
||||||
|
},
|
||||||
|
"apps": [
|
||||||
|
{
|
||||||
|
"AppName": "OfficeWindowsPC",
|
||||||
|
"Protocol": "tcp",
|
||||||
|
"SrcPort": 23389,
|
||||||
|
"PeerNode": "OFFICEPC1",
|
||||||
|
"DstPort": 3389,
|
||||||
|
"DstHost": "localhost",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"AppName": "OfficeServerSSH",
|
||||||
|
"Protocol": "tcp",
|
||||||
|
"SrcPort": 22,
|
||||||
|
"PeerNode": "OFFICEPC1",
|
||||||
|
"DstPort": 22,
|
||||||
|
"DstHost": "192.168.1.5",
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
## 升级客户端
|
## 升级客户端
|
||||||
```
|
```
|
||||||
|
|||||||
@@ -1,39 +1,72 @@
|
|||||||
# Parameters details
|
|
||||||
|
|
||||||
|
|
||||||
|
# Parameters details
|
||||||
|
In most cases, you can operate it through <https://console.openp2p.cn>. In some cases it is necessary to run manually
|
||||||
> :warning: all commands in this doc, Windows env uses "openp2p.exe", Linux env uses "./openp2p"
|
> :warning: all commands in this doc, Windows env uses "openp2p.exe", Linux env uses "./openp2p"
|
||||||
|
|
||||||
|
|
||||||
## Install and Listen
|
## Install and Listen
|
||||||
```
|
```
|
||||||
./openp2p install -node OFFICEPC1 -user USERNAME1 -password PASSWORD1
|
./openp2p install -node OFFICEPC1 -token TOKEN
|
||||||
Or
|
Or
|
||||||
./openp2p -d -node OFFICEPC1 -user USERNAME1 -password PASSWORD1
|
./openp2p -d -node OFFICEPC1 -token TOKEN
|
||||||
|
|
||||||
```
|
```
|
||||||
>* install: [recommand] will install as system service. So it will autorun when system booting.
|
>* install: [recommand] will install as system service. So it will autorun when system booting.
|
||||||
>* -d: daemon mode run once. When the worker process is found to exit unexpectedly, a new worker process will be automatically started
|
>* -d: daemon mode run once. When the worker process is found to exit unexpectedly, a new worker process will be automatically started
|
||||||
>* -node: Unique node name, unique identification
|
>* -node: Unique node name, unique identification
|
||||||
>* -user: Unique user name, the node belongs to this user
|
>* -token: See <console.openp2p.cn> "Profile"
|
||||||
>* -password: Password
|
>* -sharebandwidth: Provides bandwidth when used as a shared node, the default is 10mbps. If it is a large bandwidth of optical fiber, the larger the setting, the better the effect. -1 means not shared, the node is only used in a private P2P network. Do not join the shared P2P network, which also means that you CAN NOT use other people’s shared nodes
|
||||||
>* -sharebandwidth: Provides bandwidth when used as a shared node, the default is 10mbps. If it is a large bandwidth of optical fiber, the larger the setting, the better the effect
|
|
||||||
>* -loglevel: Need to view more debug logs, set 0; the default is 1
|
>* -loglevel: Need to view more debug logs, set 0; the default is 1
|
||||||
>* -noshare: Not shared, the node is only used in a private P2P network. Do not join the shared P2P network, which also means that you CAN NOT use other people’s shared nodes
|
|
||||||
|
|
||||||
## Connect
|
## Connect
|
||||||
```
|
```
|
||||||
./openp2p -d -node HOMEPC123 -user USERNAME1 -password PASSWORD1 -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389 -protocol tcp
|
./openp2p -d -node HOMEPC123 -token TOKEN -appname OfficeWindowsRemote -peernode OFFICEPC1 -dstip 127.0.0.1 -dstport 3389 -srcport 23389
|
||||||
Create multiple P2PApp by config file
|
Create multiple P2PApp by config file
|
||||||
./openp2p -d -f
|
./openp2p -d
|
||||||
./openp2p -f
|
|
||||||
```
|
```
|
||||||
|
>* -appname: This P2PApp name
|
||||||
>* -peernode: Target node name
|
>* -peernode: Target node name
|
||||||
>* -dstip: Target service address, default local 127.0.0.1
|
>* -dstip: Target service address, default local 127.0.0.1
|
||||||
>* -dstport: Target service port, such as windows remote desktop 3389, Linux ssh 22
|
>* -dstport: Target service port, such as windows remote desktop 3389, Linux ssh 22
|
||||||
>* -protocol: Target service protocol tcp, udp
|
>* -protocol: Target service protocol tcp, udp
|
||||||
>* -peeruser: The target user, if it is a node under the same user, no need to set
|
|
||||||
>* -peerpassword: The target password, if it is a node under the same user, no need to set
|
|
||||||
>* -f: Configuration file, if you want to configure multiple P2PApp refer to [config.json](/config.json)
|
|
||||||
|
|
||||||
|
## Config file
|
||||||
|
Generally saved in the current directory, in installation mode it will be saved to `C:\Program Files\OpenP2P\config.json` or `/usr/local/openp2p/config.json`
|
||||||
|
If you want to modify the parameters, or configure multiple P2PApps, you can manually modify the configuration file
|
||||||
|
|
||||||
|
Configuration example
|
||||||
|
```
|
||||||
|
{
|
||||||
|
"network": {
|
||||||
|
"Node": "hhd1207-222",
|
||||||
|
"Token": "TOKEN",
|
||||||
|
"ShareBandwidth": -1,
|
||||||
|
"ServerHost": "api.openp2p.cn",
|
||||||
|
"ServerPort": 27183,
|
||||||
|
"UDPPort1": 27182,
|
||||||
|
"UDPPort2": 27183
|
||||||
|
},
|
||||||
|
"apps": [
|
||||||
|
{
|
||||||
|
"AppName": "OfficeWindowsPC",
|
||||||
|
"Protocol": "tcp",
|
||||||
|
"SrcPort": 23389,
|
||||||
|
"PeerNode": "OFFICEPC1",
|
||||||
|
"DstPort": 3389,
|
||||||
|
"DstHost": "localhost",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"AppName": "OfficeServerSSH",
|
||||||
|
"Protocol": "tcp",
|
||||||
|
"SrcPort": 22,
|
||||||
|
"PeerNode": "OFFICEPC1",
|
||||||
|
"DstPort": 22,
|
||||||
|
"DstHost": "192.168.1.5",
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
```
|
||||||
## Client update
|
## Client update
|
||||||
```
|
```
|
||||||
# update local client
|
# update local client
|
||||||
|
|||||||
@@ -7,39 +7,39 @@ import (
|
|||||||
|
|
||||||
// BandwidthLimiter ...
|
// BandwidthLimiter ...
|
||||||
type BandwidthLimiter struct {
|
type BandwidthLimiter struct {
|
||||||
freeFlowTime time.Time
|
ts time.Time
|
||||||
bandwidth int // mbps
|
bw int // mbps
|
||||||
freeFlow int // bytes
|
freeBytes int // bytes
|
||||||
maxFreeFlow int // bytes
|
maxFreeBytes int // bytes
|
||||||
freeFlowMtx sync.Mutex
|
mtx sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
// mbps
|
// mbps
|
||||||
func newBandwidthLimiter(bw int) *BandwidthLimiter {
|
func newBandwidthLimiter(bw int) *BandwidthLimiter {
|
||||||
return &BandwidthLimiter{
|
return &BandwidthLimiter{
|
||||||
bandwidth: bw,
|
bw: bw,
|
||||||
freeFlowTime: time.Now(),
|
ts: time.Now(),
|
||||||
maxFreeFlow: bw * 1024 * 1024 / 8,
|
maxFreeBytes: bw * 1024 * 1024 / 8,
|
||||||
freeFlow: bw * 1024 * 1024 / 8,
|
freeBytes: bw * 1024 * 1024 / 8,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Add ...
|
// Add ...
|
||||||
func (bl *BandwidthLimiter) Add(bytes int) {
|
func (bl *BandwidthLimiter) Add(bytes int) {
|
||||||
if bl.bandwidth <= 0 {
|
if bl.bw <= 0 {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
bl.freeFlowMtx.Lock()
|
bl.mtx.Lock()
|
||||||
defer bl.freeFlowMtx.Unlock()
|
defer bl.mtx.Unlock()
|
||||||
// calc free flow 1000*1000/1024/1024=0.954; 1024*1024/1000/1000=1.048
|
// calc free flow 1000*1000/1024/1024=0.954; 1024*1024/1000/1000=1.048
|
||||||
bl.freeFlow += int(time.Now().Sub(bl.freeFlowTime) * time.Duration(bl.bandwidth) / 8 / 954)
|
bl.freeBytes += int(time.Since(bl.ts) * time.Duration(bl.bw) / 8 / 954)
|
||||||
if bl.freeFlow > bl.maxFreeFlow {
|
if bl.freeBytes > bl.maxFreeBytes {
|
||||||
bl.freeFlow = bl.maxFreeFlow
|
bl.freeBytes = bl.maxFreeBytes
|
||||||
}
|
}
|
||||||
bl.freeFlow -= bytes
|
bl.freeBytes -= bytes
|
||||||
bl.freeFlowTime = time.Now()
|
bl.ts = time.Now()
|
||||||
if bl.freeFlow < 0 {
|
if bl.freeBytes < 0 {
|
||||||
// sleep for the overflow
|
// sleep for the overflow
|
||||||
time.Sleep(time.Millisecond * time.Duration(-bl.freeFlow/(bl.bandwidth*1048/8)))
|
time.Sleep(time.Millisecond * time.Duration(-bl.freeBytes/(bl.bw*1048/8)))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2,21 +2,27 @@ package main
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"flag"
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
|
"os"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
var gConf Config
|
var gConf Config
|
||||||
|
|
||||||
|
const IntValueNotSet int = -99999999
|
||||||
|
|
||||||
type AppConfig struct {
|
type AppConfig struct {
|
||||||
// required
|
// required
|
||||||
Protocol string
|
AppName string
|
||||||
SrcPort int
|
Protocol string
|
||||||
PeerNode string
|
SrcPort int
|
||||||
DstPort int
|
PeerNode string
|
||||||
DstHost string
|
DstPort int
|
||||||
PeerUser string
|
DstHost string
|
||||||
PeerPassword string
|
PeerUser string
|
||||||
|
Enabled int // default:1
|
||||||
// runtime info
|
// runtime info
|
||||||
peerToken uint64
|
peerToken uint64
|
||||||
peerNatType int
|
peerNatType int
|
||||||
@@ -29,25 +35,60 @@ type AppConfig struct {
|
|||||||
|
|
||||||
// TODO: add loglevel, maxlogfilesize
|
// TODO: add loglevel, maxlogfilesize
|
||||||
type Config struct {
|
type Config struct {
|
||||||
Network NetworkConfig `json:"network"`
|
Network NetworkConfig `json:"network"`
|
||||||
Apps []AppConfig `json:"apps"`
|
Apps []AppConfig `json:"apps"`
|
||||||
daemonMode bool
|
LogLevel int
|
||||||
|
|
||||||
|
mtx sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) add(app AppConfig) {
|
func (c *Config) switchApp(app AppConfig, enabled int) {
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
|
for i := 0; i < len(c.Apps); i++ {
|
||||||
|
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
||||||
|
c.Apps[i].Enabled = enabled
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *Config) add(app AppConfig, force bool) {
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
if app.SrcPort == 0 || app.DstPort == 0 {
|
if app.SrcPort == 0 || app.DstPort == 0 {
|
||||||
|
gLog.Println(LevelERROR, "invalid app ", app)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
for i := 0; i < len(c.Apps); i++ {
|
for i := 0; i < len(c.Apps); i++ {
|
||||||
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
||||||
|
if force {
|
||||||
|
c.Apps[i] = app
|
||||||
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
c.Apps = append(c.Apps, app)
|
c.Apps = append(c.Apps, app)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c *Config) delete(app AppConfig) {
|
||||||
|
if app.SrcPort == 0 || app.DstPort == 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
|
for i := 0; i < len(c.Apps); i++ {
|
||||||
|
if c.Apps[i].Protocol == app.Protocol && c.Apps[i].SrcPort == app.SrcPort {
|
||||||
|
c.Apps = append(c.Apps[:i], c.Apps[i+1:]...)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (c *Config) save() {
|
func (c *Config) save() {
|
||||||
data, _ := json.MarshalIndent(c, "", "")
|
c.mtx.Lock()
|
||||||
|
defer c.mtx.Unlock()
|
||||||
|
data, _ := json.MarshalIndent(c, "", " ")
|
||||||
err := ioutil.WriteFile("config.json", data, 0644)
|
err := ioutil.WriteFile("config.json", data, 0644)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "save config.json error:", err)
|
gLog.Println(LevelERROR, "save config.json error:", err)
|
||||||
@@ -55,9 +96,13 @@ func (c *Config) save() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (c *Config) load() error {
|
func (c *Config) load() error {
|
||||||
|
c.mtx.Lock()
|
||||||
|
c.LogLevel = IntValueNotSet
|
||||||
|
c.Network.ShareBandwidth = IntValueNotSet
|
||||||
|
defer c.mtx.Unlock()
|
||||||
data, err := ioutil.ReadFile("config.json")
|
data, err := ioutil.ReadFile("config.json")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "read config.json error:", err)
|
// gLog.Println(LevelERROR, "read config.json error:", err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
err = json.Unmarshal(data, &c)
|
err = json.Unmarshal(data, &c)
|
||||||
@@ -69,10 +114,9 @@ func (c *Config) load() error {
|
|||||||
|
|
||||||
type NetworkConfig struct {
|
type NetworkConfig struct {
|
||||||
// local info
|
// local info
|
||||||
|
Token uint64
|
||||||
Node string
|
Node string
|
||||||
User string
|
User string
|
||||||
Password string
|
|
||||||
NoShare bool
|
|
||||||
localIP string
|
localIP string
|
||||||
ipv6 string
|
ipv6 string
|
||||||
hostName string
|
hostName string
|
||||||
@@ -87,3 +131,80 @@ type NetworkConfig struct {
|
|||||||
UDPPort1 int
|
UDPPort1 int
|
||||||
UDPPort2 int
|
UDPPort2 int
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func parseParams() {
|
||||||
|
serverHost := flag.String("serverhost", "api.openp2p.cn", "server host ")
|
||||||
|
// serverHost := flag.String("serverhost", "127.0.0.1", "server host ") // for debug
|
||||||
|
node := flag.String("node", "", "node name. 8-31 characters")
|
||||||
|
token := flag.Uint64("token", 0, "token")
|
||||||
|
peerNode := flag.String("peernode", "", "peer node name that you want to connect")
|
||||||
|
dstIP := flag.String("dstip", "127.0.0.1", "destination ip ")
|
||||||
|
dstPort := flag.Int("dstport", 0, "destination port ")
|
||||||
|
srcPort := flag.Int("srcport", 0, "source port ")
|
||||||
|
protocol := flag.String("protocol", "tcp", "tcp or udp")
|
||||||
|
appName := flag.String("appname", "", "app name")
|
||||||
|
flag.Bool("noshare", false, "deprecated. uses -sharebandwidth -1") // Deprecated, rm later
|
||||||
|
shareBandwidth := flag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private node no limit")
|
||||||
|
flag.Bool("f", false, "deprecated. config file") // Deprecated, rm later
|
||||||
|
daemonMode := flag.Bool("d", false, "daemonMode")
|
||||||
|
flag.Bool("bydaemon", false, "start by daemon") // Deprecated, rm later
|
||||||
|
logLevel := flag.Int("loglevel", 1, "0:debug 1:info 2:warn 3:error")
|
||||||
|
flag.Parse()
|
||||||
|
|
||||||
|
config := AppConfig{Enabled: 1}
|
||||||
|
config.PeerNode = *peerNode
|
||||||
|
config.DstHost = *dstIP
|
||||||
|
config.DstPort = *dstPort
|
||||||
|
config.SrcPort = *srcPort
|
||||||
|
config.Protocol = *protocol
|
||||||
|
config.AppName = *appName
|
||||||
|
gConf.load()
|
||||||
|
if config.SrcPort != 0 {
|
||||||
|
gConf.add(config, true)
|
||||||
|
}
|
||||||
|
gConf.mtx.Lock()
|
||||||
|
|
||||||
|
// spec paramters in commandline will always be used
|
||||||
|
flag.Visit(func(f *flag.Flag) {
|
||||||
|
if f.Name == "sharebandwidth" {
|
||||||
|
gConf.Network.ShareBandwidth = *shareBandwidth
|
||||||
|
}
|
||||||
|
if f.Name == "node" {
|
||||||
|
gConf.Network.Node = *node
|
||||||
|
}
|
||||||
|
if f.Name == "serverhost" {
|
||||||
|
gConf.Network.ServerHost = *serverHost
|
||||||
|
}
|
||||||
|
if f.Name == "loglevel" {
|
||||||
|
gConf.LogLevel = *logLevel
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
if gConf.Network.ServerHost == "" {
|
||||||
|
gConf.Network.ServerHost = *serverHost
|
||||||
|
}
|
||||||
|
if gConf.Network.Node == "" {
|
||||||
|
gConf.Network.Node = *node
|
||||||
|
}
|
||||||
|
if *token != 0 {
|
||||||
|
gConf.Network.Token = *token
|
||||||
|
}
|
||||||
|
if gConf.LogLevel == IntValueNotSet {
|
||||||
|
gConf.LogLevel = *logLevel
|
||||||
|
}
|
||||||
|
if gConf.Network.ShareBandwidth == IntValueNotSet {
|
||||||
|
gConf.Network.ShareBandwidth = *shareBandwidth
|
||||||
|
}
|
||||||
|
|
||||||
|
gConf.Network.ServerPort = 27183
|
||||||
|
gConf.Network.UDPPort1 = 27182
|
||||||
|
gConf.Network.UDPPort2 = 27183
|
||||||
|
gLog.setLevel(LogLevel(gConf.LogLevel))
|
||||||
|
gConf.mtx.Unlock()
|
||||||
|
gConf.save()
|
||||||
|
if *daemonMode {
|
||||||
|
d := daemon{}
|
||||||
|
d.run()
|
||||||
|
os.Exit(0)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -5,7 +5,9 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"os"
|
"os"
|
||||||
|
"os/exec"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/kardianos/service"
|
"github.com/kardianos/service"
|
||||||
@@ -38,7 +40,6 @@ func (d *daemon) Stop(s service.Service) error {
|
|||||||
func (d *daemon) run() {
|
func (d *daemon) run() {
|
||||||
gLog.Println(LevelINFO, "daemon run start")
|
gLog.Println(LevelINFO, "daemon run start")
|
||||||
defer gLog.Println(LevelINFO, "daemon run end")
|
defer gLog.Println(LevelINFO, "daemon run end")
|
||||||
os.Chdir(filepath.Dir(os.Args[0])) // for system service
|
|
||||||
d.running = true
|
d.running = true
|
||||||
binPath, _ := os.Executable()
|
binPath, _ := os.Executable()
|
||||||
mydir, err := os.Getwd()
|
mydir, err := os.Getwd()
|
||||||
@@ -63,10 +64,9 @@ func (d *daemon) run() {
|
|||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
args = append(args, "-bydaemon")
|
|
||||||
for {
|
for {
|
||||||
// start worker
|
// start worker
|
||||||
gLog.Println(LevelINFO, "start worker process")
|
gLog.Println(LevelINFO, "start worker process, args:", args)
|
||||||
execSpec := &os.ProcAttr{Files: []*os.File{os.Stdin, os.Stdout, os.Stderr}}
|
execSpec := &os.ProcAttr{Files: []*os.File{os.Stdin, os.Stdout, os.Stderr}}
|
||||||
p, err := os.StartProcess(binPath, args, execSpec)
|
p, err := os.StartProcess(binPath, args, execSpec)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -106,58 +106,73 @@ func (d *daemon) Control(ctrlComm string, exeAbsPath string, args []string) erro
|
|||||||
|
|
||||||
// examples:
|
// examples:
|
||||||
// listen:
|
// listen:
|
||||||
// ./openp2p install -node hhd1207-222 -user tenderiron -password 13760636579 -noshare
|
// ./openp2p install -node hhd1207-222 -token YOUR-TOKEN -sharebandwidth -1
|
||||||
// listen and build p2papp:
|
// listen and build p2papp:
|
||||||
// ./openp2p install -node hhd1207-222 -user tenderiron -password 13760636579 -noshare -peernode hhdhome-n1 -dstip 127.0.0.1 -dstport 50022 -protocol tcp -srcport 22
|
// ./openp2p install -node hhd1207-222 -token YOUR-TOKEN -sharebandwidth -1 -peernode hhdhome-n1 -dstip 127.0.0.1 -dstport 50022 -protocol tcp -srcport 22
|
||||||
func install() {
|
func install() {
|
||||||
gLog = InitLogger(filepath.Dir(os.Args[0]), "openp2p-install", LevelDEBUG, 1024*1024, LogConsole)
|
gLog.Println(LevelINFO, "install start")
|
||||||
|
defer gLog.Println(LevelINFO, "install end")
|
||||||
|
// auto uninstall
|
||||||
|
|
||||||
|
uninstall()
|
||||||
// save config file
|
// save config file
|
||||||
installFlag := flag.NewFlagSet("install", flag.ExitOnError)
|
installFlag := flag.NewFlagSet("install", flag.ExitOnError)
|
||||||
serverHost := installFlag.String("serverhost", "api.openp2p.cn", "server host ")
|
serverHost := installFlag.String("serverhost", "api.openp2p.cn", "server host ")
|
||||||
// serverHost := flag.String("serverhost", "127.0.0.1", "server host ") // for debug
|
// serverHost := flag.String("serverhost", "127.0.0.1", "server host ") // for debug
|
||||||
user := installFlag.String("user", "", "user name. 8-31 characters")
|
token := installFlag.Uint64("token", 0, "token")
|
||||||
node := installFlag.String("node", "", "node name. 8-31 characters")
|
node := installFlag.String("node", "", "node name. 8-31 characters. if not set, it will be hostname")
|
||||||
password := installFlag.String("password", "", "user password. 8-31 characters")
|
|
||||||
peerNode := installFlag.String("peernode", "", "peer node name that you want to connect")
|
peerNode := installFlag.String("peernode", "", "peer node name that you want to connect")
|
||||||
peerUser := installFlag.String("peeruser", "", "peer node user (default peeruser=user)")
|
|
||||||
peerPassword := installFlag.String("peerpassword", "", "peer node password (default peerpassword=password)")
|
|
||||||
dstIP := installFlag.String("dstip", "127.0.0.1", "destination ip ")
|
dstIP := installFlag.String("dstip", "127.0.0.1", "destination ip ")
|
||||||
dstPort := installFlag.Int("dstport", 0, "destination port ")
|
dstPort := installFlag.Int("dstport", 0, "destination port ")
|
||||||
srcPort := installFlag.Int("srcport", 0, "source port ")
|
srcPort := installFlag.Int("srcport", 0, "source port ")
|
||||||
protocol := installFlag.String("protocol", "tcp", "tcp or udp")
|
protocol := installFlag.String("protocol", "tcp", "tcp or udp")
|
||||||
noShare := installFlag.Bool("noshare", false, "disable using the huge numbers of shared nodes in OpenP2P network, your connectivity will be weak. also this node will not shared with others")
|
appName := flag.String("appname", "", "app name")
|
||||||
|
installFlag.Bool("noshare", false, "deprecated. uses -sharebandwidth -1")
|
||||||
shareBandwidth := installFlag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private node no limit")
|
shareBandwidth := installFlag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private node no limit")
|
||||||
// logLevel := installFlag.Int("loglevel", 1, "0:debug 1:info 2:warn 3:error")
|
logLevel := installFlag.Int("loglevel", 1, "0:debug 1:info 2:warn 3:error")
|
||||||
installFlag.Parse(os.Args[2:])
|
installFlag.Parse(os.Args[2:])
|
||||||
checkParams(*node, *user, *password)
|
if *node != "" && len(*node) < 8 {
|
||||||
|
gLog.Println(LevelERROR, "node name too short, it must >=8 charaters")
|
||||||
|
os.Exit(9)
|
||||||
|
}
|
||||||
|
if *node == "" { // if node name not set. use os.Hostname
|
||||||
|
hostname, _ := os.Hostname()
|
||||||
|
node = &hostname
|
||||||
|
}
|
||||||
|
gConf.load() // load old config. otherwise will clear all apps
|
||||||
|
gConf.LogLevel = *logLevel
|
||||||
gConf.Network.ServerHost = *serverHost
|
gConf.Network.ServerHost = *serverHost
|
||||||
gConf.Network.User = *user
|
gConf.Network.Token = *token
|
||||||
gConf.Network.Node = *node
|
gConf.Network.Node = *node
|
||||||
gConf.Network.Password = *password
|
gConf.Network.ServerPort = 27183
|
||||||
gConf.Network.ServerPort = 27182
|
|
||||||
gConf.Network.UDPPort1 = 27182
|
gConf.Network.UDPPort1 = 27182
|
||||||
gConf.Network.UDPPort2 = 27183
|
gConf.Network.UDPPort2 = 27183
|
||||||
gConf.Network.NoShare = *noShare
|
|
||||||
gConf.Network.ShareBandwidth = *shareBandwidth
|
gConf.Network.ShareBandwidth = *shareBandwidth
|
||||||
config := AppConfig{}
|
config := AppConfig{Enabled: 1}
|
||||||
config.PeerNode = *peerNode
|
config.PeerNode = *peerNode
|
||||||
config.PeerUser = *peerUser
|
|
||||||
config.PeerPassword = *peerPassword
|
|
||||||
config.DstHost = *dstIP
|
config.DstHost = *dstIP
|
||||||
config.DstPort = *dstPort
|
config.DstPort = *dstPort
|
||||||
config.SrcPort = *srcPort
|
config.SrcPort = *srcPort
|
||||||
config.Protocol = *protocol
|
config.Protocol = *protocol
|
||||||
gConf.add(config)
|
config.AppName = *appName
|
||||||
os.MkdirAll(defaultInstallPath, 0775)
|
if config.SrcPort != 0 {
|
||||||
err := os.Chdir(defaultInstallPath)
|
gConf.add(config, true)
|
||||||
|
}
|
||||||
|
err := os.MkdirAll(defaultInstallPath, 0775)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "MkdirAll %s error:%s", defaultInstallPath, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
err = os.Chdir(defaultInstallPath)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "cd error:", err)
|
gLog.Println(LevelERROR, "cd error:", err)
|
||||||
|
return
|
||||||
}
|
}
|
||||||
gConf.save()
|
gConf.save()
|
||||||
|
targetPath := filepath.Join(defaultInstallPath, defaultBinName)
|
||||||
|
d := daemon{}
|
||||||
// copy files
|
// copy files
|
||||||
|
|
||||||
targetPath := filepath.Join(defaultInstallPath, defaultBinName)
|
|
||||||
binPath, _ := os.Executable()
|
binPath, _ := os.Executable()
|
||||||
src, errFiles := os.Open(binPath) // can not use args[0], on Windows call openp2p is ok(=openp2p.exe)
|
src, errFiles := os.Open(binPath) // can not use args[0], on Windows call openp2p is ok(=openp2p.exe)
|
||||||
if errFiles != nil {
|
if errFiles != nil {
|
||||||
@@ -180,18 +195,14 @@ func install() {
|
|||||||
dst.Close()
|
dst.Close()
|
||||||
|
|
||||||
// install system service
|
// install system service
|
||||||
d := daemon{}
|
|
||||||
|
|
||||||
// args := []string{""}
|
// args := []string{""}
|
||||||
gLog.Println(LevelINFO, "targetPath:", targetPath)
|
gLog.Println(LevelINFO, "targetPath:", targetPath)
|
||||||
err = d.Control("install", targetPath, []string{"-d", "-f"})
|
err = d.Control("install", targetPath, []string{"-d"})
|
||||||
if err != nil {
|
if err == nil {
|
||||||
gLog.Println(LevelERROR, "install system service error:", err)
|
|
||||||
} else {
|
|
||||||
gLog.Println(LevelINFO, "install system service ok.")
|
gLog.Println(LevelINFO, "install system service ok.")
|
||||||
}
|
}
|
||||||
time.Sleep(time.Second * 2)
|
time.Sleep(time.Second * 2)
|
||||||
err = d.Control("start", targetPath, []string{"-d", "-f"})
|
err = d.Control("start", targetPath, []string{"-d"})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "start openp2p service error:", err)
|
gLog.Println(LevelERROR, "start openp2p service error:", err)
|
||||||
} else {
|
} else {
|
||||||
@@ -199,11 +210,45 @@ func install() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func installByFilename() {
|
||||||
|
params := strings.Split(filepath.Base(os.Args[0]), "-")
|
||||||
|
if len(params) < 4 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
serverHost := params[1]
|
||||||
|
token := params[2]
|
||||||
|
gLog.Println(LevelINFO, "install start")
|
||||||
|
targetPath := os.Args[0]
|
||||||
|
args := []string{"install"}
|
||||||
|
args = append(args, "-serverhost")
|
||||||
|
args = append(args, serverHost)
|
||||||
|
args = append(args, "-token")
|
||||||
|
args = append(args, token)
|
||||||
|
env := os.Environ()
|
||||||
|
cmd := exec.Command(targetPath, args...)
|
||||||
|
cmd.Stdout = os.Stdout
|
||||||
|
cmd.Stderr = os.Stderr
|
||||||
|
cmd.Stdin = os.Stdin
|
||||||
|
cmd.Env = env
|
||||||
|
err := cmd.Run()
|
||||||
|
if err != nil {
|
||||||
|
gLog.Println(LevelERROR, "install by filename, start process error:", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
gLog.Println(LevelINFO, "install end")
|
||||||
|
fmt.Println("Press the Any Key to exit")
|
||||||
|
fmt.Scanln()
|
||||||
|
os.Exit(0)
|
||||||
|
}
|
||||||
func uninstall() {
|
func uninstall() {
|
||||||
gLog = InitLogger(filepath.Dir(os.Args[0]), "openp2p-install", LevelDEBUG, 1024*1024, LogFileAndConsole)
|
gLog.Println(LevelINFO, "uninstall start")
|
||||||
|
defer gLog.Println(LevelINFO, "uninstall end")
|
||||||
d := daemon{}
|
d := daemon{}
|
||||||
d.Control("stop", "", nil)
|
err := d.Control("stop", "", nil)
|
||||||
err := d.Control("uninstall", "", nil)
|
if err != nil { // service maybe not install
|
||||||
|
return
|
||||||
|
}
|
||||||
|
err = d.Control("uninstall", "", nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "uninstall system service error:", err)
|
gLog.Println(LevelERROR, "uninstall system service error:", err)
|
||||||
} else {
|
} else {
|
||||||
@@ -211,21 +256,6 @@ func uninstall() {
|
|||||||
}
|
}
|
||||||
binPath := filepath.Join(defaultInstallPath, defaultBinName)
|
binPath := filepath.Join(defaultInstallPath, defaultBinName)
|
||||||
os.Remove(binPath + "0")
|
os.Remove(binPath + "0")
|
||||||
os.Rename(binPath, binPath+"0")
|
os.Remove(binPath)
|
||||||
os.RemoveAll(defaultInstallPath)
|
// os.RemoveAll(defaultInstallPath) // reserve config.json
|
||||||
}
|
|
||||||
|
|
||||||
func checkParams(node, user, password string) {
|
|
||||||
if len(node) < 8 {
|
|
||||||
gLog.Println(LevelERROR, "node name too short, it must >=8 charaters")
|
|
||||||
os.Exit(9)
|
|
||||||
}
|
|
||||||
if len(user) < 8 {
|
|
||||||
gLog.Println(LevelERROR, "user name too short, it must >=8 charaters")
|
|
||||||
os.Exit(9)
|
|
||||||
}
|
|
||||||
if len(password) < 8 {
|
|
||||||
gLog.Println(LevelERROR, "password too short, it must >=8 charaters")
|
|
||||||
os.Exit(9)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
|
After Width: | Height: | Size: 50 KiB |
|
After Width: | Height: | Size: 36 KiB |
|
After Width: | Height: | Size: 25 KiB |
|
After Width: | Height: | Size: 40 KiB |
|
After Width: | Height: | Size: 36 KiB |
|
After Width: | Height: | Size: 29 KiB |
|
After Width: | Height: | Size: 65 KiB |
|
After Width: | Height: | Size: 50 KiB |
@@ -0,0 +1,208 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"encoding/binary"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"os/exec"
|
||||||
|
"path/filepath"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
func handlePush(pn *P2PNetwork, subType uint16, msg []byte) error {
|
||||||
|
pushHead := PushHeader{}
|
||||||
|
err := binary.Read(bytes.NewReader(msg[openP2PHeaderSize:openP2PHeaderSize+PushHeaderSize]), binary.LittleEndian, &pushHead)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
gLog.Printf(LevelDEBUG, "handle push msg type:%d, push header:%+v", subType, pushHead)
|
||||||
|
switch subType {
|
||||||
|
case MsgPushConnectReq:
|
||||||
|
req := PushConnectReq{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong MsgPushConnectReq:%s", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
gLog.Printf(LevelINFO, "%s is connecting...", req.From)
|
||||||
|
gLog.Println(LevelDEBUG, "push connect response to ", req.From)
|
||||||
|
// verify totp token or token
|
||||||
|
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()) ||
|
||||||
|
(req.FromToken == pn.config.Token) {
|
||||||
|
gLog.Printf(LevelINFO, "Access Granted\n")
|
||||||
|
config := AppConfig{}
|
||||||
|
config.peerNatType = req.NatType
|
||||||
|
config.peerConeNatPort = req.ConeNatPort
|
||||||
|
config.peerIP = req.FromIP
|
||||||
|
config.PeerNode = req.From
|
||||||
|
// share relay node will limit bandwidth
|
||||||
|
if req.FromToken != pn.config.Token {
|
||||||
|
gLog.Printf(LevelINFO, "set share bandwidth %d mbps", pn.config.ShareBandwidth)
|
||||||
|
config.shareBandwidth = pn.config.ShareBandwidth
|
||||||
|
}
|
||||||
|
// go pn.AddTunnel(config, req.ID)
|
||||||
|
go pn.addDirectTunnel(config, req.ID)
|
||||||
|
break
|
||||||
|
}
|
||||||
|
gLog.Println(LevelERROR, "Access Denied:", req.From)
|
||||||
|
rsp := PushConnectRsp{
|
||||||
|
Error: 1,
|
||||||
|
Detail: fmt.Sprintf("connect to %s error: Access Denied", pn.config.Node),
|
||||||
|
To: req.From,
|
||||||
|
From: pn.config.Node,
|
||||||
|
}
|
||||||
|
pn.push(req.From, MsgPushConnectRsp, rsp)
|
||||||
|
case MsgPushRsp:
|
||||||
|
rsp := PushRsp{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize:], &rsp)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong pushRsp:%s", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if rsp.Error == 0 {
|
||||||
|
gLog.Printf(LevelDEBUG, "push ok, detail:%s", rsp.Detail)
|
||||||
|
} else {
|
||||||
|
gLog.Printf(LevelERROR, "push error:%d, detail:%s", rsp.Error, rsp.Detail)
|
||||||
|
}
|
||||||
|
case MsgPushAddRelayTunnelReq:
|
||||||
|
req := AddRelayTunnelReq{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong RelayNodeRsp:%s", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
config := AppConfig{}
|
||||||
|
config.PeerNode = req.RelayName
|
||||||
|
config.peerToken = req.RelayToken
|
||||||
|
go func(r AddRelayTunnelReq) {
|
||||||
|
t, errDt := pn.addDirectTunnel(config, 0)
|
||||||
|
if errDt == nil {
|
||||||
|
// notify peer relay ready
|
||||||
|
msg := TunnelMsg{ID: t.id}
|
||||||
|
pn.push(r.From, MsgPushAddRelayTunnelRsp, msg)
|
||||||
|
SaveKey(req.AppID, req.AppKey)
|
||||||
|
}
|
||||||
|
|
||||||
|
}(req)
|
||||||
|
case MsgPushUpdate:
|
||||||
|
gLog.Println(LevelINFO, "MsgPushUpdate")
|
||||||
|
update() // download new version first, then exec ./openp2p update
|
||||||
|
targetPath := filepath.Join(defaultInstallPath, defaultBinName)
|
||||||
|
args := []string{"update"}
|
||||||
|
env := os.Environ()
|
||||||
|
cmd := exec.Command(targetPath, args...)
|
||||||
|
cmd.Stdout = os.Stdout
|
||||||
|
cmd.Stderr = os.Stderr
|
||||||
|
cmd.Stdin = os.Stdin
|
||||||
|
cmd.Env = env
|
||||||
|
err := cmd.Run()
|
||||||
|
if err == nil {
|
||||||
|
os.Exit(0)
|
||||||
|
}
|
||||||
|
return err
|
||||||
|
case MsgPushRestart:
|
||||||
|
gLog.Println(LevelINFO, "MsgPushRestart")
|
||||||
|
os.Exit(0)
|
||||||
|
return err
|
||||||
|
case MsgPushReportApps:
|
||||||
|
gLog.Println(LevelINFO, "MsgPushReportApps")
|
||||||
|
req := ReportApps{}
|
||||||
|
// TODO: add the retrying apps
|
||||||
|
gConf.mtx.Lock()
|
||||||
|
defer gConf.mtx.Unlock()
|
||||||
|
for _, config := range gConf.Apps {
|
||||||
|
appActive := 0
|
||||||
|
i, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
|
if ok {
|
||||||
|
app := i.(*p2pApp)
|
||||||
|
if app.isActive() {
|
||||||
|
appActive = 1
|
||||||
|
}
|
||||||
|
}
|
||||||
|
appInfo := AppInfo{
|
||||||
|
AppName: config.AppName,
|
||||||
|
Protocol: config.Protocol,
|
||||||
|
SrcPort: config.SrcPort,
|
||||||
|
// RelayNode: relayNode,
|
||||||
|
PeerNode: config.PeerNode,
|
||||||
|
DstHost: config.DstHost,
|
||||||
|
DstPort: config.DstPort,
|
||||||
|
PeerUser: config.PeerUser,
|
||||||
|
PeerIP: config.peerIP,
|
||||||
|
PeerNatType: config.peerNatType,
|
||||||
|
RetryTime: config.retryTime.String(),
|
||||||
|
IsActive: appActive,
|
||||||
|
Enabled: config.Enabled,
|
||||||
|
}
|
||||||
|
req.Apps = append(req.Apps, appInfo)
|
||||||
|
}
|
||||||
|
pn.write(MsgReport, MsgReportApps, &req)
|
||||||
|
case MsgPushEditApp:
|
||||||
|
gLog.Println(LevelINFO, "MsgPushEditApp")
|
||||||
|
newApp := AppInfo{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize:], &newApp)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong MsgPushEditApp:%s %s", err, string(msg[openP2PHeaderSize:]))
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
oldConf := AppConfig{Enabled: 1}
|
||||||
|
// protocol0+srcPort0 exist, delApp
|
||||||
|
oldConf.AppName = newApp.AppName
|
||||||
|
oldConf.Protocol = newApp.Protocol0
|
||||||
|
oldConf.SrcPort = newApp.SrcPort0
|
||||||
|
oldConf.PeerNode = newApp.PeerNode
|
||||||
|
oldConf.DstHost = newApp.DstHost
|
||||||
|
oldConf.DstPort = newApp.DstPort
|
||||||
|
|
||||||
|
gConf.delete(oldConf)
|
||||||
|
// AddApp
|
||||||
|
newConf := oldConf
|
||||||
|
newConf.Protocol = newApp.Protocol
|
||||||
|
newConf.SrcPort = newApp.SrcPort
|
||||||
|
gConf.add(newConf, false)
|
||||||
|
gConf.save() // save quickly for the next request reportApplist
|
||||||
|
pn.DeleteApp(oldConf) // DeleteApp may cost some times, execute at the end
|
||||||
|
// autoReconnect will auto AddApp
|
||||||
|
// pn.AddApp(config)
|
||||||
|
// TODO: report result
|
||||||
|
case MsgPushEditNode:
|
||||||
|
gLog.Println(LevelINFO, "MsgPushEditNode")
|
||||||
|
req := EditNode{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize:], &req)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong MsgPushEditNode:%s %s", err, string(msg[openP2PHeaderSize:]))
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
gConf.mtx.Lock()
|
||||||
|
gConf.Network.Node = req.NewName
|
||||||
|
gConf.Network.ShareBandwidth = req.Bandwidth
|
||||||
|
gConf.mtx.Unlock()
|
||||||
|
gConf.save()
|
||||||
|
// TODO: hot reload
|
||||||
|
os.Exit(0)
|
||||||
|
case MsgPushSwitchApp:
|
||||||
|
gLog.Println(LevelINFO, "MsgPushSwitchApp")
|
||||||
|
app := AppInfo{}
|
||||||
|
err := json.Unmarshal(msg[openP2PHeaderSize:], &app)
|
||||||
|
if err != nil {
|
||||||
|
gLog.Printf(LevelERROR, "wrong MsgPushSwitchApp:%s %s", err, string(msg[openP2PHeaderSize:]))
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
config := AppConfig{Enabled: app.Enabled, SrcPort: app.SrcPort, Protocol: app.Protocol}
|
||||||
|
gLog.Println(LevelINFO, app.AppName, " switch to ", app.Enabled)
|
||||||
|
gConf.switchApp(config, app.Enabled)
|
||||||
|
if app.Enabled == 0 {
|
||||||
|
// disable APP
|
||||||
|
pn.DeleteApp(config)
|
||||||
|
}
|
||||||
|
default:
|
||||||
|
pn.msgMapMtx.Lock()
|
||||||
|
ch := pn.msgMap[pushHead.From]
|
||||||
|
pn.msgMapMtx.Unlock()
|
||||||
|
ch <- msg
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -47,7 +47,7 @@ const (
|
|||||||
type V8log struct {
|
type V8log struct {
|
||||||
loggers map[LogLevel]*log.Logger
|
loggers map[LogLevel]*log.Logger
|
||||||
files map[LogLevel]*os.File
|
files map[LogLevel]*os.File
|
||||||
llevel LogLevel
|
level LogLevel
|
||||||
stopSig chan bool
|
stopSig chan bool
|
||||||
logDir string
|
logDir string
|
||||||
mtx *sync.Mutex
|
mtx *sync.Mutex
|
||||||
@@ -92,17 +92,10 @@ func InitLogger(path string, filePrefix string, level LogLevel, maxLogSize int64
|
|||||||
return pLog
|
return pLog
|
||||||
}
|
}
|
||||||
|
|
||||||
// UninitLogger ...
|
func (vl *V8log) setLevel(level LogLevel) {
|
||||||
func (vl *V8log) UninitLogger() {
|
vl.mtx.Lock()
|
||||||
if !vl.stoped {
|
defer vl.mtx.Unlock()
|
||||||
vl.stoped = true
|
vl.level = level
|
||||||
close(vl.stopSig)
|
|
||||||
for l := range logFileNames {
|
|
||||||
if l >= vl.llevel {
|
|
||||||
vl.files[l].Close()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (vl *V8log) checkFile() {
|
func (vl *V8log) checkFile() {
|
||||||
@@ -150,7 +143,7 @@ func (vl *V8log) Printf(level LogLevel, format string, params ...interface{}) {
|
|||||||
if vl.stoped {
|
if vl.stoped {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if level < vl.llevel {
|
if level < vl.level {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
pidAndLevel := []interface{}{vl.pid, loglevel[level]}
|
pidAndLevel := []interface{}{vl.pid, loglevel[level]}
|
||||||
@@ -170,7 +163,7 @@ func (vl *V8log) Println(level LogLevel, params ...interface{}) {
|
|||||||
if vl.stoped {
|
if vl.stoped {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if level < vl.llevel {
|
if level < vl.level {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
pidAndLevel := []interface{}{vl.pid, " ", loglevel[level], " "}
|
pidAndLevel := []interface{}{vl.pid, " ", loglevel[level], " "}
|
||||||
|
|||||||
@@ -43,7 +43,7 @@ func natTest(serverHost string, serverPort int, localPort int, echoPort int) (pu
|
|||||||
// testing for public ip
|
// testing for public ip
|
||||||
if echoPort != 0 {
|
if echoPort != 0 {
|
||||||
for {
|
for {
|
||||||
gLog.Printf(LevelINFO, "public ip test start %s:%d", natRsp.IP, echoPort)
|
gLog.Printf(LevelDEBUG, "public ip test start %s:%d", natRsp.IP, echoPort)
|
||||||
conn, err := net.ListenUDP("udp", nil)
|
conn, err := net.ListenUDP("udp", nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
break
|
break
|
||||||
@@ -60,10 +60,10 @@ func natTest(serverHost string, serverPort int, localPort int, echoPort int) (pu
|
|||||||
conn.SetReadDeadline(time.Now().Add(PublicIPEchoTimeout))
|
conn.SetReadDeadline(time.Now().Add(PublicIPEchoTimeout))
|
||||||
_, _, err = conn.ReadFromUDP(buf)
|
_, _, err = conn.ReadFromUDP(buf)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
gLog.Println(LevelINFO, "public ip:YES")
|
gLog.Println(LevelDEBUG, "public ip:YES")
|
||||||
natRsp.IsPublicIP = 1
|
natRsp.IsPublicIP = 1
|
||||||
} else {
|
} else {
|
||||||
gLog.Println(LevelINFO, "public ip:NO")
|
gLog.Println(LevelDEBUG, "public ip:NO")
|
||||||
}
|
}
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,7 +1,6 @@
|
|||||||
package main
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"flag"
|
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
"os"
|
"os"
|
||||||
@@ -11,9 +10,11 @@ import (
|
|||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
rand.Seed(time.Now().UnixNano())
|
rand.Seed(time.Now().UnixNano())
|
||||||
|
binDir := filepath.Dir(os.Args[0])
|
||||||
|
os.Chdir(binDir) // for system service
|
||||||
|
gLog = InitLogger(binDir, "openp2p", LevelDEBUG, 1024*1024, LogFileAndConsole)
|
||||||
|
gLog.Println(LevelINFO, "openp2p start. version: ", OpenP2PVersion)
|
||||||
// TODO: install sub command, deamon process
|
// TODO: install sub command, deamon process
|
||||||
// groups := flag.String("groups", "", "you could join in several groups. like: GroupName1:Password1;GroupName2:Password2; group name 8-31 characters")
|
|
||||||
if len(os.Args) > 1 {
|
if len(os.Args) > 1 {
|
||||||
switch os.Args[1] {
|
switch os.Args[1] {
|
||||||
case "version", "-v", "--version":
|
case "version", "-v", "--version":
|
||||||
@@ -21,10 +22,9 @@ func main() {
|
|||||||
return
|
return
|
||||||
case "update":
|
case "update":
|
||||||
gLog = InitLogger(filepath.Dir(os.Args[0]), "openp2p", LevelDEBUG, 1024*1024, LogFileAndConsole)
|
gLog = InitLogger(filepath.Dir(os.Args[0]), "openp2p", LevelDEBUG, 1024*1024, LogFileAndConsole)
|
||||||
update()
|
|
||||||
targetPath := filepath.Join(defaultInstallPath, defaultBinName)
|
targetPath := filepath.Join(defaultInstallPath, defaultBinName)
|
||||||
d := daemon{}
|
d := daemon{}
|
||||||
err := d.Control("restart", targetPath, []string{"-d", "-f"})
|
err := d.Control("restart", targetPath, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "restart service error:", err)
|
gLog.Println(LevelERROR, "restart service error:", err)
|
||||||
} else {
|
} else {
|
||||||
@@ -38,124 +38,18 @@ func main() {
|
|||||||
uninstall()
|
uninstall()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
|
||||||
serverHost := flag.String("serverhost", "api.openp2p.cn", "server host ")
|
|
||||||
// serverHost := flag.String("serverhost", "127.0.0.1", "server host ") // for debug
|
|
||||||
user := flag.String("user", "", "user name. 8-31 characters")
|
|
||||||
node := flag.String("node", "", "node name. 8-31 characters")
|
|
||||||
password := flag.String("password", "", "user password. 8-31 characters")
|
|
||||||
peerNode := flag.String("peernode", "", "peer node name that you want to connect")
|
|
||||||
peerUser := flag.String("peeruser", "", "peer node user (default peeruser=user)")
|
|
||||||
peerPassword := flag.String("peerpassword", "", "peer node password (default peerpassword=password)")
|
|
||||||
dstIP := flag.String("dstip", "127.0.0.1", "destination ip ")
|
|
||||||
dstPort := flag.Int("dstport", 0, "destination port ")
|
|
||||||
srcPort := flag.Int("srcport", 0, "source port ")
|
|
||||||
protocol := flag.String("protocol", "tcp", "tcp or udp")
|
|
||||||
noShare := flag.Bool("noshare", false, "disable using the huge numbers of shared nodes in OpenP2P network, your connectivity will be weak. also this node will not shared with others")
|
|
||||||
shareBandwidth := flag.Int("sharebandwidth", 10, "N mbps share bandwidth limit, private node no limit")
|
|
||||||
configFile := flag.Bool("f", false, "config file")
|
|
||||||
daemonMode := flag.Bool("d", false, "daemonMode")
|
|
||||||
byDaemon := flag.Bool("bydaemon", false, "start by daemon")
|
|
||||||
logLevel := flag.Int("loglevel", 1, "0:debug 1:info 2:warn 3:error")
|
|
||||||
flag.Parse()
|
|
||||||
|
|
||||||
gLog = InitLogger(filepath.Dir(os.Args[0]), "openp2p", LogLevel(*logLevel), 1024*1024, LogFileAndConsole)
|
|
||||||
gLog.Println(LevelINFO, "openp2p start. version: ", OpenP2PVersion)
|
|
||||||
if *daemonMode {
|
|
||||||
d := daemon{}
|
|
||||||
d.run()
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if !*configFile {
|
|
||||||
// validate cmd params
|
|
||||||
checkParams(*node, *user, *password)
|
|
||||||
if *peerNode != "" {
|
|
||||||
if *dstPort == 0 {
|
|
||||||
gLog.Println(LevelERROR, "dstPort not set")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if *srcPort == 0 {
|
|
||||||
gLog.Println(LevelERROR, "srcPort not set")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
config := AppConfig{}
|
|
||||||
config.PeerNode = *peerNode
|
|
||||||
config.PeerUser = *peerUser
|
|
||||||
config.PeerPassword = *peerPassword
|
|
||||||
config.DstHost = *dstIP
|
|
||||||
config.DstPort = *dstPort
|
|
||||||
config.SrcPort = *srcPort
|
|
||||||
config.Protocol = *protocol
|
|
||||||
gLog.Println(LevelINFO, config)
|
|
||||||
if *configFile {
|
|
||||||
if err := gConf.load(); err != nil {
|
|
||||||
gLog.Println(LevelERROR, "load config error. exit.")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
} else {
|
} else {
|
||||||
gConf.add(config)
|
installByFilename()
|
||||||
gConf.Network = NetworkConfig{
|
|
||||||
Node: *node,
|
|
||||||
User: *user,
|
|
||||||
Password: *password,
|
|
||||||
NoShare: *noShare,
|
|
||||||
ServerHost: *serverHost,
|
|
||||||
ServerPort: 27182,
|
|
||||||
UDPPort1: 27182,
|
|
||||||
UDPPort2: 27183,
|
|
||||||
ipv6: "240e:3b7:621:def0:fda4:dd7f:36a1:2803", // TODO: detect real ipv6
|
|
||||||
ShareBandwidth: *shareBandwidth,
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
// gConf.save() // not change config file
|
|
||||||
gConf.daemonMode = *byDaemon
|
|
||||||
|
|
||||||
gLog.Println(LevelINFO, gConf)
|
parseParams()
|
||||||
|
gLog.Println(LevelINFO, &gConf)
|
||||||
setFirewall()
|
setFirewall()
|
||||||
network := P2PNetworkInstance(&gConf.Network)
|
network := P2PNetworkInstance(&gConf.Network)
|
||||||
if ok := network.Connect(30000); !ok {
|
if ok := network.Connect(30000); !ok {
|
||||||
gLog.Println(LevelERROR, "P2PNetwork login error")
|
gLog.Println(LevelERROR, "P2PNetwork login error")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
for _, app := range gConf.Apps {
|
|
||||||
// set default peer user password
|
|
||||||
if app.PeerPassword == "" {
|
|
||||||
app.PeerPassword = gConf.Network.Password
|
|
||||||
}
|
|
||||||
if app.PeerUser == "" {
|
|
||||||
app.PeerUser = gConf.Network.User
|
|
||||||
}
|
|
||||||
err := network.AddApp(app)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Println(LevelERROR, "addTunnel error")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// test
|
|
||||||
// go func() {
|
|
||||||
|
|
||||||
// time.Sleep(time.Second * 30)
|
|
||||||
// config := AppConfig{}
|
|
||||||
// config.PeerNode = *peerNode
|
|
||||||
// config.PeerUser = *peerUser
|
|
||||||
// config.PeerPassword = *peerPassword
|
|
||||||
// config.DstHost = *dstIP
|
|
||||||
// config.DstPort = *dstPort
|
|
||||||
// config.SrcPort = 32
|
|
||||||
// config.Protocol = *protocol
|
|
||||||
// network.AddApp(config)
|
|
||||||
// // time.Sleep(time.Second * 30)
|
|
||||||
// // network.DeleteTunnel(config)
|
|
||||||
// // time.Sleep(time.Second * 30)
|
|
||||||
// // network.DeleteTunnel(config)
|
|
||||||
// }()
|
|
||||||
|
|
||||||
// // TODO: http api
|
|
||||||
// api := ClientAPI{}
|
|
||||||
// go api.run()
|
|
||||||
gLog.Println(LevelINFO, "waiting for connection...")
|
gLog.Println(LevelINFO, "waiting for connection...")
|
||||||
forever := make(chan bool)
|
forever := make(chan bool)
|
||||||
<-forever
|
<-forever
|
||||||
|
|||||||
@@ -21,8 +21,8 @@ type overlayTCP struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (otcp *overlayTCP) run() {
|
func (otcp *overlayTCP) run() {
|
||||||
gLog.Printf(LevelINFO, "%d overlayTCP run start", otcp.id)
|
gLog.Printf(LevelDEBUG, "%d overlayTCP run start", otcp.id)
|
||||||
defer gLog.Printf(LevelINFO, "%d overlayTCP run end", otcp.id)
|
defer gLog.Printf(LevelDEBUG, "%d overlayTCP run end", otcp.id)
|
||||||
otcp.running = true
|
otcp.running = true
|
||||||
buffer := make([]byte, ReadBuffLen+PaddingSize)
|
buffer := make([]byte, ReadBuffLen+PaddingSize)
|
||||||
readBuf := buffer[:ReadBuffLen]
|
readBuf := buffer[:ReadBuffLen]
|
||||||
|
|||||||
@@ -11,16 +11,17 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type p2pApp struct {
|
type p2pApp struct {
|
||||||
config AppConfig
|
config AppConfig
|
||||||
listener net.Listener
|
listener net.Listener
|
||||||
tunnel *P2PTunnel
|
tunnel *P2PTunnel
|
||||||
rtid uint64
|
rtid uint64
|
||||||
hbTime time.Time
|
relayNode string
|
||||||
hbMtx sync.Mutex
|
hbTime time.Time
|
||||||
running bool
|
hbMtx sync.Mutex
|
||||||
id uint64
|
running bool
|
||||||
key uint64
|
id uint64
|
||||||
wg sync.WaitGroup
|
key uint64
|
||||||
|
wg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
||||||
func (app *p2pApp) isActive() bool {
|
func (app *p2pApp) isActive() bool {
|
||||||
@@ -72,11 +73,10 @@ func (app *p2pApp) listenTCP() error {
|
|||||||
otcp.appKeyBytes = encryptKey
|
otcp.appKeyBytes = encryptKey
|
||||||
}
|
}
|
||||||
app.tunnel.overlayConns.Store(otcp.id, &otcp)
|
app.tunnel.overlayConns.Store(otcp.id, &otcp)
|
||||||
gLog.Printf(LevelINFO, "Accept overlayID:%d", otcp.id)
|
gLog.Printf(LevelDEBUG, "Accept overlayID:%d", otcp.id)
|
||||||
// tell peer connect
|
// tell peer connect
|
||||||
req := OverlayConnectReq{ID: otcp.id,
|
req := OverlayConnectReq{ID: otcp.id,
|
||||||
User: app.config.PeerUser,
|
Token: app.tunnel.pn.config.Token,
|
||||||
Password: app.config.PeerPassword,
|
|
||||||
DstIP: app.config.DstHost,
|
DstIP: app.config.DstHost,
|
||||||
DstPort: app.config.DstPort,
|
DstPort: app.config.DstPort,
|
||||||
Protocol: app.config.Protocol,
|
Protocol: app.config.Protocol,
|
||||||
@@ -108,7 +108,9 @@ func (app *p2pApp) listen() error {
|
|||||||
go app.relayHeartbeatLoop()
|
go app.relayHeartbeatLoop()
|
||||||
}
|
}
|
||||||
for app.running {
|
for app.running {
|
||||||
if app.config.Protocol == "tcp" {
|
if app.config.Protocol == "udp" {
|
||||||
|
app.listenTCP()
|
||||||
|
} else {
|
||||||
app.listenTCP()
|
app.listenTCP()
|
||||||
}
|
}
|
||||||
time.Sleep(time.Second * 5)
|
time.Sleep(time.Second * 5)
|
||||||
|
|||||||
@@ -37,7 +37,7 @@ type P2PNetwork struct {
|
|||||||
msgMapMtx sync.Mutex
|
msgMapMtx sync.Mutex
|
||||||
config NetworkConfig
|
config NetworkConfig
|
||||||
allTunnels sync.Map
|
allTunnels sync.Map
|
||||||
apps sync.Map
|
apps sync.Map //key: protocol+srcport; value: p2pApp
|
||||||
limiter *BandwidthLimiter
|
limiter *BandwidthLimiter
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -63,7 +63,7 @@ func P2PNetworkInstance(config *NetworkConfig) *P2PNetwork {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) run() {
|
func (pn *P2PNetwork) run() {
|
||||||
go pn.autoReconnectApp()
|
go pn.autorunApp()
|
||||||
heartbeatTimer := time.NewTicker(NetworkHeartbeatTime)
|
heartbeatTimer := time.NewTicker(NetworkHeartbeatTime)
|
||||||
for pn.running {
|
for pn.running {
|
||||||
select {
|
select {
|
||||||
@@ -93,55 +93,57 @@ func (pn *P2PNetwork) Connect(timeout int) bool {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) autoReconnectApp() {
|
func (pn *P2PNetwork) runAll() {
|
||||||
gLog.Println(LevelINFO, "autoReconnectApp start")
|
gConf.mtx.Lock()
|
||||||
retryApps := make([]AppConfig, 0)
|
defer gConf.mtx.Unlock()
|
||||||
|
for _, config := range gConf.Apps {
|
||||||
|
if config.Enabled == 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if config.AppName == "" {
|
||||||
|
config.AppName = fmt.Sprintf("%s%d", config.Protocol, config.SrcPort)
|
||||||
|
}
|
||||||
|
appExist := false
|
||||||
|
appActive := false
|
||||||
|
i, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
|
if ok {
|
||||||
|
app := i.(*p2pApp)
|
||||||
|
appExist = true
|
||||||
|
if app.isActive() {
|
||||||
|
appActive = true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if appExist && appActive {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if appExist && !appActive {
|
||||||
|
gLog.Printf(LevelINFO, "detect app %s disconnect, reconnecting...", config.AppName)
|
||||||
|
pn.DeleteApp(config)
|
||||||
|
if config.retryTime.Add(time.Minute * 15).Before(time.Now()) {
|
||||||
|
config.retryNum = 0
|
||||||
|
}
|
||||||
|
config.retryNum++
|
||||||
|
config.retryTime = time.Now()
|
||||||
|
if config.retryNum > MaxRetry {
|
||||||
|
gLog.Printf(LevelERROR, "app %s%d retry more than %d times, exit.", config.Protocol, config.SrcPort, MaxRetry)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
}
|
||||||
|
go pn.AddApp(config)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
func (pn *P2PNetwork) autorunApp() {
|
||||||
|
gLog.Println(LevelINFO, "autorunApp start")
|
||||||
|
// TODO: use gConf to check reconnect
|
||||||
for pn.running {
|
for pn.running {
|
||||||
time.Sleep(time.Second)
|
time.Sleep(time.Second)
|
||||||
if !pn.online {
|
if !pn.online {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if len(retryApps) > 0 {
|
pn.runAll()
|
||||||
gLog.Printf(LevelINFO, "retryApps len=%d", len(retryApps))
|
time.Sleep(time.Second * 10)
|
||||||
thisRound := make([]AppConfig, 0)
|
|
||||||
for i := 0; i < len(retryApps); i++ {
|
|
||||||
// reset retryNum when running 15min continuously
|
|
||||||
if retryApps[i].retryTime.Add(time.Minute * 15).Before(time.Now()) {
|
|
||||||
retryApps[i].retryNum = 0
|
|
||||||
}
|
|
||||||
retryApps[i].retryNum++
|
|
||||||
retryApps[i].retryTime = time.Now()
|
|
||||||
if retryApps[i].retryNum > MaxRetry {
|
|
||||||
gLog.Printf(LevelERROR, "app %s%d retry more than %d times, exit.", retryApps[i].Protocol, retryApps[i].SrcPort, MaxRetry)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
pn.DeleteApp(retryApps[i])
|
|
||||||
if err := pn.AddApp(retryApps[i]); err != nil {
|
|
||||||
gLog.Printf(LevelERROR, "AddApp %s%d error:%s", retryApps[i].Protocol, retryApps[i].SrcPort, err)
|
|
||||||
thisRound = append(thisRound, retryApps[i])
|
|
||||||
time.Sleep(RetryInterval)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
retryApps = thisRound
|
|
||||||
}
|
|
||||||
pn.apps.Range(func(_, i interface{}) bool {
|
|
||||||
app := i.(*p2pApp)
|
|
||||||
if app.isActive() {
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
gLog.Printf(LevelINFO, "detect app %s%d disconnect,last hb %s reconnecting...", app.config.Protocol, app.config.SrcPort, app.hbTime)
|
|
||||||
config := app.config
|
|
||||||
// clear peerinfo
|
|
||||||
config.peerConeNatPort = 0
|
|
||||||
config.peerIP = ""
|
|
||||||
config.peerNatType = 0
|
|
||||||
config.peerToken = 0
|
|
||||||
pn.DeleteApp(config)
|
|
||||||
retryApps = append(retryApps, config)
|
|
||||||
return true
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
gLog.Println(LevelINFO, "autoReconnectApp end")
|
gLog.Println(LevelINFO, "autorunApp end")
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) addRelayTunnel(config AppConfig, appid uint64, appkey uint64) (*P2PTunnel, uint64, error) {
|
func (pn *P2PNetwork) addRelayTunnel(config AppConfig, appid uint64, appkey uint64) (*P2PTunnel, uint64, error) {
|
||||||
@@ -198,21 +200,17 @@ func (pn *P2PNetwork) addRelayTunnel(config AppConfig, appid uint64, appkey uint
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
||||||
gLog.Printf(LevelINFO, "addApp %s%d to %s:%s:%d start", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
gLog.Printf(LevelINFO, "addApp %s to %s:%s:%d start", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
||||||
defer gLog.Printf(LevelINFO, "addApp %s%d to %s:%s:%d end", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
defer gLog.Printf(LevelINFO, "addApp %s to %s:%s:%d end", config.AppName, config.PeerNode, config.DstHost, config.DstPort)
|
||||||
if !pn.online {
|
if !pn.online {
|
||||||
return errors.New("P2PNetwork offline")
|
return errors.New("P2PNetwork offline")
|
||||||
}
|
}
|
||||||
// check if app already exist?
|
// check if app already exist?
|
||||||
appExist := false
|
appExist := false
|
||||||
pn.apps.Range(func(_, i interface{}) bool {
|
_, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
app := i.(*p2pApp)
|
if ok {
|
||||||
if app.config.Protocol == config.Protocol && app.config.SrcPort == config.SrcPort {
|
appExist = true
|
||||||
appExist = true
|
}
|
||||||
return false
|
|
||||||
}
|
|
||||||
return true
|
|
||||||
})
|
|
||||||
if appExist {
|
if appExist {
|
||||||
return errors.New("P2PApp already exist")
|
return errors.New("P2PApp already exist")
|
||||||
}
|
}
|
||||||
@@ -221,7 +219,7 @@ func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
|||||||
t, err := pn.addDirectTunnel(config, 0)
|
t, err := pn.addDirectTunnel(config, 0)
|
||||||
var rtid uint64
|
var rtid uint64
|
||||||
relayNode := ""
|
relayNode := ""
|
||||||
peerNatType := 100
|
peerNatType := NATUnknown
|
||||||
peerIP := ""
|
peerIP := ""
|
||||||
errMsg := ""
|
errMsg := ""
|
||||||
if err != nil && err == ErrorHandshake {
|
if err != nil && err == ErrorHandshake {
|
||||||
@@ -247,7 +245,6 @@ func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
|||||||
PeerNode: config.PeerNode,
|
PeerNode: config.PeerNode,
|
||||||
DstPort: config.DstPort,
|
DstPort: config.DstPort,
|
||||||
DstHost: config.DstHost,
|
DstHost: config.DstHost,
|
||||||
PeerUser: config.PeerUser,
|
|
||||||
PeerNatType: peerNatType,
|
PeerNatType: peerNatType,
|
||||||
PeerIP: peerIP,
|
PeerIP: peerIP,
|
||||||
ShareBandwidth: pn.config.ShareBandwidth,
|
ShareBandwidth: pn.config.ShareBandwidth,
|
||||||
@@ -255,15 +252,18 @@ func (pn *P2PNetwork) AddApp(config AppConfig) error {
|
|||||||
Version: OpenP2PVersion,
|
Version: OpenP2PVersion,
|
||||||
}
|
}
|
||||||
pn.write(MsgReport, MsgReportConnect, &req)
|
pn.write(MsgReport, MsgReportConnect, &req)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
app := p2pApp{
|
app := p2pApp{
|
||||||
id: appID,
|
id: appID,
|
||||||
key: appKey,
|
key: appKey,
|
||||||
tunnel: t,
|
tunnel: t,
|
||||||
config: config,
|
config: config,
|
||||||
rtid: rtid,
|
rtid: rtid,
|
||||||
hbTime: time.Now()}
|
relayNode: relayNode,
|
||||||
pn.apps.Store(appID, &app)
|
hbTime: time.Now()}
|
||||||
|
pn.apps.Store(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort), &app)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
go app.listen()
|
go app.listen()
|
||||||
}
|
}
|
||||||
@@ -274,22 +274,18 @@ func (pn *P2PNetwork) DeleteApp(config AppConfig) {
|
|||||||
gLog.Printf(LevelINFO, "DeleteApp %s%d start", config.Protocol, config.SrcPort)
|
gLog.Printf(LevelINFO, "DeleteApp %s%d start", config.Protocol, config.SrcPort)
|
||||||
defer gLog.Printf(LevelINFO, "DeleteApp %s%d end", config.Protocol, config.SrcPort)
|
defer gLog.Printf(LevelINFO, "DeleteApp %s%d end", config.Protocol, config.SrcPort)
|
||||||
// close the apps of this config
|
// close the apps of this config
|
||||||
pn.apps.Range(func(_, i interface{}) bool {
|
i, ok := pn.apps.Load(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
|
if ok {
|
||||||
app := i.(*p2pApp)
|
app := i.(*p2pApp)
|
||||||
if app.config.Protocol == config.Protocol && app.config.SrcPort == config.SrcPort {
|
gLog.Printf(LevelINFO, "app %s exist, delete it", fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
gLog.Printf(LevelINFO, "app %s exist, delete it", fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
app.close()
|
||||||
app := i.(*p2pApp)
|
pn.apps.Delete(fmt.Sprintf("%s%d", config.Protocol, config.SrcPort))
|
||||||
app.close()
|
}
|
||||||
pn.apps.Delete(app.id)
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
return true
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) addDirectTunnel(config AppConfig, tid uint64) (*P2PTunnel, error) {
|
func (pn *P2PNetwork) addDirectTunnel(config AppConfig, tid uint64) (*P2PTunnel, error) {
|
||||||
gLog.Printf(LevelINFO, "addDirectTunnel %s%d to %s:%s:%d start", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
gLog.Printf(LevelDEBUG, "addDirectTunnel %s%d to %s:%s:%d start", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
||||||
defer gLog.Printf(LevelINFO, "addDirectTunnel %s%d to %s:%s:%d end", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
defer gLog.Printf(LevelDEBUG, "addDirectTunnel %s%d to %s:%s:%d end", config.Protocol, config.SrcPort, config.PeerNode, config.DstHost, config.DstPort)
|
||||||
isClient := false
|
isClient := false
|
||||||
// client side tid=0, assign random uint64
|
// client side tid=0, assign random uint64
|
||||||
if tid == 0 {
|
if tid == 0 {
|
||||||
@@ -377,10 +373,10 @@ func (pn *P2PNetwork) init() error {
|
|||||||
pn.config.natType = NATSymmetric
|
pn.config.natType = NATSymmetric
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelINFO, "detect NAT type error:", err)
|
gLog.Println(LevelDEBUG, "detect NAT type error:", err)
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
gLog.Println(LevelINFO, "detect NAT type:", pn.config.natType, " publicIP:", pn.config.publicIP)
|
gLog.Println(LevelDEBUG, "detect NAT type:", pn.config.natType, " publicIP:", pn.config.publicIP)
|
||||||
gatewayURL := fmt.Sprintf("%s:%d", pn.config.ServerHost, pn.config.ServerPort)
|
gatewayURL := fmt.Sprintf("%s:%d", pn.config.ServerHost, pn.config.ServerPort)
|
||||||
forwardPath := "/openp2p/v1/login"
|
forwardPath := "/openp2p/v1/login"
|
||||||
config := tls.Config{InsecureSkipVerify: true} // let's encrypt root cert "DST Root CA X3" expired at 2021/09/29. many old system(windows server 2008 etc) will not trust our cert
|
config := tls.Config{InsecureSkipVerify: true} // let's encrypt root cert "DST Root CA X3" expired at 2021/09/29. many old system(windows server 2008 etc) will not trust our cert
|
||||||
@@ -388,16 +384,10 @@ func (pn *P2PNetwork) init() error {
|
|||||||
u := url.URL{Scheme: "wss", Host: gatewayURL, Path: forwardPath}
|
u := url.URL{Scheme: "wss", Host: gatewayURL, Path: forwardPath}
|
||||||
q := u.Query()
|
q := u.Query()
|
||||||
q.Add("node", pn.config.Node)
|
q.Add("node", pn.config.Node)
|
||||||
q.Add("user", pn.config.User)
|
q.Add("token", fmt.Sprintf("%d", pn.config.Token))
|
||||||
q.Add("password", pn.config.Password)
|
|
||||||
q.Add("version", OpenP2PVersion)
|
q.Add("version", OpenP2PVersion)
|
||||||
q.Add("nattype", fmt.Sprintf("%d", pn.config.natType))
|
q.Add("nattype", fmt.Sprintf("%d", pn.config.natType))
|
||||||
|
q.Add("sharebandwidth", fmt.Sprintf("%d", pn.config.ShareBandwidth))
|
||||||
noShareStr := "false"
|
|
||||||
if pn.config.NoShare {
|
|
||||||
noShareStr = "true"
|
|
||||||
}
|
|
||||||
q.Add("noshare", noShareStr)
|
|
||||||
u.RawQuery = q.Encode()
|
u.RawQuery = q.Encode()
|
||||||
var ws *websocket.Conn
|
var ws *websocket.Conn
|
||||||
ws, _, err = websocket.DefaultDialer.Dial(u.String(), nil)
|
ws, _, err = websocket.DefaultDialer.Dial(u.String(), nil)
|
||||||
@@ -425,7 +415,7 @@ func (pn *P2PNetwork) init() error {
|
|||||||
Version: OpenP2PVersion,
|
Version: OpenP2PVersion,
|
||||||
}
|
}
|
||||||
rsp := netInfo()
|
rsp := netInfo()
|
||||||
gLog.Println(LevelINFO, rsp)
|
gLog.Println(LevelDEBUG, "netinfo:", rsp)
|
||||||
if rsp != nil && rsp.Country != "" {
|
if rsp != nil && rsp.Country != "" {
|
||||||
if len(rsp.IP) == net.IPv6len {
|
if len(rsp.IP) == net.IPv6len {
|
||||||
pn.config.ipv6 = rsp.IP.String()
|
pn.config.ipv6 = rsp.IP.String()
|
||||||
@@ -434,7 +424,7 @@ func (pn *P2PNetwork) init() error {
|
|||||||
req.NetInfo = *rsp
|
req.NetInfo = *rsp
|
||||||
}
|
}
|
||||||
pn.write(MsgReport, MsgReportBasic, &req)
|
pn.write(MsgReport, MsgReportBasic, &req)
|
||||||
gLog.Println(LevelINFO, "P2PNetwork init ok")
|
gLog.Println(LevelDEBUG, "P2PNetwork init ok")
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -466,13 +456,20 @@ func (pn *P2PNetwork) handleMessage(t int, msg []byte) {
|
|||||||
pn.running = false
|
pn.running = false
|
||||||
} else {
|
} else {
|
||||||
pn.serverTs = rsp.Ts
|
pn.serverTs = rsp.Ts
|
||||||
|
pn.config.Token = rsp.Token
|
||||||
|
pn.config.User = rsp.User
|
||||||
|
gConf.mtx.Lock()
|
||||||
|
gConf.Network.Token = rsp.Token
|
||||||
|
gConf.Network.User = rsp.User
|
||||||
|
gConf.mtx.Unlock()
|
||||||
|
gConf.save()
|
||||||
pn.localTs = time.Now().Unix()
|
pn.localTs = time.Now().Unix()
|
||||||
gLog.Printf(LevelINFO, "login ok. Server ts=%d, local ts=%d", rsp.Ts, pn.localTs)
|
gLog.Printf(LevelINFO, "login ok. user=%s,Server ts=%d, local ts=%d", rsp.User, rsp.Ts, pn.localTs)
|
||||||
}
|
}
|
||||||
case MsgHeartbeat:
|
case MsgHeartbeat:
|
||||||
gLog.Printf(LevelDEBUG, "P2PNetwork heartbeat ok")
|
gLog.Printf(LevelDEBUG, "P2PNetwork heartbeat ok")
|
||||||
case MsgPush:
|
case MsgPush:
|
||||||
pn.handlePush(head.SubType, msg)
|
handlePush(pn, head.SubType, msg)
|
||||||
default:
|
default:
|
||||||
pn.msgMapMtx.Lock()
|
pn.msgMapMtx.Lock()
|
||||||
ch := pn.msgMap[0]
|
ch := pn.msgMap[0]
|
||||||
@@ -483,7 +480,7 @@ func (pn *P2PNetwork) handleMessage(t int, msg []byte) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) readLoop() {
|
func (pn *P2PNetwork) readLoop() {
|
||||||
gLog.Printf(LevelINFO, "P2PNetwork readLoop start")
|
gLog.Printf(LevelDEBUG, "P2PNetwork readLoop start")
|
||||||
pn.wg.Add(1)
|
pn.wg.Add(1)
|
||||||
defer pn.wg.Done()
|
defer pn.wg.Done()
|
||||||
for pn.running {
|
for pn.running {
|
||||||
@@ -497,7 +494,7 @@ func (pn *P2PNetwork) readLoop() {
|
|||||||
}
|
}
|
||||||
pn.handleMessage(t, msg)
|
pn.handleMessage(t, msg)
|
||||||
}
|
}
|
||||||
gLog.Printf(LevelINFO, "P2PNetwork readLoop end")
|
gLog.Printf(LevelDEBUG, "P2PNetwork readLoop end")
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) write(mainType uint16, subType uint16, packet interface{}) error {
|
func (pn *P2PNetwork) write(mainType uint16, subType uint16, packet interface{}) error {
|
||||||
@@ -567,12 +564,15 @@ func (pn *P2PNetwork) read(node string, mainType uint16, subType uint16, timeout
|
|||||||
} else {
|
} else {
|
||||||
nodeID = nodeNameToID(node)
|
nodeID = nodeNameToID(node)
|
||||||
}
|
}
|
||||||
|
pn.msgMapMtx.Lock()
|
||||||
|
ch := pn.msgMap[nodeID]
|
||||||
|
pn.msgMapMtx.Unlock()
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-time.After(timeout):
|
case <-time.After(timeout):
|
||||||
gLog.Printf(LevelERROR, "wait msg%d:%d timeout", mainType, subType)
|
gLog.Printf(LevelERROR, "wait msg%d:%d timeout", mainType, subType)
|
||||||
return
|
return
|
||||||
case msg := <-pn.msgMap[nodeID]:
|
case msg := <-ch:
|
||||||
head = &openP2PHeader{}
|
head = &openP2PHeader{}
|
||||||
err := binary.Read(bytes.NewReader(msg[:openP2PHeaderSize]), binary.LittleEndian, head)
|
err := binary.Read(bytes.NewReader(msg[:openP2PHeaderSize]), binary.LittleEndian, head)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -592,106 +592,12 @@ func (pn *P2PNetwork) read(node string, mainType uint16, subType uint16, timeout
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pn *P2PNetwork) handlePush(subType uint16, msg []byte) error {
|
|
||||||
pushHead := PushHeader{}
|
|
||||||
err := binary.Read(bytes.NewReader(msg[openP2PHeaderSize:openP2PHeaderSize+PushHeaderSize]), binary.LittleEndian, &pushHead)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
gLog.Printf(LevelDEBUG, "handle push msg type:%d, push header:%+v", subType, pushHead)
|
|
||||||
switch subType {
|
|
||||||
case MsgPushConnectReq:
|
|
||||||
req := PushConnectReq{}
|
|
||||||
err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Printf(LevelERROR, "wrong MsgPushConnectReq:%s", err)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
gLog.Printf(LevelINFO, "%s is connecting...", req.From)
|
|
||||||
gLog.Println(LevelDEBUG, "push connect response to ", req.From)
|
|
||||||
// verify token or name&password
|
|
||||||
if VerifyTOTP(req.Token, pn.config.User, pn.config.Password, time.Now().Unix()+(pn.serverTs-pn.localTs)) || // localTs may behind, auto adjust ts
|
|
||||||
VerifyTOTP(req.Token, pn.config.User, pn.config.Password, time.Now().Unix()) ||
|
|
||||||
(req.User == pn.config.User && req.Password == pn.config.Password) {
|
|
||||||
gLog.Printf(LevelINFO, "Access Granted\n")
|
|
||||||
config := AppConfig{}
|
|
||||||
config.peerNatType = req.NatType
|
|
||||||
config.peerConeNatPort = req.ConeNatPort
|
|
||||||
config.peerIP = req.FromIP
|
|
||||||
config.PeerNode = req.From
|
|
||||||
// share relay node will limit bandwidth
|
|
||||||
if req.User != pn.config.User || req.Password != pn.config.Password {
|
|
||||||
gLog.Printf(LevelINFO, "set share bandwidth %d mbps", pn.config.ShareBandwidth)
|
|
||||||
config.shareBandwidth = pn.config.ShareBandwidth
|
|
||||||
}
|
|
||||||
// go pn.AddTunnel(config, req.ID)
|
|
||||||
go pn.addDirectTunnel(config, req.ID)
|
|
||||||
break
|
|
||||||
}
|
|
||||||
gLog.Println(LevelERROR, "Access Denied:", req.From)
|
|
||||||
rsp := PushConnectRsp{
|
|
||||||
Error: 1,
|
|
||||||
Detail: fmt.Sprintf("connect to %s error: Access Denied", pn.config.Node),
|
|
||||||
To: req.From,
|
|
||||||
From: pn.config.Node,
|
|
||||||
}
|
|
||||||
pn.push(req.From, MsgPushConnectRsp, rsp)
|
|
||||||
case MsgPushRsp:
|
|
||||||
rsp := PushRsp{}
|
|
||||||
err := json.Unmarshal(msg[openP2PHeaderSize:], &rsp)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Printf(LevelERROR, "wrong pushRsp:%s", err)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if rsp.Error == 0 {
|
|
||||||
gLog.Printf(LevelDEBUG, "push ok, detail:%s", rsp.Detail)
|
|
||||||
} else {
|
|
||||||
gLog.Printf(LevelERROR, "push error:%d, detail:%s", rsp.Error, rsp.Detail)
|
|
||||||
}
|
|
||||||
case MsgPushAddRelayTunnelReq:
|
|
||||||
req := AddRelayTunnelReq{}
|
|
||||||
err := json.Unmarshal(msg[openP2PHeaderSize+PushHeaderSize:], &req)
|
|
||||||
if err != nil {
|
|
||||||
gLog.Printf(LevelERROR, "wrong RelayNodeRsp:%s", err)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
config := AppConfig{}
|
|
||||||
config.PeerNode = req.RelayName
|
|
||||||
config.peerToken = req.RelayToken
|
|
||||||
// set user password, maybe the relay node is your private node
|
|
||||||
config.PeerUser = pn.config.User
|
|
||||||
config.PeerPassword = pn.config.Password
|
|
||||||
go func(r AddRelayTunnelReq) {
|
|
||||||
t, errDt := pn.addDirectTunnel(config, 0)
|
|
||||||
if errDt == nil {
|
|
||||||
// notify peer relay ready
|
|
||||||
msg := TunnelMsg{ID: t.id}
|
|
||||||
pn.push(r.From, MsgPushAddRelayTunnelRsp, msg)
|
|
||||||
SaveKey(req.AppID, req.AppKey)
|
|
||||||
}
|
|
||||||
|
|
||||||
}(req)
|
|
||||||
case MsgPushUpdate:
|
|
||||||
update()
|
|
||||||
if gConf.daemonMode {
|
|
||||||
os.Exit(0)
|
|
||||||
}
|
|
||||||
default:
|
|
||||||
pn.msgMapMtx.Lock()
|
|
||||||
ch := pn.msgMap[pushHead.From]
|
|
||||||
pn.msgMapMtx.Unlock()
|
|
||||||
ch <- msg
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (pn *P2PNetwork) updateAppHeartbeat(appID uint64) {
|
func (pn *P2PNetwork) updateAppHeartbeat(appID uint64) {
|
||||||
pn.apps.Range(func(id, i interface{}) bool {
|
pn.apps.Range(func(id, i interface{}) bool {
|
||||||
key := id.(uint64)
|
app := i.(*p2pApp)
|
||||||
if key != appID {
|
if app.id != appID {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
app := i.(*p2pApp)
|
|
||||||
app.updateHeartbeat()
|
app.updateHeartbeat()
|
||||||
return false
|
return false
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -52,13 +52,12 @@ func (t *P2PTunnel) init() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) connect() error {
|
func (t *P2PTunnel) connect() error {
|
||||||
gLog.Printf(LevelINFO, "start p2pTunnel to %s ", t.config.PeerNode)
|
gLog.Printf(LevelDEBUG, "start p2pTunnel to %s ", t.config.PeerNode)
|
||||||
t.isServer = false
|
t.isServer = false
|
||||||
req := PushConnectReq{
|
req := PushConnectReq{
|
||||||
User: t.config.PeerUser,
|
|
||||||
Password: t.config.PeerPassword,
|
|
||||||
Token: t.config.peerToken,
|
Token: t.config.peerToken,
|
||||||
From: t.pn.config.Node,
|
From: t.pn.config.Node,
|
||||||
|
FromToken: t.pn.config.Token,
|
||||||
FromIP: t.pn.config.publicIP,
|
FromIP: t.pn.config.publicIP,
|
||||||
ConeNatPort: t.coneNatPort,
|
ConeNatPort: t.coneNatPort,
|
||||||
NatType: t.pn.config.natType,
|
NatType: t.pn.config.natType,
|
||||||
@@ -144,7 +143,7 @@ func (t *P2PTunnel) handshake() error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
gLog.Println(LevelINFO, "handshake to ", t.config.PeerNode)
|
gLog.Println(LevelDEBUG, "handshake to ", t.config.PeerNode)
|
||||||
var err error
|
var err error
|
||||||
// TODO: handle NATNone, nodes with public ip has no punching
|
// TODO: handle NATNone, nodes with public ip has no punching
|
||||||
if (t.pn.config.natType == NATCone && t.config.peerNatType == NATCone) || (t.pn.config.natType == NATNone || t.config.peerNatType == NATNone) {
|
if (t.pn.config.natType == NATCone && t.config.peerNatType == NATCone) || (t.pn.config.natType == NATNone || t.config.peerNatType == NATNone) {
|
||||||
@@ -163,7 +162,7 @@ func (t *P2PTunnel) handshake() error {
|
|||||||
gLog.Println(LevelERROR, "punch handshake error:", err)
|
gLog.Println(LevelERROR, "punch handshake error:", err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
gLog.Printf(LevelINFO, "handshake to %s ok", t.config.PeerNode)
|
gLog.Printf(LevelDEBUG, "handshake to %s ok", t.config.PeerNode)
|
||||||
err = t.run()
|
err = t.run()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, err)
|
gLog.Println(LevelERROR, err)
|
||||||
@@ -198,7 +197,7 @@ func (t *P2PTunnel) run() error {
|
|||||||
gLog.Println(LevelDEBUG, string(buff))
|
gLog.Println(LevelDEBUG, string(buff))
|
||||||
}
|
}
|
||||||
qConn.WriteBytes(MsgP2P, MsgTunnelHandshakeAck, []byte("OpenP2P,hello2"))
|
qConn.WriteBytes(MsgP2P, MsgTunnelHandshakeAck, []byte("OpenP2P,hello2"))
|
||||||
gLog.Println(LevelINFO, "quic connection ok")
|
gLog.Println(LevelDEBUG, "quic connection ok")
|
||||||
t.conn = qConn
|
t.conn = qConn
|
||||||
t.setRun(true)
|
t.setRun(true)
|
||||||
go t.readLoop()
|
go t.readLoop()
|
||||||
@@ -216,7 +215,7 @@ func (t *P2PTunnel) run() error {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
t.pn.read(t.config.PeerNode, MsgPush, MsgPushQuicConnect, time.Second*5)
|
t.pn.read(t.config.PeerNode, MsgPush, MsgPushQuicConnect, time.Second*5)
|
||||||
gLog.Println(LevelINFO, "quic dial to ", t.ra.String())
|
gLog.Println(LevelDEBUG, "quic dial to ", t.ra.String())
|
||||||
qConn, e := dialQuic(conn, t.ra, TunnelIdleTimeout)
|
qConn, e := dialQuic(conn, t.ra, TunnelIdleTimeout)
|
||||||
if e != nil {
|
if e != nil {
|
||||||
return fmt.Errorf("quic dial to %s error:%s", t.ra.String(), e)
|
return fmt.Errorf("quic dial to %s error:%s", t.ra.String(), e)
|
||||||
@@ -233,7 +232,7 @@ func (t *P2PTunnel) run() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
gLog.Println(LevelINFO, "rtt=", time.Since(handshakeBegin))
|
gLog.Println(LevelINFO, "rtt=", time.Since(handshakeBegin))
|
||||||
gLog.Println(LevelINFO, "quic connection ok")
|
gLog.Println(LevelDEBUG, "quic connection ok")
|
||||||
t.conn = qConn
|
t.conn = qConn
|
||||||
t.setRun(true)
|
t.setRun(true)
|
||||||
go t.readLoop()
|
go t.readLoop()
|
||||||
@@ -243,7 +242,7 @@ func (t *P2PTunnel) run() error {
|
|||||||
|
|
||||||
func (t *P2PTunnel) readLoop() {
|
func (t *P2PTunnel) readLoop() {
|
||||||
decryptData := make([]byte, ReadBuffLen+PaddingSize) // 16 bytes for padding
|
decryptData := make([]byte, ReadBuffLen+PaddingSize) // 16 bytes for padding
|
||||||
gLog.Printf(LevelINFO, "%d tunnel readloop start", t.id)
|
gLog.Printf(LevelDEBUG, "%d tunnel readloop start", t.id)
|
||||||
for t.isRuning() {
|
for t.isRuning() {
|
||||||
t.conn.SetReadDeadline(time.Now().Add(TunnelIdleTimeout))
|
t.conn.SetReadDeadline(time.Now().Add(TunnelIdleTimeout))
|
||||||
head, body, err := t.conn.ReadMessage()
|
head, body, err := t.conn.ReadMessage()
|
||||||
@@ -326,14 +325,14 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
gLog.Printf(LevelERROR, "wrong MsgOverlayConnectReq:%s", err)
|
gLog.Printf(LevelERROR, "wrong MsgOverlayConnectReq:%s", err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
// app connect only accept user/password, avoid someone using the share relay node's token
|
// app connect only accept token(not relay totp token), avoid someone using the share relay node's token
|
||||||
if req.User != t.pn.config.User || req.Password != t.pn.config.Password {
|
if req.Token != t.pn.config.Token {
|
||||||
gLog.Println(LevelERROR, "Access Denied:", req.User)
|
gLog.Println(LevelERROR, "Access Denied:", req.Token)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
overlayID := req.ID
|
overlayID := req.ID
|
||||||
gLog.Printf(LevelINFO, "App:%d overlayID:%d connect %+v", req.AppID, overlayID, req)
|
gLog.Printf(LevelDEBUG, "App:%d overlayID:%d connect %+v", req.AppID, overlayID, req)
|
||||||
if req.Protocol == "tcp" {
|
if req.Protocol == "tcp" {
|
||||||
conn, err := net.DialTimeout("tcp", fmt.Sprintf("%s:%d", req.DstIP, req.DstPort), time.Second*5)
|
conn, err := net.DialTimeout("tcp", fmt.Sprintf("%s:%d", req.DstIP, req.DstPort), time.Second*5)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -368,7 +367,7 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
overlayID := req.ID
|
overlayID := req.ID
|
||||||
gLog.Printf(LevelINFO, "%d disconnect overlay connection %d", t.id, overlayID)
|
gLog.Printf(LevelDEBUG, "%d disconnect overlay connection %d", t.id, overlayID)
|
||||||
i, ok := t.overlayConns.Load(overlayID)
|
i, ok := t.overlayConns.Load(overlayID)
|
||||||
if ok {
|
if ok {
|
||||||
otcp := i.(*overlayTCP)
|
otcp := i.(*overlayTCP)
|
||||||
@@ -379,13 +378,13 @@ func (t *P2PTunnel) readLoop() {
|
|||||||
}
|
}
|
||||||
t.setRun(false)
|
t.setRun(false)
|
||||||
t.conn.Close()
|
t.conn.Close()
|
||||||
gLog.Printf(LevelINFO, "%d tunnel readloop end", t.id)
|
gLog.Printf(LevelDEBUG, "%d tunnel readloop end", t.id)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) writeLoop() {
|
func (t *P2PTunnel) writeLoop() {
|
||||||
tc := time.NewTicker(TunnelHeartbeatTime)
|
tc := time.NewTicker(TunnelHeartbeatTime)
|
||||||
defer tc.Stop()
|
defer tc.Stop()
|
||||||
defer gLog.Printf(LevelINFO, "%d tunnel writeloop end", t.id)
|
defer gLog.Printf(LevelDEBUG, "%d tunnel writeloop end", t.id)
|
||||||
for t.isRuning() {
|
for t.isRuning() {
|
||||||
select {
|
select {
|
||||||
case <-tc.C:
|
case <-tc.C:
|
||||||
@@ -402,7 +401,7 @@ func (t *P2PTunnel) writeLoop() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (t *P2PTunnel) listen() error {
|
func (t *P2PTunnel) listen() error {
|
||||||
gLog.Printf(LevelINFO, "p2ptunnel wait for connecting")
|
gLog.Printf(LevelDEBUG, "p2ptunnel wait for connecting")
|
||||||
t.isServer = true
|
t.isServer = true
|
||||||
return t.handshake()
|
return t.handshake()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -10,7 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
const OpenP2PVersion = "0.97.1"
|
const OpenP2PVersion = "1.0.0"
|
||||||
const ProducnName string = "openp2p"
|
const ProducnName string = "openp2p"
|
||||||
|
|
||||||
type openP2PHeader struct {
|
type openP2PHeader struct {
|
||||||
@@ -79,6 +79,10 @@ const (
|
|||||||
MsgPushUpdate = 6
|
MsgPushUpdate = 6
|
||||||
MsgPushReportApps = 7
|
MsgPushReportApps = 7
|
||||||
MsgPushQuicConnect = 8
|
MsgPushQuicConnect = 8
|
||||||
|
MsgPushEditApp = 9
|
||||||
|
MsgPushSwitchApp = 10
|
||||||
|
MsgPushRestart = 11
|
||||||
|
MsgPushEditNode = 12
|
||||||
)
|
)
|
||||||
|
|
||||||
// MsgP2P sub type message
|
// MsgP2P sub type message
|
||||||
@@ -109,6 +113,7 @@ const (
|
|||||||
MsgReportBasic = iota
|
MsgReportBasic = iota
|
||||||
MsgReportQuery
|
MsgReportQuery
|
||||||
MsgReportConnect
|
MsgReportConnect
|
||||||
|
MsgReportApps
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -128,6 +133,7 @@ const (
|
|||||||
RetryInterval = time.Second * 30
|
RetryInterval = time.Second * 30
|
||||||
PublicIPEchoTimeout = time.Second * 3
|
PublicIPEchoTimeout = time.Second * 3
|
||||||
NatTestTimeout = time.Second * 10
|
NatTestTimeout = time.Second * 10
|
||||||
|
ClientAPITimeout = time.Second * 10
|
||||||
)
|
)
|
||||||
|
|
||||||
// NATNone has public ip
|
// NATNone has public ip
|
||||||
@@ -135,6 +141,7 @@ const (
|
|||||||
NATNone = 0
|
NATNone = 0
|
||||||
NATCone = 1
|
NATCone = 1
|
||||||
NATSymmetric = 2
|
NATSymmetric = 2
|
||||||
|
NATUnknown = 314
|
||||||
)
|
)
|
||||||
|
|
||||||
func newMessage(mainType uint16, subType uint16, packet interface{}) ([]byte, error) {
|
func newMessage(mainType uint16, subType uint16, packet interface{}) ([]byte, error) {
|
||||||
@@ -163,9 +170,8 @@ func nodeNameToID(name string) uint64 {
|
|||||||
|
|
||||||
type PushConnectReq struct {
|
type PushConnectReq struct {
|
||||||
From string `json:"from,omitempty"`
|
From string `json:"from,omitempty"`
|
||||||
User string `json:"user,omitempty"`
|
FromToken uint64 `json:"fromToken,omitempty"` //my token
|
||||||
Password string `json:"password,omitempty"`
|
Token uint64 `json:"token,omitempty"` // totp token
|
||||||
Token uint64 `json:"token,omitempty"`
|
|
||||||
ConeNatPort int `json:"coneNatPort,omitempty"`
|
ConeNatPort int `json:"coneNatPort,omitempty"`
|
||||||
NatType int `json:"natType,omitempty"`
|
NatType int `json:"natType,omitempty"`
|
||||||
FromIP string `json:"fromIP,omitempty"`
|
FromIP string `json:"fromIP,omitempty"`
|
||||||
@@ -189,6 +195,8 @@ type PushRsp struct {
|
|||||||
type LoginRsp struct {
|
type LoginRsp struct {
|
||||||
Error int `json:"error,omitempty"`
|
Error int `json:"error,omitempty"`
|
||||||
Detail string `json:"detail,omitempty"`
|
Detail string `json:"detail,omitempty"`
|
||||||
|
User string `json:"user,omitempty"`
|
||||||
|
Token uint64 `json:"token,omitempty"`
|
||||||
Ts int64 `json:"ts,omitempty"`
|
Ts int64 `json:"ts,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -209,8 +217,7 @@ type P2PHandshakeReq struct {
|
|||||||
|
|
||||||
type OverlayConnectReq struct {
|
type OverlayConnectReq struct {
|
||||||
ID uint64 `json:"id,omitempty"`
|
ID uint64 `json:"id,omitempty"`
|
||||||
User string `json:"user,omitempty"`
|
Token uint64 `json:"token,omitempty"` // not totp token
|
||||||
Password string `json:"password,omitempty"`
|
|
||||||
DstIP string `json:"dstIP,omitempty"`
|
DstIP string `json:"dstIP,omitempty"`
|
||||||
DstPort int `json:"dstPort,omitempty"`
|
DstPort int `json:"dstPort,omitempty"`
|
||||||
Protocol string `json:"protocol,omitempty"`
|
Protocol string `json:"protocol,omitempty"`
|
||||||
@@ -271,6 +278,32 @@ type ReportConnect struct {
|
|||||||
Version string `json:"version,omitempty"`
|
Version string `json:"version,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type AppInfo struct {
|
||||||
|
AppName string `json:"appName,omitempty"`
|
||||||
|
Error string `json:"error,omitempty"`
|
||||||
|
Protocol string `json:"protocol,omitempty"`
|
||||||
|
SrcPort int `json:"srcPort,omitempty"`
|
||||||
|
Protocol0 string `json:"protocol0,omitempty"`
|
||||||
|
SrcPort0 int `json:"srcPort0,omitempty"`
|
||||||
|
NatType int `json:"natType,omitempty"`
|
||||||
|
PeerNode string `json:"peerNode,omitempty"`
|
||||||
|
DstPort int `json:"dstPort,omitempty"`
|
||||||
|
DstHost string `json:"dstHost,omitempty"`
|
||||||
|
PeerUser string `json:"peerUser,omitempty"`
|
||||||
|
PeerNatType int `json:"peerNatType,omitempty"`
|
||||||
|
PeerIP string `json:"peerIP,omitempty"`
|
||||||
|
ShareBandwidth int `json:"shareBandWidth,omitempty"`
|
||||||
|
RelayNode string `json:"relayNode,omitempty"`
|
||||||
|
Version string `json:"version,omitempty"`
|
||||||
|
RetryTime string `json:"retryTime,omitempty"`
|
||||||
|
IsActive int `json:"isActive,omitempty"`
|
||||||
|
Enabled int `json:"enabled,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type ReportApps struct {
|
||||||
|
Apps []AppInfo
|
||||||
|
}
|
||||||
|
|
||||||
type UpdateInfo struct {
|
type UpdateInfo struct {
|
||||||
Error int `json:"error,omitempty"`
|
Error int `json:"error,omitempty"`
|
||||||
ErrorDetail string `json:"errorDetail,omitempty"`
|
ErrorDetail string `json:"errorDetail,omitempty"`
|
||||||
@@ -295,3 +328,17 @@ type NetInfo struct {
|
|||||||
ASNOrg string `json:"asn_org,omitempty"`
|
ASNOrg string `json:"asn_org,omitempty"`
|
||||||
Hostname string `json:"hostname,omitempty"`
|
Hostname string `json:"hostname,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type ProfileInfo struct {
|
||||||
|
User string `json:"user,omitempty"`
|
||||||
|
Password string `json:"password,omitempty"`
|
||||||
|
Email string `json:"email,omitempty"`
|
||||||
|
Phone string `json:"phone,omitempty"`
|
||||||
|
Token string `json:"token,omitempty"`
|
||||||
|
Addtime string `json:"addtime,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type EditNode struct {
|
||||||
|
NewName string `json:"newName,omitempty"`
|
||||||
|
Bandwidth int `json:"bandwidth,omitempty"`
|
||||||
|
}
|
||||||
|
|||||||
@@ -99,7 +99,7 @@ func (conn *quicConn) Accept() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func listenQuic(addr string, idleTimeout time.Duration) (*quicConn, error) {
|
func listenQuic(addr string, idleTimeout time.Duration) (*quicConn, error) {
|
||||||
gLog.Println(LevelINFO, "quic listen on ", addr)
|
gLog.Println(LevelDEBUG, "quic listen on ", addr)
|
||||||
listener, err := quic.ListenAddr(addr, generateTLSConfig(),
|
listener, err := quic.ListenAddr(addr, generateTLSConfig(),
|
||||||
&quic.Config{Versions: quicVersion, MaxIdleTimeout: idleTimeout, DisablePathMTUDiscovery: true})
|
&quic.Config{Versions: quicVersion, MaxIdleTimeout: idleTimeout, DisablePathMTUDiscovery: true})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -8,9 +8,11 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const TOTPStep = 30 // 30s
|
const TOTPStep = 30 // 30s
|
||||||
func GenTOTP(user string, password string, ts int64) uint64 {
|
func GenTOTP(token uint64, ts int64) uint64 {
|
||||||
step := ts / TOTPStep
|
step := ts / TOTPStep
|
||||||
mac := hmac.New(sha256.New, []byte(user+password))
|
tbuff := make([]byte, 8)
|
||||||
|
binary.LittleEndian.PutUint64(tbuff, token)
|
||||||
|
mac := hmac.New(sha256.New, tbuff)
|
||||||
b := make([]byte, 8)
|
b := make([]byte, 8)
|
||||||
binary.LittleEndian.PutUint64(b, uint64(step))
|
binary.LittleEndian.PutUint64(b, uint64(step))
|
||||||
mac.Write(b)
|
mac.Write(b)
|
||||||
@@ -19,11 +21,11 @@ func GenTOTP(user string, password string, ts int64) uint64 {
|
|||||||
return num
|
return num
|
||||||
}
|
}
|
||||||
|
|
||||||
func VerifyTOTP(code uint64, user string, password string, ts int64) bool {
|
func VerifyTOTP(code uint64, token uint64, ts int64) bool {
|
||||||
if code == 0 {
|
if code == 0 {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
if code == GenTOTP(user, password, ts) || code == GenTOTP(user, password, ts-TOTPStep) || code == GenTOTP(user, password, ts+TOTPStep) {
|
if code == GenTOTP(token, ts) || code == GenTOTP(token, ts-TOTPStep) || code == GenTOTP(token, ts+TOTPStep) {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
|
|||||||
@@ -9,24 +9,24 @@ import (
|
|||||||
func TestTOTP(t *testing.T) {
|
func TestTOTP(t *testing.T) {
|
||||||
for i := 0; i < 20; i++ {
|
for i := 0; i < 20; i++ {
|
||||||
ts := time.Now().Unix()
|
ts := time.Now().Unix()
|
||||||
code := GenTOTP("testuser1", "testpassword1", ts)
|
code := GenTOTP(13666999958022769123, ts)
|
||||||
t.Log(code)
|
t.Log(code)
|
||||||
if !VerifyTOTP(code, "testuser1", "testpassword1", ts) {
|
if !VerifyTOTP(code, 13666999958022769123, ts) {
|
||||||
t.Error("TOTP error")
|
t.Error("TOTP error")
|
||||||
}
|
}
|
||||||
if !VerifyTOTP(code, "testuser1", "testpassword1", ts-10) {
|
if !VerifyTOTP(code, 13666999958022769123, ts-10) {
|
||||||
t.Error("TOTP error")
|
t.Error("TOTP error")
|
||||||
}
|
}
|
||||||
if !VerifyTOTP(code, "testuser1", "testpassword1", ts+10) {
|
if !VerifyTOTP(code, 13666999958022769123, ts+10) {
|
||||||
t.Error("TOTP error")
|
t.Error("TOTP error")
|
||||||
}
|
}
|
||||||
if VerifyTOTP(code, "testuser1", "testpassword1", ts+60) {
|
if VerifyTOTP(code, 13666999958022769123, ts+60) {
|
||||||
t.Error("TOTP error")
|
t.Error("TOTP error")
|
||||||
}
|
}
|
||||||
if VerifyTOTP(code, "testuser2", "testpassword1", ts+1) {
|
if VerifyTOTP(code, 13666999958022769124, ts+1) {
|
||||||
t.Error("TOTP error")
|
t.Error("TOTP error")
|
||||||
}
|
}
|
||||||
if VerifyTOTP(code, "testuser1", "testpassword2", ts+1) {
|
if VerifyTOTP(code, 13666999958022769125, ts+1) {
|
||||||
t.Error("TOTP error")
|
t.Error("TOTP error")
|
||||||
}
|
}
|
||||||
time.Sleep(time.Second)
|
time.Sleep(time.Second)
|
||||||
|
|||||||
@@ -16,18 +16,9 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
// type updateFileInfo struct {
|
|
||||||
// Name string `json:"name,omitempty"`
|
|
||||||
// RelativePath string `json:"relativePath,omitempty"`
|
|
||||||
// Length int64 `json:"length,omitempty"`
|
|
||||||
// URL string `json:"url,omitempty"`
|
|
||||||
// Hash string `json:"hash,omitempty"`
|
|
||||||
// }
|
|
||||||
|
|
||||||
func update() {
|
func update() {
|
||||||
gLog.Println(LevelINFO, "update start")
|
gLog.Println(LevelINFO, "update start")
|
||||||
defer gLog.Println(LevelINFO, "update end")
|
defer gLog.Println(LevelINFO, "update end")
|
||||||
// TODO: download from gitee. save flow
|
|
||||||
c := http.Client{
|
c := http.Client{
|
||||||
Transport: &http.Transport{
|
Transport: &http.Transport{
|
||||||
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
|
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
|
||||||
@@ -36,7 +27,7 @@ func update() {
|
|||||||
}
|
}
|
||||||
goos := runtime.GOOS
|
goos := runtime.GOOS
|
||||||
goarch := runtime.GOARCH
|
goarch := runtime.GOARCH
|
||||||
rsp, err := c.Get(fmt.Sprintf("https://openp2p.cn:27182/api/v1/update?fromver=%s&os=%s&arch=%s", OpenP2PVersion, goos, goarch))
|
rsp, err := c.Get(fmt.Sprintf("https://openp2p.cn:27183/api/v1/update?fromver=%s&os=%s&arch=%s", OpenP2PVersion, goos, goarch))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "update:query update list failed:", err)
|
gLog.Println(LevelERROR, "update:query update list failed:", err)
|
||||||
return
|
return
|
||||||
@@ -61,7 +52,6 @@ func update() {
|
|||||||
gLog.Println(LevelERROR, "update error:", updateInfo.Error, updateInfo.ErrorDetail)
|
gLog.Println(LevelERROR, "update error:", updateInfo.Error, updateInfo.ErrorDetail)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
os.MkdirAll("download", 0666)
|
|
||||||
err = updateFile(updateInfo.Url, "", "openp2p")
|
err = updateFile(updateInfo.Url, "", "openp2p")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
gLog.Println(LevelERROR, "update: download failed:", err)
|
gLog.Println(LevelERROR, "update: download failed:", err)
|
||||||
@@ -112,6 +102,7 @@ func updateFile(url string, checksum string, dst string) error {
|
|||||||
os.Rename(os.Args[0]+"0", os.Args[0])
|
os.Rename(os.Args[0]+"0", os.Args[0])
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
os.Remove(tmpFile)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -133,11 +124,6 @@ func unzip(dst, src string) (err error) {
|
|||||||
for _, f := range archive.File {
|
for _, f := range archive.File {
|
||||||
filePath := filepath.Join(dst, f.Name)
|
filePath := filepath.Join(dst, f.Name)
|
||||||
fmt.Println("unzipping file ", filePath)
|
fmt.Println("unzipping file ", filePath)
|
||||||
|
|
||||||
// if !strings.HasPrefix(filePath, filepath.Clean(dst)+string(os.PathSeparator)) {
|
|
||||||
// fmt.Println("invalid file path")
|
|
||||||
// return
|
|
||||||
// }
|
|
||||||
if f.FileInfo().IsDir() {
|
if f.FileInfo().IsDir() {
|
||||||
fmt.Println("creating directory...")
|
fmt.Println("creating directory...")
|
||||||
os.MkdirAll(filePath, os.ModePerm)
|
os.MkdirAll(filePath, os.ModePerm)
|
||||||
|
|||||||