mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-29 16:06:36 +08:00
增加出口负载均衡
This commit is contained in:
@@ -17,6 +17,8 @@ public class ForwardDto {
|
||||
|
||||
@NotBlank(message = "远程地址不能为空")
|
||||
private String remoteAddr;
|
||||
|
||||
private String strategy;
|
||||
|
||||
/**
|
||||
* 入口端口(可选,为空时自动分配)
|
||||
|
||||
@@ -23,6 +23,8 @@ public class ForwardUpdateDto {
|
||||
|
||||
@NotBlank(message = "远程地址不能为空")
|
||||
private String remoteAddr;
|
||||
|
||||
private String strategy;
|
||||
|
||||
/**
|
||||
* 入口端口(可选,为空时自动分配)
|
||||
|
||||
@@ -84,6 +84,7 @@ public class ForwardWithTunnelDto {
|
||||
*/
|
||||
private Long outFlow;
|
||||
|
||||
private String strategy;
|
||||
// /**
|
||||
// * 入口端口开始
|
||||
// */
|
||||
|
||||
@@ -31,21 +31,21 @@ public class GostUtil {
|
||||
return WebSocketServer.send_msg(node_id, req, "DeleteLimiters");
|
||||
}
|
||||
|
||||
public static GostDto AddService(Long node_id, String name, Integer in_port, Integer limiter, String remoteAddr, Integer fow_type, Tunnel tunnel) {
|
||||
public static GostDto AddService(Long node_id, String name, Integer in_port, Integer limiter, String remoteAddr, Integer fow_type, Tunnel tunnel, String strategy) {
|
||||
JSONArray services = new JSONArray();
|
||||
String[] protocols = {"tcp", "udp"};
|
||||
for (String protocol : protocols) {
|
||||
JSONObject service = createServiceConfig(name, in_port, limiter, remoteAddr, protocol, fow_type, tunnel);
|
||||
JSONObject service = createServiceConfig(name, in_port, limiter, remoteAddr, protocol, fow_type, tunnel, strategy);
|
||||
services.add(service);
|
||||
}
|
||||
return WebSocketServer.send_msg(node_id, services, "AddService");
|
||||
}
|
||||
|
||||
public static GostDto UpdateService(Long node_id, String name, Integer in_port, Integer limiter, String remoteAddr, Integer fow_type, Tunnel tunnel) {
|
||||
public static GostDto UpdateService(Long node_id, String name, Integer in_port, Integer limiter, String remoteAddr, Integer fow_type, Tunnel tunnel, String strategy) {
|
||||
JSONArray services = new JSONArray();
|
||||
String[] protocols = {"tcp", "udp"};
|
||||
for (String protocol : protocols) {
|
||||
JSONObject service = createServiceConfig(name, in_port, limiter, remoteAddr, protocol, fow_type, tunnel);
|
||||
JSONObject service = createServiceConfig(name, in_port, limiter, remoteAddr, protocol, fow_type, tunnel, strategy);
|
||||
services.add(service);
|
||||
}
|
||||
return WebSocketServer.send_msg(node_id, services, "UpdateService");
|
||||
@@ -60,6 +60,90 @@ public class GostUtil {
|
||||
return WebSocketServer.send_msg(node_id, data, "DeleteService");
|
||||
}
|
||||
|
||||
public static GostDto AddRemoteService(Long node_id, String name, Integer out_port, String remoteAddr, String protocol, String strategy) {
|
||||
JSONObject data = new JSONObject();
|
||||
data.put("name", name + "_tls");
|
||||
data.put("addr", ":" + out_port);
|
||||
JSONObject handler = new JSONObject();
|
||||
handler.put("type", "relay");
|
||||
data.put("handler", handler);
|
||||
JSONObject listener = new JSONObject();
|
||||
listener.put("type", protocol);
|
||||
data.put("listener", listener);
|
||||
JSONObject forwarder = new JSONObject();
|
||||
JSONArray nodes = new JSONArray();
|
||||
|
||||
String[] split = remoteAddr.split(",");
|
||||
int num = 1;
|
||||
for (String addr : split) {
|
||||
JSONObject node = new JSONObject();
|
||||
node.put("name", "node_" + num );
|
||||
node.put("addr", addr);
|
||||
nodes.add(node);
|
||||
num ++;
|
||||
}
|
||||
if (strategy == null || strategy.equals("")){
|
||||
strategy = "fifo";
|
||||
}
|
||||
forwarder.put("nodes", nodes);
|
||||
JSONObject selector = new JSONObject();
|
||||
selector.put("strategy", strategy);
|
||||
selector.put("maxFails", 1);
|
||||
selector.put("failTimeout", "10s");
|
||||
forwarder.put("selector", selector);
|
||||
|
||||
data.put("forwarder", forwarder);
|
||||
JSONArray services = new JSONArray();
|
||||
services.add(data);
|
||||
return WebSocketServer.send_msg(node_id, services, "AddService");
|
||||
}
|
||||
|
||||
public static GostDto UpdateRemoteService(Long node_id, String name, Integer out_port, String remoteAddr,String protocol, String strategy) {
|
||||
JSONObject data = new JSONObject();
|
||||
data.put("name", name + "_tls");
|
||||
data.put("addr", ":" + out_port);
|
||||
JSONObject handler = new JSONObject();
|
||||
handler.put("type", "relay");
|
||||
data.put("handler", handler);
|
||||
JSONObject listener = new JSONObject();
|
||||
listener.put("type", protocol);
|
||||
data.put("listener", listener);
|
||||
JSONObject forwarder = new JSONObject();
|
||||
JSONArray nodes = new JSONArray();
|
||||
|
||||
String[] split = remoteAddr.split(",");
|
||||
int num = 1;
|
||||
for (String addr : split) {
|
||||
JSONObject node = new JSONObject();
|
||||
node.put("name", "node_" + num );
|
||||
node.put("addr", addr);
|
||||
nodes.add(node);
|
||||
num ++;
|
||||
}
|
||||
if (strategy == null || strategy.equals("")){
|
||||
strategy = "fifo";
|
||||
}
|
||||
forwarder.put("nodes", nodes);
|
||||
JSONObject selector = new JSONObject();
|
||||
selector.put("strategy", strategy);
|
||||
selector.put("maxFails", 1);
|
||||
selector.put("failTimeout", "10s");
|
||||
forwarder.put("selector", selector);
|
||||
|
||||
data.put("forwarder", forwarder);
|
||||
JSONArray services = new JSONArray();
|
||||
services.add(data);
|
||||
return WebSocketServer.send_msg(node_id, services, "UpdateService");
|
||||
}
|
||||
|
||||
public static GostDto DeleteRemoteService(Long node_id, String name) {
|
||||
JSONArray data = new JSONArray();
|
||||
data.add(name + "_tls");
|
||||
JSONObject req = new JSONObject();
|
||||
req.put("services", data);
|
||||
return WebSocketServer.send_msg(node_id, req, "DeleteService");
|
||||
}
|
||||
|
||||
public static GostDto PauseService(Long node_id, String name) {
|
||||
JSONObject data = new JSONObject();
|
||||
JSONArray services = new JSONArray();
|
||||
@@ -162,60 +246,6 @@ public class GostUtil {
|
||||
return WebSocketServer.send_msg(node_id, data, "DeleteChains");
|
||||
}
|
||||
|
||||
public static GostDto AddRemoteService(Long node_id, String name, Integer out_port, String remoteAddr, String protocol) {
|
||||
JSONObject data = new JSONObject();
|
||||
data.put("name", name + "_tls");
|
||||
data.put("addr", ":" + out_port);
|
||||
JSONObject handler = new JSONObject();
|
||||
handler.put("type", "relay");
|
||||
data.put("handler", handler);
|
||||
JSONObject listener = new JSONObject();
|
||||
listener.put("type", protocol);
|
||||
data.put("listener", listener);
|
||||
JSONObject forwarder = new JSONObject();
|
||||
JSONArray nodes = new JSONArray();
|
||||
JSONObject node = new JSONObject();
|
||||
node.put("name", name + "_node");
|
||||
node.put("addr", remoteAddr);
|
||||
nodes.add(node);
|
||||
forwarder.put("nodes", nodes);
|
||||
data.put("forwarder", forwarder);
|
||||
JSONArray services = new JSONArray();
|
||||
services.add(data);
|
||||
return WebSocketServer.send_msg(node_id, services, "AddService");
|
||||
}
|
||||
|
||||
public static GostDto UpdateRemoteService(Long node_id, String name, Integer out_port, String remoteAddr) {
|
||||
JSONObject data = new JSONObject();
|
||||
data.put("name", name + "_tls");
|
||||
data.put("addr", ":" + out_port);
|
||||
JSONObject handler = new JSONObject();
|
||||
handler.put("type", "relay");
|
||||
data.put("handler", handler);
|
||||
JSONObject listener = new JSONObject();
|
||||
listener.put("type", "tls");
|
||||
data.put("listener", listener);
|
||||
JSONObject forwarder = new JSONObject();
|
||||
JSONArray nodes = new JSONArray();
|
||||
JSONObject node = new JSONObject();
|
||||
node.put("name", name + "_node");
|
||||
node.put("addr", remoteAddr);
|
||||
nodes.add(node);
|
||||
forwarder.put("nodes", nodes);
|
||||
data.put("forwarder", forwarder);
|
||||
JSONArray services = new JSONArray();
|
||||
services.add(data);
|
||||
return WebSocketServer.send_msg(node_id, services, "UpdateService");
|
||||
}
|
||||
|
||||
public static GostDto DeleteRemoteService(Long node_id, String name) {
|
||||
JSONArray data = new JSONArray();
|
||||
data.add(name + "_tls");
|
||||
JSONObject req = new JSONObject();
|
||||
req.put("services", data);
|
||||
return WebSocketServer.send_msg(node_id, req, "DeleteService");
|
||||
}
|
||||
|
||||
private static JSONObject createLimiterData(Long name, String speed) {
|
||||
JSONObject data = new JSONObject();
|
||||
data.put("name", name.toString());
|
||||
@@ -225,7 +255,7 @@ public class GostUtil {
|
||||
return data;
|
||||
}
|
||||
|
||||
private static JSONObject createServiceConfig(String name, Integer in_port, Integer limiter, String remoteAddr, String protocol, Integer fow_type, Tunnel tunnel) {
|
||||
private static JSONObject createServiceConfig(String name, Integer in_port, Integer limiter, String remoteAddr, String protocol, Integer fow_type, Tunnel tunnel, String strategy) {
|
||||
JSONObject service = new JSONObject();
|
||||
service.put("name", name + "_" + protocol);
|
||||
if (Objects.equals(protocol, "tcp")){
|
||||
@@ -249,7 +279,7 @@ public class GostUtil {
|
||||
|
||||
// 端口转发需要配置转发器
|
||||
if (isPortForwarding(fow_type)) {
|
||||
JSONObject forwarder = createForwarder(protocol, remoteAddr);
|
||||
JSONObject forwarder = createForwarder(protocol, remoteAddr, strategy);
|
||||
service.put("forwarder", forwarder);
|
||||
}
|
||||
|
||||
@@ -274,14 +304,31 @@ public class GostUtil {
|
||||
return listener;
|
||||
}
|
||||
|
||||
private static JSONObject createForwarder(String protocol, String remoteAddr) {
|
||||
private static JSONObject createForwarder(String protocol, String remoteAddr, String strategy) {
|
||||
JSONObject forwarder = new JSONObject();
|
||||
JSONArray nodes = new JSONArray();
|
||||
JSONObject node = new JSONObject();
|
||||
node.put("name", protocol);
|
||||
node.put("addr", remoteAddr);
|
||||
nodes.add(node);
|
||||
|
||||
String[] split = remoteAddr.split(",");
|
||||
int num = 1;
|
||||
for (String addr : split) {
|
||||
JSONObject node = new JSONObject();
|
||||
node.put("name", "node_" + num );
|
||||
node.put("addr", addr);
|
||||
nodes.add(node);
|
||||
num ++;
|
||||
}
|
||||
|
||||
if (strategy == null || strategy.equals("")){
|
||||
strategy = "fifo";
|
||||
}
|
||||
|
||||
forwarder.put("nodes", nodes);
|
||||
|
||||
JSONObject selector = new JSONObject();
|
||||
selector.put("strategy", strategy);
|
||||
selector.put("maxFails", 1);
|
||||
selector.put("failTimeout", "10s");
|
||||
forwarder.put("selector", selector);
|
||||
return forwarder;
|
||||
}
|
||||
|
||||
|
||||
@@ -32,6 +32,8 @@ public class Forward extends BaseEntity{
|
||||
|
||||
private String remoteAddr;
|
||||
|
||||
private String strategy;
|
||||
|
||||
private Long inFlow;
|
||||
|
||||
private Long outFlow;
|
||||
|
||||
@@ -659,7 +659,7 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
|
||||
}
|
||||
|
||||
// 创建主服务
|
||||
R serviceResult = createMainService(nodeInfo.getInNode(), serviceName, forward, limiter, tunnel.getType(), tunnel);
|
||||
R serviceResult = createMainService(nodeInfo.getInNode(), serviceName, forward, limiter, tunnel.getType(), tunnel, forward.getStrategy());
|
||||
if (serviceResult.getCode() != 0) {
|
||||
GostUtil.DeleteChains(nodeInfo.getInNode().getId(), serviceName);
|
||||
if (nodeInfo.getOutNode() != null) {
|
||||
@@ -693,7 +693,7 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
|
||||
}
|
||||
|
||||
// 更新主服务
|
||||
R serviceResult = updateMainService(nodeInfo.getInNode(), serviceName, forward, limiter, tunnel.getType(), tunnel);
|
||||
R serviceResult = updateMainService(nodeInfo.getInNode(), serviceName, forward, limiter, tunnel.getType(), tunnel, forward.getStrategy());
|
||||
if (serviceResult.getCode() != 0) {
|
||||
updateForwardStatusToError(forward);
|
||||
return serviceResult;
|
||||
@@ -825,18 +825,15 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
|
||||
* 创建远程服务
|
||||
*/
|
||||
private R createRemoteService(Node outNode, String serviceName, Forward forward, String protocol) {
|
||||
GostDto result = GostUtil.AddRemoteService(outNode.getId(),
|
||||
serviceName, forward.getOutPort(),
|
||||
forward.getRemoteAddr(), protocol);
|
||||
GostDto result = GostUtil.AddRemoteService(outNode.getId(), serviceName, forward.getOutPort(), forward.getRemoteAddr(), protocol, forward.getStrategy());
|
||||
return isGostOperationSuccess(result) ? R.ok() : R.err(result.getMsg());
|
||||
}
|
||||
|
||||
/**
|
||||
* 创建主服务
|
||||
*/
|
||||
private R createMainService(Node inNode, String serviceName, Forward forward, Integer limiter, Integer tunnelType, Tunnel tunnel) {
|
||||
GostDto result = GostUtil.AddService(inNode.getId(), serviceName,
|
||||
forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel);
|
||||
private R createMainService(Node inNode, String serviceName, Forward forward, Integer limiter, Integer tunnelType, Tunnel tunnel, String strategy) {
|
||||
GostDto result = GostUtil.AddService(inNode.getId(), serviceName, forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel, strategy);
|
||||
return isGostOperationSuccess(result) ? R.ok() : R.err(result.getMsg());
|
||||
}
|
||||
|
||||
@@ -863,11 +860,11 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
|
||||
// 创建新远程服务
|
||||
GostDto createResult = GostUtil.UpdateRemoteService(outNode.getId(),
|
||||
serviceName, forward.getOutPort(),
|
||||
forward.getRemoteAddr());
|
||||
forward.getRemoteAddr(), protocol, forward.getStrategy());
|
||||
if (createResult.getMsg().contains(GOST_NOT_FOUND_MSG)) {
|
||||
createResult = GostUtil.AddRemoteService(outNode.getId(),
|
||||
serviceName, forward.getOutPort(),
|
||||
forward.getRemoteAddr(),protocol);
|
||||
forward.getRemoteAddr(),protocol, forward.getStrategy());
|
||||
}
|
||||
return isGostOperationSuccess(createResult) ? R.ok() : R.err(createResult.getMsg());
|
||||
}
|
||||
@@ -875,14 +872,11 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
|
||||
/**
|
||||
* 更新主服务
|
||||
*/
|
||||
private R updateMainService(Node inNode, String serviceName, Forward forward, Integer limiter, Integer tunnelType, Tunnel tunnel) {
|
||||
GostDto result = GostUtil.UpdateService(inNode.getId(), serviceName,
|
||||
forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel);
|
||||
private R updateMainService(Node inNode, String serviceName, Forward forward, Integer limiter, Integer tunnelType, Tunnel tunnel, String strategy) {
|
||||
GostDto result = GostUtil.UpdateService(inNode.getId(), serviceName, forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel, strategy);
|
||||
|
||||
if (result.getMsg().contains(GOST_NOT_FOUND_MSG)) {
|
||||
result = GostUtil.AddService(inNode.getId(), serviceName,
|
||||
forward.getInPort(), limiter, forward.getRemoteAddr(),
|
||||
tunnelType, tunnel);
|
||||
result = GostUtil.AddService(inNode.getId(), serviceName, forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel, strategy);
|
||||
}
|
||||
|
||||
return isGostOperationSuccess(result) ? R.ok() : R.err(result.getMsg());
|
||||
|
||||
@@ -727,7 +727,7 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
|
||||
// 3. 根据隧道类型执行不同的诊断策略
|
||||
if (tunnel.getType() == TUNNEL_TYPE_PORT_FORWARD) {
|
||||
// 端口转发:只给入口节点发送诊断指令,ping谷歌DNS
|
||||
DiagnosisResult inResult = performPingDiagnosisWithConnectionCheck(inNode, "8.8.8.8", "入口->外网");
|
||||
DiagnosisResult inResult = performPingDiagnosisWithConnectionCheck(inNode, "www.google.com", "入口->外网");
|
||||
results.add(inResult);
|
||||
} else {
|
||||
// 隧道转发:入口ping出口,出口ping谷歌DNS
|
||||
@@ -735,7 +735,7 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
|
||||
results.add(inToOutResult);
|
||||
|
||||
// 先检查出口节点的真实连接状态,然后再进行诊断
|
||||
DiagnosisResult outToExternalResult = performPingDiagnosisWithConnectionCheck(outNode, "8.8.8.8", "出口->外网");
|
||||
DiagnosisResult outToExternalResult = performPingDiagnosisWithConnectionCheck(outNode, "www.google.com", "出口->外网");
|
||||
results.add(outToExternalResult);
|
||||
}
|
||||
|
||||
|
||||
@@ -499,8 +499,7 @@ public class UserTunnelServiceImpl extends ServiceImpl<UserTunnelMapper, UserTun
|
||||
String serviceName = buildServiceName(forward.getId(), Long.valueOf(userId), userTunnel.getId());
|
||||
|
||||
// 6. 更新入口节点的主服务限速配置(使用批量UpdateService接口)
|
||||
GostUtil.UpdateService(inNode.getId(), serviceName, forward.getInPort(), speedId,
|
||||
forward.getRemoteAddr(), tunnel.getType(), tunnel);
|
||||
GostUtil.UpdateService(inNode.getId(), serviceName, forward.getInPort(), speedId, forward.getRemoteAddr(), tunnel.getType(), tunnel, forward.getStrategy());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user