diff --git a/.classpath b/.classpath
index 8856a7e..7a21a4e 100644
--- a/.classpath
+++ b/.classpath
@@ -4,5 +4,6 @@
+
diff --git a/KLALB/KNECloud/iplist.json b/KLALB/KNECloud/iplist.json
deleted file mode 100644
index 836089b..0000000
--- a/KLALB/KNECloud/iplist.json
+++ /dev/null
@@ -1,64 +0,0 @@
-{
- "name" : "KNECloud",
- "services" : [ {
- "name" : "KNE官网",
- "protocol" : "HTTP",
- "localaddress" : "127.0.0.1:8081"
- }, {
- "name" : "KNE云盘",
- "protocol" : "HTTP",
- "localaddress" : "127.0.0.1:5212"
- }, {
- "name" : "服务器远程桌面",
- "protocol" : "RDP",
- "localaddress" : "127.0.0.1:3389"
- }, {
- "name" : "MC",
- "protocol" : "Minecraft",
- "localaddress" : "127.0.0.1:25565"
- } ],
- "tunnels" : [ {
- "name" : "Openfrp-直连线路",
- "remoteaddress" : "127.0.0.1:4566",
- "frpc" : [ "[common]", "server_addr = 180.76.147.250", "server_port = 8120", "tcp_mux = true", "protocol = tcp", "dns_server = 223.5.5.5", "user = bb28264e0abf3bcbec3a2180c5a279b8", "token = I2KMo1HvRxuvurv2", "[ts1xx]", "type = xtcp", "role = visitor", "server_name = zl1", "bind_addr = 127.0.0.1", "bind_port = 4566", "sk = knecloud" ]
- }, {
- "name" : "Openfrp-直连线路",
- "remoteaddress" : "127.0.0.1:4566"
- }, {
- "name" : "Openfrp-直连线路",
- "remoteaddress" : "127.0.0.1:4566"
- }, {
- "name" : "Openfrp-直连线路",
- "remoteaddress" : "127.0.0.1:4566"
- },{
- "name" : "NULL-宿迁联通",
- "remoteaddress" : "153.36.240.12:65529"
- }, {
- "name" : "Openfrp-北京多线-3",
- "remoteaddress" : "cn-bj-bgp-3.openfrp.top:65529"
- }, {
- "name" : "Openfrp-北京多线-7",
- "remoteaddress" : "180.76.147.250:65529"
- }, {
- "name" : "Sakurafrp-安徽电信-1",
- "remoteaddress" : "cn-ah-dx-1.natfrp.cloud:65529"
- }, {
- "name" : "Sakurafrp-南宁电信-1",
- "remoteaddress" : "cn-nn-dx-1.natfrp.cloud:65529"
- }, {
- "name" : "Sakurafrp-武汉电信-1",
- "remoteaddress" : "cn-wh-dx-1.natfrp.cloud:65529"
- }, {
- "name" : "Sakurafrp-枣庄BGP-10",
- "remoteaddress" : "cn-zz-bgp-10.natfrp.cloud:23330"
- }, {
- "name" : "Sakurafrp-枣庄BGP-7",
- "remoteaddress" : "cn-zz-bgp-7.natfrp.cloud:33336"
- }, {
- "name" : "HYHBX-深圳阿里-1",
- "remoteaddress" : "119.23.238.91:12413"
- }, {
- "name" : "HYHBX-上海腾讯-2",
- "remoteaddress" : "43.143.109.64:49965"
- } ]
-}
\ No newline at end of file
diff --git a/KLALB/KNECloud/services.json b/KLALB/KNECloud/services.json
deleted file mode 100644
index fd6b9a9..0000000
--- a/KLALB/KNECloud/services.json
+++ /dev/null
@@ -1,28 +0,0 @@
-{
- "services": [
- {
- "name": "KNE官网",
- "protocol": "HTTP",
- "localaddress": "127.0.0.1:8081",
- "bindaddress": "0.0.0.0:8081"
- },
- {
- "name": "KNE云盘",
- "protocol": "HTTP",
- "localaddress": "127.0.0.1:5212",
- "bindaddress": "0.0.0.0:4568"
- },
- {
- "name": "服务器远程桌面",
- "protocol": "RDP",
- "localaddress": "127.0.0.1:3389",
- "bindaddress": "0.0.0.0:3389"
- },
- {
- "name": "Minecraft",
- "protocol": "Minecraft",
- "localaddress": "127.0.0.1:25565",
- "bindaddress": "0.0.0.0:25586"
- }
- ]
-}
\ No newline at end of file
diff --git a/KLALBClient.bat b/KLALBClient.bat
deleted file mode 100644
index 2193e21..0000000
--- a/KLALBClient.bat
+++ /dev/null
@@ -1 +0,0 @@
-@java -cp ".\bin" org.kne.cloud.network.klalb.KLALBClient %*
\ No newline at end of file
diff --git a/KLALB协议规范V2.0.docx b/KLALB协议规范V2.0.docx
index 84a04c3..0375c53 100644
Binary files a/KLALB协议规范V2.0.docx and b/KLALB协议规范V2.0.docx differ
diff --git a/client.cfg b/client.cfg
new file mode 100644
index 0000000..f61a830
--- /dev/null
+++ b/client.cfg
@@ -0,0 +1,20 @@
+//加速线路
+{KLALBRemote}->43.248.189.107:35000
+{KLALBRemote}->cn-bj-bgp-3.openfrp.top:65529
+{KLALBRemote}->180.76.147.250:65529
+{KLALBRemote}->cn-ah-dx-1.natfrp.cloud:65529
+{KLALBRemote}->cn-nn-dx-1.natfrp.cloud:65529
+{KLALBRemote}->cn-wh-dx-1.natfrp.cloud:65529
+{KLALBRemote}->cn-zz-bgp-10.natfrp.cloud:23330
+{KLALBRemote}->cn-zz-bgp-7.natfrp.cloud:33336
+{KLALBRemote}->43.143.109.64:49965
+{KLALBRemote}->frp.104300.xyz:49965
+{KLALBRemote}->us.afrps.cn:49966
+{KLALBRemote}->hk.afrps.cn:49966
+{KLALBRemote}->la.afrps.cn:49966
+{KLALBRemote}->frp.freefrp.net:49965
+{KLALBRemote}->frp1.freefrp.net:49965
+{KLALBRemote}->frp2.freefrp.net:49965
+{KLALBRemote}->frp4.freefrp.net:49965
+
+0.0.0.0:25565->{KLALBVirtual}[171d:a999:e697:4b23:ae52:9f29:4532:e1ee]:25565
\ No newline at end of file
diff --git a/klalbs4.json b/klalbs4.json
deleted file mode 100644
index 2c9908b..0000000
--- a/klalbs4.json
+++ /dev/null
@@ -1,29 +0,0 @@
-{
- "name" : "KNECloud",
- "services" : [ {
- "name" : "KNE官网",
- "protocol" : "HTTP",
- "localaddress" : "127.0.0.1:8081"
- }, {
- "name" : "KNE云盘",
- "protocol" : "HTTP",
- "localaddress" : "127.0.0.1:5212"
- }, {
- "name" : "服务器远程桌面",
- "protocol" : "RDP",
- "localaddress" : "127.0.0.1:3389"
- } ],
- "tunnels" : [ {
- "name" : "Openfrp-韩国春川-1",
- "remoteaddress" : "kr-nc-bgp-1.openfrp.top:65529"
- }, {
- "name" : "Openfrp-北京多线-3",
- "remoteaddress" : "cn-bj-bgp-3.openfrp.top:65529"
- }, {
- "name" : "Openfrp-北京多线-7",
- "remoteaddress" : "180.76.147.250:65529"
- }, {
- "name" : "Openfrp-美国-8「BGP」",
- "remoteaddress" : "us-los-bgp-8.openfrp.top:65529"
- }]
-}
diff --git a/kserver.ini b/kserver.ini
index 2a66e37..7e44418 100644
--- a/kserver.ini
+++ b/kserver.ini
@@ -1 +1,3 @@
-port=4569
+virtualip=67c:72ce:765e:4db1:a02e:92fa:1959:29eb
+bind=0.0.0.0:4569
+local=127.0.0.1:5212
diff --git a/lines.cfg b/lines.cfg
new file mode 100644
index 0000000..64f20cf
--- /dev/null
+++ b/lines.cfg
@@ -0,0 +1,29 @@
+0.0.0.0:35000(RDP)->192.168.1.233:3389
+0.0.0.0:35000(HTTPS)->192.168.1.233:8444
+0.0.0.0:35000(HTTP)->192.168.1.233:5212
+0.0.0.0:35000(SSH)->192.168.1.236:22
+0.0.0.0:35000(KLALB)->{KLALBRemote}
+0.0.0.0:35000->127.0.0.1:35001 //Winds服(普通线路)
+
+{KLALBVirtual}[::0]:25565->127.0.0.1:35001 //Winds服(加速线路)
+
+//加速线路
+{KLALBRemote}->43.248.189.107:35000
+{KLALBRemote}->cn-bj-bgp-3.openfrp.top:65529
+{KLALBRemote}->180.76.147.250:65529
+{KLALBRemote}->cn-ah-dx-1.natfrp.cloud:65529
+{KLALBRemote}->cn-nn-dx-1.natfrp.cloud:65529
+{KLALBRemote}->cn-wh-dx-1.natfrp.cloud:65529
+{KLALBRemote}->cn-zz-bgp-10.natfrp.cloud:23330
+{KLALBRemote}->cn-zz-bgp-7.natfrp.cloud:33336
+{KLALBRemote}->43.143.109.64:49965
+{KLALBRemote}->frp.104300.xyz:49965
+{KLALBRemote}->us.afrps.cn:49966
+{KLALBRemote}->hk.afrps.cn:49966
+{KLALBRemote}->la.afrps.cn:49966
+{KLALBRemote}->frp.freefrp.net:49965
+{KLALBRemote}->frp1.freefrp.net:49965
+{KLALBRemote}->frp2.freefrp.net:49965
+{KLALBRemote}->frp4.freefrp.net:49965
+
+0.0.0.0:25565->{KLALBVirtual}[171d:a999:e697:4b23:ae52:9f29:4532:e1ee]:25565
\ No newline at end of file
diff --git a/ports.ini b/ports.ini
deleted file mode 100644
index 8fec8bc..0000000
--- a/ports.ini
+++ /dev/null
@@ -1,7 +0,0 @@
-7462
-RDP$192.168.1.233:3389
-HTTPS$192.168.1.233:8444
-HTTP$192.168.1.235:80
-SSH$192.168.1.236:22
-KLALB$192.168.1.233:4569
-DEFAULT$192.168.1.236:22
\ No newline at end of file
diff --git a/server.cfg b/server.cfg
new file mode 100644
index 0000000..bbaf70b
--- /dev/null
+++ b/server.cfg
@@ -0,0 +1,8 @@
+0.0.0.0:35000(RDP)->192.168.1.233:3389
+0.0.0.0:35000(HTTPS)->192.168.1.233:8444
+0.0.0.0:35000(HTTP)->192.168.1.233:5212
+0.0.0.0:35000(SSH)->192.168.1.236:22
+0.0.0.0:35000(KLALB)->{KLALBRemote}
+0.0.0.0:35000->127.0.0.1:35001 //Winds服(普通线路)
+
+{KLALBVirtual}[::0]:25565->127.0.0.1:35001 //Winds服(加速线路)
\ No newline at end of file
diff --git a/src/org/kne/cloud/network/klalb/ACKTPacket.java b/src/org/kne/cloud/network/klalb/ACKTPacket.java
index 7f6303d..b920bc8 100644
--- a/src/org/kne/cloud/network/klalb/ACKTPacket.java
+++ b/src/org/kne/cloud/network/klalb/ACKTPacket.java
@@ -17,7 +17,7 @@ public class ACKTPacket extends KLALBPacket {
}
public ACKTPacket() {
- super();
+ super(ACKT);
}
@Override
diff --git a/src/org/kne/cloud/network/klalb/DATATPacket.java b/src/org/kne/cloud/network/klalb/DATATPacket.java
index 9097af1..bddeeaa 100644
--- a/src/org/kne/cloud/network/klalb/DATATPacket.java
+++ b/src/org/kne/cloud/network/klalb/DATATPacket.java
@@ -22,7 +22,7 @@ public class DATATPacket extends KLALBPacket {
}
public DATATPacket() {
- super();
+ super(DATAT);
}
public int getSport() {
diff --git a/src/org/kne/cloud/network/klalb/KLALBController.java b/src/org/kne/cloud/network/klalb/KLALBController.java
index 1718947..46661d3 100644
--- a/src/org/kne/cloud/network/klalb/KLALBController.java
+++ b/src/org/kne/cloud/network/klalb/KLALBController.java
@@ -11,18 +11,24 @@ import java.net.Socket;
import java.net.SocketAddress;
import java.net.SocketTimeoutException;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
+import java.util.Scanner;
import java.util.Set;
import java.util.Timer;
import java.util.UUID;
import java.util.Vector;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
+import java.util.function.Supplier;
+import org.kne.cloud.network.mport.MultipurposeSocketAddress;
+import org.kne.cloud.network.mport.NetworkService;
+import org.kne.cloud.network.mport.ProxyProfileEntry;
import org.kne.cloud.network.mport.ThreadTool;
import javassist.ClassPool;
@@ -33,6 +39,29 @@ import javassist.bytecode.CodeAttribute;
import javassist.bytecode.CodeIterator;
public class KLALBController {
+ private LineManager lineManager;
+ private Supplier selflineTableSupplier=()->{return null;};
+
+
+
+
+ public Supplier getSelflineTableSupplier() {
+ return selflineTableSupplier;
+ }
+
+ public void setSelflineTableSupplier(Supplier selflineTableSupplier) {
+ this.selflineTableSupplier = selflineTableSupplier;
+ }
+
+ public LineManager getLineManager() {
+ if(lineManager==null)
+ lineManager=createLineManager();
+ return lineManager;
+ }
+
+ protected LineManager createLineManager() {
+ return new LineManager(this);
+ }
private Inet6Address self;
@@ -70,6 +99,7 @@ public class KLALBController {
}
public void addRemoteSocket(KLALBRemoteSocket krs) {
+ krs.setController(this);
CountDownLatch cdl = new CountDownLatch(1);
krs.setPacketReceiver((rec) -> {
try {
@@ -139,6 +169,14 @@ public class KLALBController {
VADDRPacket var = (VADDRPacket) rec;
setLine(var.getVaddr(), krs);
cdl.countDown();
+ }else if(rec instanceof LINESPacket) {
+ LINESPacket lpt=(LINESPacket) rec;
+ String s=lpt.getLines();
+ Scanner scn=new Scanner(s);
+ while(scn.hasNext()) {
+ String sn=scn.nextLine();
+ getLineManager().addHostPort(new MultipurposeSocketAddress(sn));
+ }
}
} catch (NoRouteToHostException e) {
e.printStackTrace();
@@ -146,8 +184,12 @@ public class KLALBController {
});
krs.setCloseListener((x) -> {
removeLine(krs);
+ cdl.countDown();
});
krs.sendPacket(new VADDRPacket(self), 65537);
+ String selflineTable=selflineTableSupplier.get();
+ if(selflineTable!=null)
+ krs.sendPacket(new LINESPacket(selflineTable), 65537);
try {
cdl.await();
} catch (InterruptedException e) {
@@ -155,14 +197,14 @@ public class KLALBController {
}
}
- public KLALBController() {
- this.self = KLALBUtils.uuidToIP(UUID.randomUUID());
- }
-
public KLALBController(Inet6Address self) {
this.self = self;
}
+ public KLALBController() {
+ this.self=KLALBUtils.uuidToIP(UUID.randomUUID());
+ }
+
private Map bindmap = new ConcurrentHashMap<>();
protected Map getBindmap() {
@@ -225,7 +267,8 @@ public class KLALBController {
int count0 = Math.min(count, l.size());
List l2 = (List) ((ArrayList) l).clone();
LineDecitionComparator ldc = new LineDecitionComparator(l2, priority);
- l2.sort(ldc);
+ Collections.sort(l2,ldc);
+ //System.out.println(l2);
for (Iterator iterator = l2.iterator(); iterator.hasNext();) {
KLALBRemoteSocket krst = (KLALBRemoteSocket) iterator.next();
krst.sendPacket(packet, priority);
@@ -236,7 +279,16 @@ public class KLALBController {
}
}
-
+ protected void removeFromSend(Inet6Address addr,KLALBPacket klalbPacket) {
+ synchronized (routes) {
+ List l = routes.get(addr);
+ if (l != null && !l.isEmpty()) {
+ l.forEach((x)->{
+ x.remoeFromSendQueue(klalbPacket);
+ });
+ }
+ }
+ }
protected boolean checkIsBind(KLALBVirtualSocketImpl klalbVirtualSocketImpl) {
return bindmap.containsValue(klalbVirtualSocketImpl);
}
@@ -245,4 +297,37 @@ private Timer t=new Timer("数据包发送计时器", true);
return t;
}
+ public void registerToProxyTypeAs(String proxyname) {
+ MultipurposeSocketAddress.getSocketFactoryRegister().put(proxyname, new KLALBVirtualSocketFactory(this));
+ MultipurposeSocketAddress.getServerSocketFactoryRegister().put(proxyname, new KLALBVirtualServerSocketFactory(this));
+ ProxyProfileEntry.getRegister().put(proxyname, new KLALBNetworkService());
+ }
+
+ private class KLALBNetworkService implements NetworkService{
+
+ @Override
+ public void listen(MultipurposeSocketAddress msa) {
+ // TODO 自动生成的方法存根
+
+ }
+
+ @Override
+ public void unlisten(MultipurposeSocketAddress msc) {
+ // TODO 自动生成的方法存根
+
+ }
+
+ @Override
+ public void connect(MultipurposeSocketAddress msa) {
+ getLineManager().addHostPort(msa);
+ }
+
+ @Override
+ public void unconnect(MultipurposeSocketAddress msc) {
+ getLineManager().removeHostPort(msc);
+ }
+
+ }
+
+
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBInputStream.java b/src/org/kne/cloud/network/klalb/KLALBInputStream.java
index 6e024b8..af6489b 100644
--- a/src/org/kne/cloud/network/klalb/KLALBInputStream.java
+++ b/src/org/kne/cloud/network/klalb/KLALBInputStream.java
@@ -60,6 +60,10 @@ public class KLALBInputStream extends DataInputStream {
klp=new VADDRPacket();
klp.readFromStream(this);
return klp;
+ case LINES:
+ klp=new LINESPacket();
+ klp.readFromStream(this);
+ return klp;
}
throw new StreamCorruptedException("unknown package type:"+type);
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBMain.java b/src/org/kne/cloud/network/klalb/KLALBMain.java
new file mode 100644
index 0000000..b9440ab
--- /dev/null
+++ b/src/org/kne/cloud/network/klalb/KLALBMain.java
@@ -0,0 +1,72 @@
+package org.kne.cloud.network.klalb;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Scanner;
+
+import org.kne.cloud.network.klalb.LineManager.LineEntry;
+import org.kne.cloud.network.mport.MultipurposeSocketAddress;
+import org.kne.cloud.network.mport.ProxyProfileAnalyser;
+import org.kne.cloud.network.mport.ProxyProfileExecutor;
+import org.kne.cloud.network.mport.TCPListener;
+
+public class KLALBMain {
+ public static TCPListener tcpl;
+ public static ServerPropties sp;
+ public static KLALBController kc;
+
+ public static ProxyProfileExecutor pfa;
+ public static File f=new File("lines.cfg");
+ public static void main(String[] args) throws IOException {
+
+ System.out.println("KLALB负载均衡V2.0");
+ sp=new ServerPropties();
+
+ System.out.println("虚拟地址:"+sp.getVirtualIP().getHostAddress());
+ kc=new KLALBController(sp.getVirtualIP());
+ kc.registerToProxyTypeAs("KLALB");
+
+
+
+
+ pfa=new ProxyProfileExecutor();
+ pfa.load(f);
+
+ Scanner scn=new Scanner(System.in);
+ while(true) {
+ String s=scn.next();
+ String[]sc=s.split(" ");
+ switch(sc[0]) {
+ case "help":
+ System.out.print("state:查看线路状态");
+ System.out.println("reload:重新加载线路配置");
+ break;
+ case "reloadlines":
+ pfa.load(f);
+ System.out.println("重新加载线路配置成功");
+ break;
+ case "state":
+ LineManager le=kc.getLineManager();
+ for (Iterator> iterator = le.entrySet().iterator(); iterator.hasNext();) {
+ java.util.Map.Entry hostPort = iterator.next();
+ System.out.println(hostPort.getValue() .toString());
+ }
+ break;
+ default:
+ System.out.println("未知命令,请输入help以查询指令说明");
+ }
+ }
+ }
+ private static void openPort(int port) throws IOException {
+ if(tcpl!=null)
+ tcpl.close();
+ tcpl=new TCPListener(new MultipurposeSocketAddress("0.0.0.0", port));
+ tcpl.setCon((soc)->{
+ KLALBRemoteSocket krs=new KLALBRemoteSocket(soc);
+ kc.addRemoteSocket(krs);
+ });
+ tcpl.open();
+ }
+}
diff --git a/src/org/kne/cloud/network/klalb/KLALBPacket.java b/src/org/kne/cloud/network/klalb/KLALBPacket.java
index 0d99eb6..26d5fd7 100644
--- a/src/org/kne/cloud/network/klalb/KLALBPacket.java
+++ b/src/org/kne/cloud/network/klalb/KLALBPacket.java
@@ -17,12 +17,9 @@ public abstract class KLALBPacket{
public static final int ACKT=6;
public static final int DATAU=7;
public static final int VADDR=8;
-
+ public static final int LINES=9;
private int type;
- public KLALBPacket() {
- super();
- }
public KLALBPacket(int type) {
super();
this.type = type;
@@ -41,6 +38,7 @@ public abstract class KLALBPacket{
}
public long getLength() {
- return 4;
+ return 1;
}
+
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java b/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java
index 811ef0f..1c0b26d 100644
--- a/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java
+++ b/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java
@@ -2,7 +2,10 @@ package org.kne.cloud.network.klalb;
import java.io.IOException;
import java.net.Inet6Address;
+import java.net.NoRouteToHostException;
import java.net.Socket;
+import java.net.SocketException;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
@@ -12,12 +15,35 @@ import org.kne.cloud.network.mport.ThreadTool;
public class KLALBRemoteSocket extends MonitoredSocket {
private Inet6Address remoteVaddr;
+ private KLALBController controller;
+
+ public KLALBController getController() {
+ return controller;
+ }
+ protected void setController(KLALBController controller) {
+ this.controller = controller;
+ }
public Inet6Address getRemoteVaddr() {
return remoteVaddr;
}
-
- public KLALBRemoteSocket(Socket socket) {
- super(socket);
+ public String toString() {
+ return super.toString()+ getMonitor().toString();
+ }
+public KLALBRemoteSocket(Socket socket) {
+ this(socket,new Monitor());
+}
+ public KLALBRemoteSocket(Socket socket,Monitor m) {
+ super(socket,m);
+ try {
+ socket.setSendBufferSize(32768);
+ socket.setReceiveBufferSize(65535);
+ //System.out.println(socket.getSendBufferSize()+" "+socket.getReceiveBufferSize());
+ socket.setTcpNoDelay(true);
+ socket.setSoTimeout(10000);
+ } catch (SocketException e1) {
+ // TODO 自动生成的 catch 块
+ e1.printStackTrace();
+ }
ThreadTool.makeVDaemonThreadIfSupport("远程接收线程", () -> {
try {
@@ -25,11 +51,18 @@ public class KLALBRemoteSocket extends MonitoredSocket {
KLALBPacket kpp = getInputStream().readPacket();
//System.out.println("RX:" + kpp);
if(kpp instanceof PINGPacket) {
- sendPacket(new PONGPacket(), 65540);
+ sendPacket(new PONGPacket(((PINGPacket) kpp).getTime()), 65540);
}else if(kpp instanceof PONGPacket) {
- latency.set(System.nanoTime()-pingstarttime);
- pinging=false;
+ PONGPacket png=(PONGPacket) kpp;
+ getMonitor().updateLatency (System.nanoTime()- png.getTime());
+ try {
+ socket.setSoTimeout(1100);
+ } catch (SocketException e1) {
+ // TODO 自动生成的 catch 块
+ e1.printStackTrace();
+ }
}else {
+
if(kpp instanceof VADDRPacket) {
remoteVaddr=((VADDRPacket) kpp).getVaddr();
}
@@ -40,7 +73,7 @@ public class KLALBRemoteSocket extends MonitoredSocket {
}
}
} catch (IOException e) {
- e.printStackTrace();
+ //e.printStackTrace();
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
@@ -55,18 +88,18 @@ public class KLALBRemoteSocket extends MonitoredSocket {
try {
while (!isClosed()) {
- if(System.nanoTime()-pingstarttime>100000000000L) {
- pinging=false;
- }
- if(checkPingTime()&&!pinging) {
- getOutputStream().writePacket(new PINGPacket());
- pinging =true;
- pingstarttime=System.nanoTime();
+ if(checkPingTime()) {
+ getOutputStream().writePacket(new PINGPacket(System.nanoTime()));
}else {
PQItem pqi=sendQueue.poll();
if(pqi!=null) {
KLALBPacket kpp=pqi.getPacket();
+ //if(!(kpp instanceof PONGPacket))
//System.out.println("TX:" + kpp);
+ if(kpp instanceof PONGPacket) {
+ PONGPacket pp=(PONGPacket) kpp;
+ pp.redeltaTime(System.nanoTime()-pqi.getAddtime());
+ }
getOutputStream().writePacket(kpp);
}else {
Thread.sleep(1);
@@ -74,7 +107,7 @@ public class KLALBRemoteSocket extends MonitoredSocket {
}
}
} catch (IOException e) {
- e.printStackTrace();
+ // e.printStackTrace();
} catch (InterruptedException e) {
} finally {
try {
@@ -87,11 +120,9 @@ public class KLALBRemoteSocket extends MonitoredSocket {
}
- private volatile long pingstarttime=System.nanoTime();
- private volatile boolean pinging=false;
private volatile long time = System.nanoTime();
- private volatile long pingInterval=1000000000L;
+ private volatile long pingInterval=500000000L;
public long getPingInterval() {
return pingInterval;
@@ -112,11 +143,8 @@ public class KLALBRemoteSocket extends MonitoredSocket {
}
- private volatile boolean waitingforPing = false;
- public long getLatency() {
- return latency.get();
- }
+
private KLALBInputStream kis;
private KLALBOutputStream kos;
@@ -135,7 +163,7 @@ public class KLALBRemoteSocket extends MonitoredSocket {
return kos;
}
- private AtomicLong latency = new AtomicLong();
+
private volatile Consumer rec;
private Consumer clo;
private PriorityBlockingQueue sendQueue=new PriorityBlockingQueue();
@@ -160,9 +188,22 @@ public class KLALBRemoteSocket extends MonitoredSocket {
}
@Override
public synchronized void close() throws IOException {
- super.close();
if(!isClosed()&&clo!=null)
clo.accept(this);
+ cdl.countDown();
+ super.close();
+ sendQueue.forEach((x)->{
+ try {
+ controller.sendPacketToAddress(remoteVaddr, x.getPacket(), x.getPriority()+1);
+ } catch (NoRouteToHostException e) {
+ }
+ });
+
+ }
+
+ @Override
+ public boolean isClosed() {
+ return cdl.getCount()<=0;
}
private AtomicLong al=new AtomicLong(0);
@@ -170,6 +211,10 @@ public class KLALBRemoteSocket extends MonitoredSocket {
private KLALBPacket packet;
private long index=al.getAndIncrement();
private int priority=5;
+ private long addtime=System.nanoTime();
+ public long getAddtime() {
+ return addtime;
+ }
public PQItem(KLALBPacket value, int priority) {
super();
this.packet = value;
@@ -201,4 +246,17 @@ public class KLALBRemoteSocket extends MonitoredSocket {
}
}
}
+ private CountDownLatch cdl=new CountDownLatch(1);
+ public void waitforlose() {
+ try {
+ cdl.await();
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ }
+ public void remoeFromSendQueue(KLALBPacket klalbPacket) {
+ sendQueue.removeIf((x)->{
+ return x.getPacket().equals(klalbPacket);
+ });
+ }
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBTest1.java b/src/org/kne/cloud/network/klalb/KLALBTest1.java
deleted file mode 100644
index 07f1f60..0000000
--- a/src/org/kne/cloud/network/klalb/KLALBTest1.java
+++ /dev/null
@@ -1,52 +0,0 @@
-package org.kne.cloud.network.klalb;
-
-import java.io.IOException;
-import java.net.Inet6Address;
-import java.net.InetAddress;
-import java.net.InetSocketAddress;
-import java.net.Socket;
-import java.net.SocketException;
-import java.net.UnknownHostException;
-
-import org.kne.cloud.network.mport.SocketBridge;
-import org.kne.cloud.network.mport.TCPListener;
-
-public class KLALBTest1 {
-
- public static void main(String[] args) throws UnknownHostException, IOException, InterruptedException {
- KLALBRemoteSocket s1=new KLALBRemoteSocket( new Socket("127.0.0.1", 8899));
- KLALBRemoteSocket s2=new KLALBRemoteSocket(new Socket("127.0.0.1", 8899));
- KLALBRemoteSocket s3=new KLALBRemoteSocket(new Socket("127.0.0.1", 8899));
- KLALBController kc=new KLALBController((Inet6Address) Inet6Address.getByName("::2"));
- kc.addRemoteSocket(s1);
- kc.addRemoteSocket(s2);
- kc.addRemoteSocket(s3);
- TCPListener tl=new TCPListener(7463);
- tl.open();
- tl.setCon((x)->{
- KLALBVirtualSocket kvs;
- try {
- kvs = new KLALBVirtualSocket(kc);
- kvs.connect(new InetSocketAddress(InetAddress.getByName("::9"), 90));
- new SocketBridge(kvs, x).run();
- } catch (SocketException e) {
- // TODO 自动生成的 catch 块
- e.printStackTrace();
- } catch (UnknownHostException e) {
- // TODO 自动生成的 catch 块
- e.printStackTrace();
- } catch (IOException e) {
- // TODO 自动生成的 catch 块
- e.printStackTrace();
- }
- });
-
- System.out.println("c");
- Thread.sleep(20000);
- s2.close();
- while(true) {
- Thread.sleep(1000);
- }
- }
-
-}
diff --git a/src/org/kne/cloud/network/klalb/KLALBTest2.java b/src/org/kne/cloud/network/klalb/KLALBTest2.java
deleted file mode 100644
index 6307cfa..0000000
--- a/src/org/kne/cloud/network/klalb/KLALBTest2.java
+++ /dev/null
@@ -1,36 +0,0 @@
-package org.kne.cloud.network.klalb;
-
-import java.io.IOException;
-import java.net.*;
-
-import org.kne.cloud.network.mport.SocketBridge;
-import org.kne.cloud.network.mport.TCPListener;
-
-public class KLALBTest2 {
- public static void main(String[] args) throws IOException, InterruptedException {
- ServerSocket ssk=new ServerSocket(8899);
- KLALBRemoteSocket s1=new KLALBRemoteSocket(ssk.accept());
- KLALBRemoteSocket s2=new KLALBRemoteSocket(ssk.accept());
- KLALBRemoteSocket s3=new KLALBRemoteSocket(ssk.accept());
- KLALBController cont=new KLALBController((Inet6Address) Inet6Address.getByName("::9"));
- cont.addRemoteSocket(s1);
- cont.addRemoteSocket(s2);
- cont.addRemoteSocket(s3);
- KLALBVirtualServerSocket ksr=new KLALBVirtualServerSocket(cont);
- TCPListener tl=new TCPListener(90);
- tl.setServerSocketFactory(new KLALBVirtualServerSocketFactory(cont));
- tl.setCon((x)->{
- try {
- new SocketBridge(x, new Socket("127.0.0.1",5212)).run();
- } catch (UnknownHostException e) {
- e.printStackTrace();
- } catch (IOException e) {
- e.printStackTrace();
- }
- });
- tl.open();
-
- Thread.sleep(20000);
- s3.close();
- }
-}
diff --git a/src/org/kne/cloud/network/klalb/KLALBUtils.java b/src/org/kne/cloud/network/klalb/KLALBUtils.java
index f873fd2..3e4a5bb 100644
--- a/src/org/kne/cloud/network/klalb/KLALBUtils.java
+++ b/src/org/kne/cloud/network/klalb/KLALBUtils.java
@@ -29,4 +29,29 @@ public class KLALBUtils {
public static void main(String[] args) {
System.out.println(uuidToIP(new UUID(-1,-1)));
}
+ public static String bytesUnit(long v) {
+ if(v>=1024L*1024*1024*1024*1024) {
+ return format( v/(1024.0*1024.0*1024.0*1024.0*1024.0))+"P";
+ }else if(v>=1024L*1024*1024*1024) {
+ return format (v/(1024.0*1024.0*1024.0*1024.0))+"T";
+ }else if(v>=1024L*1024*1024) {
+ return format( v/(1024.0*1024.0*1024.0))+"G";
+ }else if(v>=1024L*1024) {
+ return format( v/(1024.0*1024.0))+"M";
+ }else if(v>=1024L) {
+ return format( v/(1024.0))+"K";
+ }else {
+ return v+"B";
+ }
+ }
+ private static String format( double d) {
+ String fmt=String.format( "%.2f",d);
+ if(fmt.length()>4) {
+ fmt=fmt.substring(0, 4);
+ }
+ if(fmt.endsWith(".")) {
+ fmt=fmt.substring(0, fmt.length()-1);
+ }
+ return fmt;
+ }
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBVirtualSocket.java b/src/org/kne/cloud/network/klalb/KLALBVirtualSocket.java
index 79414de..57d702b 100644
--- a/src/org/kne/cloud/network/klalb/KLALBVirtualSocket.java
+++ b/src/org/kne/cloud/network/klalb/KLALBVirtualSocket.java
@@ -1,7 +1,7 @@
package org.kne.cloud.network.klalb;
-import java.net.SocketException;
-import java.net.SocketImpl;
+import java.io.IOException;
+import java.net.*;
import org.kne.cloud.network.mport.VirtualSocket;
@@ -11,11 +11,52 @@ public class KLALBVirtualSocket extends VirtualSocket {
public KLALBVirtualSocket(KLALBController controler) throws SocketException {
super(controler.createVirtualImpl());
- this.controler=controler;
+ this.controler = controler;
}
protected KLALBController getControler() {
return controler;
}
+ public KLALBVirtualSocket(KLALBController controler,String host, int port) throws UnknownHostException, IOException {
+ this(controler,host != null ? new InetSocketAddress(host, port)
+ : new InetSocketAddress(InetAddress.getByName(null), port), (SocketAddress) null);
+ }
+
+ public KLALBVirtualSocket(KLALBController controler,SocketAddress address, SocketAddress localAddr) throws IOException {
+ this(controler);
+ // backward compatibility
+ if (address == null)
+ throw new NullPointerException();
+
+ try {
+ if (localAddr != null)
+ bind(localAddr);
+ connect(address);
+ } catch (IOException | IllegalArgumentException | SecurityException e) {
+ try {
+ close();
+ } catch (IOException ce) {
+ e.addSuppressed(ce);
+ }
+ throw e;
+ }
+ }
+
+ public KLALBVirtualSocket(KLALBController controler,String host, int port, InetAddress localAddr, int localPort) throws IOException {
+ this(controler,host != null ? new InetSocketAddress(host, port)
+ : new InetSocketAddress(InetAddress.getByName(null), port),
+ new InetSocketAddress(localAddr, localPort));
+ }
+
+ public KLALBVirtualSocket(KLALBController controler,InetAddress address, int port, InetAddress localAddr, int localPort) throws IOException {
+ this(controler,address != null ? new InetSocketAddress(address, port) : null,
+ new InetSocketAddress(localAddr, localPort));
+ }
+
+ public KLALBVirtualSocket(KLALBController controler,InetAddress address, int port) throws IOException {
+ this(controler,address != null ? new InetSocketAddress(address, port) : null,
+ (SocketAddress) null);
+ }
+
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBVirtualSocketFactory.java b/src/org/kne/cloud/network/klalb/KLALBVirtualSocketFactory.java
new file mode 100644
index 0000000..51f8064
--- /dev/null
+++ b/src/org/kne/cloud/network/klalb/KLALBVirtualSocketFactory.java
@@ -0,0 +1,39 @@
+package org.kne.cloud.network.klalb;
+
+import java.io.IOException;
+import java.net.InetAddress;
+import java.net.Socket;
+import java.net.UnknownHostException;
+
+import javax.net.SocketFactory;
+
+public class KLALBVirtualSocketFactory extends SocketFactory {
+private KLALBController controller;
+
+ public KLALBVirtualSocketFactory(KLALBController controller) {
+ super();
+ this.controller = controller;
+ }
+ @Override
+ public Socket createSocket(String arg0, int arg1) throws IOException, UnknownHostException {
+
+ return new KLALBVirtualSocket(controller, arg0, arg1);
+ }
+
+ @Override
+ public Socket createSocket(InetAddress arg0, int arg1) throws IOException {
+ return new KLALBVirtualSocket(controller, arg0, arg1);
+ }
+
+ @Override
+ public Socket createSocket(String arg0, int arg1, InetAddress arg2, int arg3)
+ throws IOException, UnknownHostException {
+ return new KLALBVirtualSocket(controller, arg0, arg1, arg2, arg3);
+ }
+
+ @Override
+ public Socket createSocket(InetAddress arg0, int arg1, InetAddress arg2, int arg3) throws IOException {
+ return new KLALBVirtualSocket(controller, arg0, arg1, arg2, arg3);
+ }
+
+}
diff --git a/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java b/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java
index 52d16fc..093b1c5 100644
--- a/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java
+++ b/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java
@@ -17,14 +17,18 @@ import java.net.SocketException;
import java.net.SocketImpl;
import java.net.SocketOptions;
import java.net.SocketTimeoutException;
+import java.net.UnknownHostException;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingDeque;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.LinkedBlockingDeque;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.atomic.AtomicReferenceArray;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
@@ -51,13 +55,36 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
this.outputchachesize = outputchachesize;
}
- private int inputchachesize = 65535 * 50;
- private int outputchachesize = 65535 * 50;
- private Inet6Address bindaddr;
+ private int inputchachesize = 5773 * 100;
+ private int outputchachesize = 5773 * 100;
+ private Inet6Address bindaddr;{
+ try {
+ bindaddr=(Inet6Address) Inet6Address.getByName("::0");
+ } catch (UnknownHostException e) {
+ e.printStackTrace();
+ }
+ }
private BlockingQueue backlogQueue;
- private BlockingDeque sendDeque = new LinkedBlockingDeque<>();
+ private Queue sendDeque = new ConcurrentLinkedQueue<>();
private List sendlist=new Vector<>();
+ private TimerTask sendCheckTask=new SendCheckTask();
+ private class SendCheckTask extends TimerTask{
+
+ public void run() {
+ synchronized (sendlist) {
+ for (int i = 0; i < sendlist.size(); i++) {
+ try {
+ sendlist.get(i).check(i);
+ } catch (IOException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }
+ }
+ }
+ }
+
+ }
private boolean succeed, refused;
protected boolean isListening() {
@@ -126,16 +153,18 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
}
}
controller.sendPacketToAddress(from, new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp.getNumber(),
- countInputBytes() < inputchachesize), 32768);
+ countInputBytes() < inputchachesize), 32768,2);
} else if (u instanceof ACKTPacket) {
ACKTPacket ackt = (ACKTPacket) u;
avaliable = ackt.isAvaliable();
+ AtomicReferencekl=new AtomicReference<>();
sendlist.removeIf((tsk)->{
- boolean b=tsk.getKp().getNumber()==ackt.getNumber();
- if(b)
- tsk.cancel();
+ boolean b=tsk.getPacket().getNumber()==ackt.getNumber();
+ if(b)
+ kl.set(tsk.getPacket());
return b;
});
+ controller.removeFromSend(from,kl.get());
}
} catch (IOException e) {
e.printStackTrace();
@@ -155,7 +184,7 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
private int countOutputBytes() {
AtomicInteger i = new AtomicInteger(0);
sendlist.forEach((c) -> {
- i.addAndGet(c.getKp().getData().length);
+ i.addAndGet(c.getPacket().getData().length);
});
return i.get();
}
@@ -170,17 +199,27 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
public KLALBVirtualSocketImpl(KLALBController kc) {
super();
this.controller = kc;
+ kc.getTimer().schedule(sendCheckTask, 50, 50);
}
@Override
public void setOption(int optID, Object value) throws SocketException {
- // TODO 自动生成的方法存根
-
+ if(optID==SocketOptions.SO_RCVBUF) {
+ inputchachesize=(int) value;
+ }else if(optID==SocketOptions.SO_SNDBUF) {
+ outputchachesize=(int)value;
+ }
}
@Override
public Object getOption(int optID) throws SocketException {
- // TODO 自动生成的方法存根
+ if(optID==SocketOptions.SO_BINDADDR) {
+ return bindaddr;
+ }else if(optID==SocketOptions.SO_RCVBUF) {
+ return inputchachesize;
+ }else if(optID==SocketOptions.SO_SNDBUF) {
+ return outputchachesize;
+ }
return null;
}
@@ -255,6 +294,7 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
@Override
protected void listen(int backlog) throws IOException {
backlogQueue = new ArrayBlockingQueue<>(backlog);
+ address=bindaddr;
}
private Map accepts = new ConcurrentHashMap<>();
@@ -268,7 +308,8 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
try {
KLALBVirtualSocketImpl kvsi = (KLALBVirtualSocketImpl) s;
InetSocketAddress isa = backlogQueue.take();
-
+ kvsi.inputchachesize=inputchachesize;
+ kvsi.outputchachesize=outputchachesize;
kvsi.port = isa.getPort();
kvsi.address = isa.getAddress();
kvsi.localport = localport;
@@ -381,21 +422,21 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
private class KVSIOutputStream extends OutputStream {
- private ByteArrayOutputStream bos = new ByteArrayOutputStream(65535);
-
+ private byte[] cache=new byte[65535];
+ private int count=0;
@Override
public void write(int b) throws IOException {
if (isClosed())
throw new SocketException("Socket is closed");
- bos.write(b);
- if (bos.size() >= 65535) {
+ cache[count++]=(byte) b;
+ if (count >=5773 ) {//1429?
flush();
}
}
@Override
public void flush() throws IOException {
- if (bos.size() > 0) {
+ if (count > 0) {
try {
while (!avaliable) {
Thread.sleep(1);
@@ -411,20 +452,27 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
// TODO 自动生成的 catch 块
e.printStackTrace();
}
- SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, bos.toByteArray()), 5,5);
- x.addToTimer(1000);
+ byte[]ba;
+ if(count==cache.length) {
+ ba=cache;
+ cache=new byte[cache.length];
+ }else {
+ ba=Arrays.copyOf(cache, count);
+ }
+ SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, ba), 5,10);
+ x.run();
sendlist.add(x);
/*controller.sendPacketToAddress((Inet6Address) address,
new DATATPacket(localport, port, outputcount++, bos.toByteArray()), 5);*/
}
- bos.reset();
+ count=0;
}
@Override
public void close() throws IOException {
flush();
- SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, new byte[0]), 5,5);
- x.addToTimer(1000);
+ SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, new byte[0]), 5,10);
+ x.run();
sendlist.add(x);
/*controller.sendPacketToAddress((Inet6Address) address,
new DATATPacket(localport, port, outputcount++, new byte[0]), 5);*/
@@ -480,6 +528,7 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
acceptedSocketCloseListener.accept(this);
else
controller.unbind(this);
+ sendCheckTask.cancel();
}
private boolean closed;
diff --git a/src/org/kne/cloud/network/klalb/LINESPacket.java b/src/org/kne/cloud/network/klalb/LINESPacket.java
new file mode 100644
index 0000000..8e586fe
--- /dev/null
+++ b/src/org/kne/cloud/network/klalb/LINESPacket.java
@@ -0,0 +1,60 @@
+package org.kne.cloud.network.klalb;
+
+import java.io.DataInput;
+import java.io.DataOutput;
+import java.io.IOException;
+import java.net.Inet6Address;
+
+public class LINESPacket extends KLALBPacket {
+ private String lines;
+
+
+ public String getLines() {
+ return lines;
+ }
+
+ public LINESPacket(String lines) {
+ super(LINES);
+ this.lines=lines;
+ }
+
+ public LINESPacket() {
+ super(LINES);
+ }
+ private int UTFlength(String str) {
+ int strlen = str.length();
+ int utflen = 0;
+ for (int i = 0; i < strlen; i++) {
+ int c = str.charAt(i);
+ if ((c >= 0x0001) && (c <= 0x007F)) {
+ utflen++;
+ } else if (c > 0x07FF) {
+ utflen += 3;
+ } else {
+ utflen += 2;
+ }
+ }
+ return utflen;
+ }
+ @Override
+ public String toString() {
+ return "LINES\n"+lines;
+ }
+
+ @Override
+ public long getLength() {
+ return super.getLength()+2+UTFlength(lines);
+ }
+
+ @Override
+ protected void writeToStream(DataOutput dto) throws IOException {
+ super.writeToStream(dto);
+ dto.writeUTF(lines);
+ }
+
+ @Override
+ protected void readFromStream(DataInput din) throws IOException {
+ super.readFromStream(din);
+ lines=din.readUTF();
+ }
+}
diff --git a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java
index 5c32c5b..263a909 100644
--- a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java
+++ b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java
@@ -12,14 +12,15 @@ public class LineDecitionComparator implements Comparator {
public LineDecitionComparator (List krs,int priority) {
for (Iterator iterator = krs.iterator(); iterator.hasNext();) {
KLALBRemoteSocket klalbRemoteSocket = (KLALBRemoteSocket) iterator.next();
- long x=klalbRemoteSocket.getLatency()>>1;
+ long x=klalbRemoteSocket.getMonitor().getLatency()>>1;
+ x+=klalbRemoteSocket.getMonitor().getJitter()>>1;
AtomicLong al=new AtomicLong(0);
klalbRemoteSocket.getSendQueue().forEach((v)->{
if(v.getPriority()>=priority) {
al.addAndGet(v.getPacket().getLength());
}
});
- long speed=klalbRemoteSocket.getOutSpeed();
+ long speed=klalbRemoteSocket.getMonitor(). getOutSpeedMax();
if(speed==0) {
if(al.get()>0) {
x=Long.MAX_VALUE;
@@ -27,6 +28,7 @@ public class LineDecitionComparator implements Comparator {
}else {
x+=al.get()*1000000000/speed;
}
+ //System.out.println(x);
predictedlatencys.put(klalbRemoteSocket, x);
}
}
diff --git a/src/org/kne/cloud/network/klalb/LineManager.java b/src/org/kne/cloud/network/klalb/LineManager.java
new file mode 100644
index 0000000..e251fea
--- /dev/null
+++ b/src/org/kne/cloud/network/klalb/LineManager.java
@@ -0,0 +1,154 @@
+package org.kne.cloud.network.klalb;
+
+import java.io.IOException;
+import java.net.UnknownHostException;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Timer;
+import java.util.TimerTask;
+import java.util.TreeMap;
+import java.util.TreeSet;
+import java.util.Hashtable;
+
+import org.kne.cloud.network.klalb.LineManager.LineEntry;
+import org.kne.cloud.network.mport.MultipurposeSocketAddress;
+import org.kne.cloud.network.mport.ThreadTool;
+
+public class LineManager extends Hashtable{
+ private Object lock=new Object();
+ private KLALBController controller;
+ public LineManager(KLALBController controller) {
+ super();
+ this.controller = controller;
+ }
+ /*public void applyList(HostPortList hpl) {
+ HostPortList lh=(HostPortList) hpl.clone();
+ for (Iterator iterator = hm.iterator(); iterator.hasNext();) {
+ LineEntry hostPort = (LineEntry) iterator.next();
+ lh.remove(hostPort.getHp());
+ }
+ for (Iterator iterator = lh.iterator(); iterator.hasNext();) {
+ MultipurposeSocketAddress hostPort = (MultipurposeSocketAddress) iterator.next();
+ LineEntry le=new LineEntry(hostPort);
+ hm.add(le);
+ ThreadTool.makeVDaemonThreadIfSupport("线路监视线程", le).start();
+ }
+
+ for (Iterator iterator = hm.iterator(); iterator.hasNext();) {
+ LineEntry hostPort = (LineEntry) iterator.next();
+ if(!hpl.contains(hostPort.getHp())) {
+ hostPort.close();
+ iterator.remove();
+ }
+ }
+
+ }*/
+ public void retry() {
+ synchronized (lock) {
+ lock.notifyAll();
+ }
+ }
+ public synchronized void addHostPort(MultipurposeSocketAddress v) {
+
+ LineEntry le=new LineEntry(v);
+ if(putIfAbsent(v,le)==null)
+ le.open();
+ }
+ public synchronized void removeHostPort(MultipurposeSocketAddress v) {
+ for (Iterator> iterator = entrySet().iterator(); iterator.hasNext();) {
+ java.util.Map.Entry hostPort = iterator.next();
+ if(hostPort.getValue().equals(v)) {
+ hostPort.getValue().close();
+ iterator.remove();
+ }
+ }
+ }
+
+
+
+ public KLALBController getController() {
+ return controller;
+ }
+
+ public class LineEntry implements Runnable{
+ private MultipurposeSocketAddress hp;
+ private KLALBRemoteSocket krs;
+ private volatile boolean flag=true;
+ private volatile boolean online=false;
+ private Monitor monitor=new Monitor();
+ public boolean isOnline() {
+ return online;
+ }
+
+ public MultipurposeSocketAddress getHp() {
+ return hp;
+ }
+
+ public LineEntry(MultipurposeSocketAddress hp) {
+ super();
+ this.hp = hp;
+ }
+ public String toString() {
+ StringBuilder sb=new StringBuilder();
+ sb.append(online?"●在线 ":"○离线 ");
+ sb.append(hp.toString());
+ sb.append('\t');
+ sb.append(monitor.toString());
+ return sb.toString();
+ }
+
+ public String toString2() {
+ StringBuilder sb=new StringBuilder();
+
+ sb.append(hp.toString());
+ sb.append('\n');
+ sb.append(online?"●在线\t":"○离线\t");
+ sb.append(monitor.toString2());
+ return sb.toString();
+ }
+
+ @Override
+ public void run() {
+ while(true){
+ try {
+ krs=new KLALBRemoteSocket(hp.connectSocket(),monitor);
+ //System.out.println("open");
+ controller.addRemoteSocket(krs);
+ online=true;
+ krs.waitforlose();
+ //System.out.println("close");
+ online=false;
+ } catch (UnknownHostException e) {
+ } catch (IOException e) {
+ //e.printStackTrace();
+ }
+ if(!flag) {
+ break;
+ }
+ try {
+ synchronized (lock) {
+ lock.wait(10000);
+ }
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ }
+
+ }
+ public void open() {
+
+ ThreadTool.makeVDaemonThreadIfSupport("线路监视线程", this).start();
+ }
+ public void close() {
+ flag=false;
+ if(krs!=null)
+ try {
+ krs.close();
+ } catch (IOException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }
+
+ }
+ }
+}
diff --git a/src/org/kne/cloud/network/klalb/Monitor.java b/src/org/kne/cloud/network/klalb/Monitor.java
new file mode 100644
index 0000000..cf2eda3
--- /dev/null
+++ b/src/org/kne/cloud/network/klalb/Monitor.java
@@ -0,0 +1,106 @@
+package org.kne.cloud.network.klalb;
+
+import java.util.concurrent.atomic.AtomicLong;
+import static org.kne.cloud.network.klalb.KLALBUtils.*;
+public class Monitor {
+
+ private AtomicLong inTraffic=new AtomicLong();
+ private AtomicLong outTraffic=new AtomicLong();
+
+
+ private AtomicLong inTrafficOld=new AtomicLong();
+ private AtomicLong outTrafficOld=new AtomicLong();
+
+ private volatile long inSpeed;
+ private volatile long outSpeed;
+
+ private volatile long inSpeedMax;
+ private volatile long outSpeedMax;
+
+ private volatile long latency ;
+ private volatile long jitter ;
+ public long getInSpeedMax() {
+ return inSpeedMax;
+ }
+ public void updateLatency(long newlatency) {
+ long vj=Math.abs( latency-newlatency);;
+ if(jitterinSpeedMax) {
+ inSpeedMax=inSpeed;
+ }else {
+ inSpeedMax=(inSpeedMax*999+inSpeed)/1000;
+ }
+ if(outSpeed>outSpeedMax) {
+ outSpeedMax=outSpeed;
+ }else {
+ outSpeedMax=(outSpeedMax*999+outSpeed)/1000;
+ }
+ }
+
+ public long getLatency() {
+ return latency;
+ }
+ public String toString() {
+ StringBuilder sb=new StringBuilder();
+ sb.append(bytesUnit(outTraffic.get())).append("\u2191\t").append(bytesUnit(inTraffic.get())).append("\u2193\t").append(bytesUnit(outSpeed)).append("/s\u2191\t").append(bytesUnit(inSpeed)).append("/s\u2193\t").append(latency/1000000).append("ms");
+ return sb.toString();
+ }
+
+ public String toString2() {
+ StringBuilder sb=new StringBuilder();
+ sb.append(bytesUnit(outTraffic.get())).append("\t").append(bytesUnit(inTraffic.get())).append("\t").append(bytesUnit(outSpeed)).append("/s\t").append(bytesUnit(inSpeed)).append("/s\t").append(latency/1000000).append("ms").append("\t").append(jitter/1000000).append("ms");
+ return sb.toString();
+ }
+
+}
diff --git a/src/org/kne/cloud/network/klalb/MonitoredInputStream.java b/src/org/kne/cloud/network/klalb/MonitoredInputStream.java
index dd8b082..f5db60a 100644
--- a/src/org/kne/cloud/network/klalb/MonitoredInputStream.java
+++ b/src/org/kne/cloud/network/klalb/MonitoredInputStream.java
@@ -1,15 +1,15 @@
package org.kne.cloud.network.klalb;
+import java.io.FilterInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.util.concurrent.atomic.AtomicLong;
-public class MonitoredInputStream extends InputStream {
+public class MonitoredInputStream extends FilterInputStream {
- private InputStream in;
private AtomicLong[] total;
public MonitoredInputStream(InputStream inputStream,AtomicLong... v) {
- this.in=inputStream;
+ super(inputStream);
this.total=v;
}
diff --git a/src/org/kne/cloud/network/klalb/MonitoredOutputStream.java b/src/org/kne/cloud/network/klalb/MonitoredOutputStream.java
index 16764b8..49d4a39 100644
--- a/src/org/kne/cloud/network/klalb/MonitoredOutputStream.java
+++ b/src/org/kne/cloud/network/klalb/MonitoredOutputStream.java
@@ -1,15 +1,15 @@
package org.kne.cloud.network.klalb;
+import java.io.FilterOutputStream;
import java.io.IOException;
import java.io.OutputStream;
import java.util.concurrent.atomic.AtomicLong;
-public class MonitoredOutputStream extends OutputStream {
+public class MonitoredOutputStream extends FilterOutputStream {
- private OutputStream out;
private AtomicLong[] total;
public MonitoredOutputStream(OutputStream outputStream,AtomicLong ...v) {
- this.out=outputStream;
+ super(outputStream);
this.total=v;
}
diff --git a/src/org/kne/cloud/network/klalb/MonitoredSocket.java b/src/org/kne/cloud/network/klalb/MonitoredSocket.java
index 58cc454..a19e8d8 100644
--- a/src/org/kne/cloud/network/klalb/MonitoredSocket.java
+++ b/src/org/kne/cloud/network/klalb/MonitoredSocket.java
@@ -11,68 +11,49 @@ import java.util.concurrent.atomic.AtomicLong;
import org.kne.cloud.network.mport.FilterSocket;
public class MonitoredSocket extends FilterSocket {
+ public Monitor getMonitor() {
+ return monitor;
+ }
+
+
+
private static Timer t=new Timer("带宽测量线程",true);
-
+ private Monitor monitor;
private TimerTask ptt=new TimerTask() {
@Override
public void run() {
- runMonitor();
+ runMonitor();
}
};
public MonitoredSocket(Socket socket) {
+ this(socket,new Monitor());
+ }
+ public MonitoredSocket(Socket socket,Monitor monitor) {
super(socket);
t.scheduleAtFixedRate(ptt, 1000, 1000);
+ this.monitor=monitor;
}
protected void runMonitor() {
- inSpeed.set( inTraffic.get()-inTrafficOld.get());
- outSpeed.set( outTraffic.get()-outTrafficOld.get());
- inTrafficOld.set(inTraffic.get());
- outTrafficOld.set(outTraffic.get());
+ monitor.updateSpeed();
}
- public long getInTraffic() {
- return inTraffic.get();
- }
-
- public long getOutTraffic() {
- return outTraffic.get();
- }
-
-
-
- public long getInSpeed() {
- return inSpeed.get();
- }
-
-
- public long getOutSpeed() {
- return outSpeed.get();
- }
+
private InputStream in;
private OutputStream out;
- private AtomicLong inTraffic=new AtomicLong();
- private AtomicLong outTraffic=new AtomicLong();
-
-
- private AtomicLong inTrafficOld=new AtomicLong();
- private AtomicLong outTrafficOld=new AtomicLong();
-
- private AtomicLong inSpeed=new AtomicLong();
- private AtomicLong outSpeed=new AtomicLong();
@Override
public InputStream getInputStream() throws IOException {
if(in==null)
- in=new MonitoredInputStream(socket.getInputStream(), inTraffic);
+ in=new MonitoredInputStream(socket.getInputStream(), monitor.getInTrafficAL());
return in;
}
@Override
public OutputStream getOutputStream() throws IOException {
if(out==null) {
- out=new MonitoredOutputStream(socket.getOutputStream(), outTraffic);
+ out=new MonitoredOutputStream(socket.getOutputStream(), monitor.getOutTrafficAL());
}
return out;
}
diff --git a/src/org/kne/cloud/network/klalb/PINGPacket.java b/src/org/kne/cloud/network/klalb/PINGPacket.java
index eb748b2..6a2fec7 100644
--- a/src/org/kne/cloud/network/klalb/PINGPacket.java
+++ b/src/org/kne/cloud/network/klalb/PINGPacket.java
@@ -1,13 +1,41 @@
package org.kne.cloud.network.klalb;
+import java.io.DataInput;
+import java.io.DataOutput;
+import java.io.IOException;
+
public class PINGPacket extends KLALBPacket {
-
+ private long time;
+
+ @Override
+ protected void writeToStream(DataOutput dto) throws IOException {
+ super.writeToStream(dto);
+ dto.writeLong(time);
+ }
+
+ public long getTime() {
+ return time;
+ }
+
+ @Override
+ protected void readFromStream(DataInput din) throws IOException {
+ super.readFromStream(din);
+ time=din.readLong();
+ }
+
+ @Override
+ public long getLength() {
+ return super.getLength()+8;
+ }
public PINGPacket() {
super(PING);
}
-
+ public PINGPacket(long time) {
+ super(PING);
+ this.time=time;
+ }
@Override
public String toString() {
return "PING";
diff --git a/src/org/kne/cloud/network/klalb/PONGPacket.java b/src/org/kne/cloud/network/klalb/PONGPacket.java
index a110faa..aa45750 100644
--- a/src/org/kne/cloud/network/klalb/PONGPacket.java
+++ b/src/org/kne/cloud/network/klalb/PONGPacket.java
@@ -1,12 +1,43 @@
package org.kne.cloud.network.klalb;
-public class PONGPacket extends KLALBPacket {
+import java.io.DataInput;
+import java.io.DataOutput;
+import java.io.IOException;
+public class PONGPacket extends KLALBPacket {
+ public long getTime() {
+ return time;
+ }
+ private long time;
+
+ @Override
+ protected void writeToStream(DataOutput dto) throws IOException {
+ super.writeToStream(dto);
+ dto.writeLong(time);
+ }
+
+ @Override
+ protected void readFromStream(DataInput din) throws IOException {
+ super.readFromStream(din);
+ time=din.readLong();
+ }
+
+ @Override
+ public long getLength() {
+ return super.getLength()+8;
+ }
public PONGPacket() {
super(PONG);
}
+ public PONGPacket(long time) {
+ super(PONG);
+ this.time=time;
+ }
@Override
public String toString() {
return "PONG";
}
+ public void redeltaTime(long delta) {
+ time+=delta;
+ }
}
diff --git a/src/org/kne/cloud/network/klalb/RSTPacket.java b/src/org/kne/cloud/network/klalb/RSTPacket.java
index 0b6dd8a..6a0ba15 100644
--- a/src/org/kne/cloud/network/klalb/RSTPacket.java
+++ b/src/org/kne/cloud/network/klalb/RSTPacket.java
@@ -18,7 +18,7 @@ public class RSTPacket extends KLALBPacket {
}
public RSTPacket() {
- super();
+ super(RST);
}
public int getSport() {
diff --git a/src/org/kne/cloud/network/klalb/SACKTPacket.java b/src/org/kne/cloud/network/klalb/SACKTPacket.java
index 1ad5652..1561b02 100644
--- a/src/org/kne/cloud/network/klalb/SACKTPacket.java
+++ b/src/org/kne/cloud/network/klalb/SACKTPacket.java
@@ -26,7 +26,7 @@ public class SACKTPacket extends KLALBPacket {
}
public SACKTPacket() {
- super();
+ super(SACKT);
}
@Override
diff --git a/src/org/kne/cloud/network/klalb/SYNTPacket.java b/src/org/kne/cloud/network/klalb/SYNTPacket.java
index b22e319..eccd342 100644
--- a/src/org/kne/cloud/network/klalb/SYNTPacket.java
+++ b/src/org/kne/cloud/network/klalb/SYNTPacket.java
@@ -21,7 +21,7 @@ public class SYNTPacket extends KLALBPacket {
}
public SYNTPacket() {
- super();
+ super(SYNT);
}
@Override
diff --git a/src/org/kne/cloud/network/klalb/SendTask.java b/src/org/kne/cloud/network/klalb/SendTask.java
index 2adf250..5d078e6 100644
--- a/src/org/kne/cloud/network/klalb/SendTask.java
+++ b/src/org/kne/cloud/network/klalb/SendTask.java
@@ -1,20 +1,33 @@
package org.kne.cloud.network.klalb;
+import java.io.IOException;
import java.net.Inet6Address;
import java.net.NoRouteToHostException;
import java.util.Timer;
import java.util.TimerTask;
+import java.util.concurrent.atomic.AtomicInteger;
-public class SendTask extends TimerTask {
+public class SendTask {
private int maxcount;
- private int count;
+ private AtomicInteger count=new AtomicInteger();
private int priority;
private Inet6Address address;
+ private volatile long starttime=System.nanoTime();
+ private void resetTime() {
+ starttime=System.nanoTime();
+ }
+ private boolean checkTime(long limit) {
+ long curr=System.nanoTime();
+ boolean b=curr-starttime>limit;
+ if(b)
+ starttime=curr;
+ return b;
+ }
public SendTask( KLALBController controllker,Inet6Address address,DATATPacket kp,int priority,int maxcount) {
super();
this.maxcount = maxcount;
this.controllker = controllker;
- this.kp = kp;
+ this.packet = kp;
this.address=address;
this.priority=priority;
}
@@ -28,13 +41,13 @@ public class SendTask extends TimerTask {
}
- public DATATPacket getKp() {
- return kp;
+ public DATATPacket getPacket() {
+ return packet;
}
private KLALBController controllker;
- private DATATPacket kp;
+ private DATATPacket packet;
public int getMaxcount() {
@@ -43,25 +56,34 @@ public class SendTask extends TimerTask {
-
-
- @Override
- public void run() {
+ public void check(int number) throws IOException {
+ long limit= (1<=maxcount) {
+ throw new IOException("send error!");
+ }
try {
+ if(count.get()==0) {
controllker.sendPacketToAddress((Inet6Address) address,
- kp, priority+count);
+ packet, priority+count.get(),1);
+ }else {
+ controllker.sendPacketToAddress((Inet6Address) address,
+ packet, priority+count.get(),2);
+ }
} catch (NoRouteToHostException e) {
e.printStackTrace();
}
- count++;
- if(count>1) {
- System.out.println("重传:"+kp);
- }
- if(count>=maxcount) {
- cancel();
+ resetTime();
+ count.incrementAndGet();
+ if(count.get()>1) {
+ System.out.println("第"+(count.get()-1)+"次重传:"+packet);
}
+
}
- public void addToTimer(int timeout) {
- controllker.getTimer().scheduleAtFixedRate(this, 0, timeout);
- }
+
}
diff --git a/src/org/kne/cloud/network/klalb/ServerPropties.java b/src/org/kne/cloud/network/klalb/ServerPropties.java
new file mode 100644
index 0000000..4c2359b
--- /dev/null
+++ b/src/org/kne/cloud/network/klalb/ServerPropties.java
@@ -0,0 +1,58 @@
+package org.kne.cloud.network.klalb;
+
+import java.io.File;
+import java.io.FileNotFoundException;
+import java.io.FileReader;
+import java.io.IOException;
+import java.io.PrintStream;
+import java.net.Inet6Address;
+import java.net.UnknownHostException;
+import java.util.Properties;
+import java.util.UUID;
+
+public class ServerPropties extends Properties {
+
+ /**
+ *
+ */
+ private static final long serialVersionUID = 1L;
+ public ServerPropties() throws IOException {
+ super();
+ load();
+ }
+ public void load() throws IOException {
+ File f=new File("kserver.ini");
+ if(!f.exists()||f.length()==0) {
+ f.createNewFile();
+ PrintStream ps=new PrintStream(f);
+ ps.println("virtualip="+ KLALBUtils.uuidToIP(UUID.randomUUID()).getHostAddress());
+ ps.println("bind=0.0.0.0:4568");
+ ps.println("local=127.0.0.1:25565");
+ ps.close();
+ }
+ FileReader fr = null;
+ try {
+ fr=new FileReader(f);
+ load(fr);
+ }finally {
+ if(fr!=null)
+ fr.close();
+ }
+ }
+ public Inet6Address getVirtualIP() {
+ try {
+ return (Inet6Address) Inet6Address.getByName((String) get("virtualip"));
+ } catch (UnknownHostException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }
+ return null;
+ }
+ public String getBind() {
+ return (String) get("bind");
+ }
+
+ public String getLocal() {
+ return (String) get("local");
+ }
+}
diff --git a/src/org/kne/cloud/network/klalb/SimpleKLALBClient.java b/src/org/kne/cloud/network/klalb/SimpleKLALBClient.java
new file mode 100644
index 0000000..a825ce3
--- /dev/null
+++ b/src/org/kne/cloud/network/klalb/SimpleKLALBClient.java
@@ -0,0 +1,73 @@
+package org.kne.cloud.network.klalb;
+
+import java.io.File;
+import java.io.IOException;
+import java.net.ServerSocket;
+import java.net.UnknownHostException;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Properties;
+import java.util.Scanner;
+
+import org.kne.cloud.network.klalb.LineManager.LineEntry;
+import org.kne.cloud.network.mport.MultipurposeSocketAddress;
+import org.kne.cloud.network.mport.SocketBridge;
+import org.kne.cloud.network.mport.TCPListener;
+import org.kne.util.AutoProperties;
+
+public class SimpleKLALBClient {
+public static void main(String[] args) throws UnknownHostException, IOException {
+ Scanner scn=new Scanner(System.in);
+ Properties def=new Properties();
+ def.setProperty("server", "");
+ def.setProperty("local", "");
+ AutoProperties ap=new AutoProperties(new File("klalbclient.ini"),def);
+ KLALBController kc=new KLALBController();
+ KLALBRemoteSocket kr=new KLALBRemoteSocket(new MultipurposeSocketAddress(ap.getProperty("server")).connectSocket());
+ kc.addRemoteSocket(kr);
+ kr.close();
+ System.out.println("连接成功:"+ap.getProperty("server"));
+ TCPListener tl=new TCPListener(new MultipurposeSocketAddress(ap.getProperty("local")));
+ tl.setCon((s)->{
+ KLALBVirtualSocket kvs=null;
+ try {
+ kvs=new KLALBVirtualSocket(kc, kr.getRemoteVaddr(), 23333);
+ new SocketBridge(s, kvs).run();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }finally {
+
+ try {
+ s.close();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ if(kvs!=null)
+ try {
+ kvs.close();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ tl.open();
+ System.out.println("提示:输入state并回车可以查看当前线路状态");
+ while(true) {
+ String s=scn.nextLine();
+ switch(s) {
+ case "state":
+ System.out.println("状态\t上传流量\t下载流量\t上传速度\t下载速度\t延迟\t抖动");
+ LineManager le=kc.getLineManager();
+ for (Iterator> iterator = le.entrySet().iterator(); iterator.hasNext();) {
+ java.util.Map.Entry hostPort = iterator.next();
+ System.out.println(hostPort.getValue() .toString2());
+ //System.out.println();
+ }
+ break;
+ default :
+ System.out.println("未知命令!");
+
+ }
+ }
+}
+}
diff --git a/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java
new file mode 100644
index 0000000..8a2544b
--- /dev/null
+++ b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java
@@ -0,0 +1,180 @@
+package org.kne.cloud.network.klalb;
+
+import java.io.BufferedReader;
+import java.io.File;
+import java.io.FileReader;
+import java.io.IOException;
+import java.net.InetAddress;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.net.UnknownHostException;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Scanner;
+
+import org.kne.cloud.network.klalb.LineManager.LineEntry;
+import org.kne.cloud.network.mport.MultipurposeSocketAddress;
+import org.kne.cloud.network.mport.Protocol;
+import org.kne.cloud.network.mport.ProtocolDetectorServerSocketFactory;
+import org.kne.cloud.network.mport.ProtocolDetectorSocket;
+import org.kne.cloud.network.mport.ProxyProfileAnalyser;
+import org.kne.cloud.network.mport.ProxyProfileExecutor;
+import org.kne.cloud.network.mport.SocketBridge;
+import org.kne.cloud.network.mport.TCPListener;
+
+public class SimpleKLALBServer {
+ public static TCPListener tcpl,tcpl2;
+ public static ServerPropties sp;
+ public static KLALBController kc;
+ static {
+ MultipurposeSocketAddress.getServerSocketFactoryRegister().put("DETTCP", new ProtocolDetectorServerSocketFactory());
+ }
+ public static void main(String[] args) throws IOException {
+
+ System.out.println("KLALB负载均衡V2.0");
+ sp=new ServerPropties();
+
+ System.out.println("虚拟地址:"+sp.getVirtualIP().getHostAddress());
+ kc=new KLALBController(sp.getVirtualIP());
+ kc.setSelflineTableSupplier(()->{
+ return fileRead("linetable.txt");
+ });
+ kc.registerToProxyTypeAs("KLALB");
+
+ System.out.println("开放端口:"+sp.getBind());
+ openPort(sp.getBind());
+ System.out.println("本地服务:"+sp.getLocal());
+ openLocalPort(sp.getLocal());
+
+
+
+ Scanner scn=new Scanner(System.in);
+ while(true) {
+ String s=scn.next();
+ String[]sc=s.split(" ");
+ switch(sc[0]) {
+ case "help":
+ System.out.print("state:查看线路状态");
+ System.out.println("reload:重新加载线路配置");
+ break;
+ case "state":
+ LineManager le=kc.getLineManager();
+ for (Iterator> iterator = le.entrySet().iterator(); iterator.hasNext();) {
+ java.util.Map.Entry hostPort = iterator.next();
+ System.out.println(hostPort.getValue() .toString());
+ }
+ break;
+ default:
+ System.out.println("未知命令,请输入help以查询指令说明");
+ }
+ }
+ }
+ private static void openPort(String bip) throws IOException {
+ if(tcpl!=null)
+ tcpl.close();
+ MultipurposeSocketAddress mpsa=new MultipurposeSocketAddress(bip,"DETTCP");
+ tcpl=new TCPListener(mpsa);
+ tcpl.setCon((soc)->{
+ ProtocolDetectorSocket pds=(ProtocolDetectorSocket) soc;
+ if(!pds.getProtocolStack().isEmpty()&&pds.getProtocolStack().pop().getName().equals("KLALB")) {
+ KLALBRemoteSocket krs=new KLALBRemoteSocket(pds);
+ kc.addRemoteSocket(krs);
+ }else {
+ MultipurposeSocketAddress mpsa2=new MultipurposeSocketAddress(sp.getLocal());
+ Socket s=null;
+ try {
+ s=mpsa2.connectSocket();
+ new SocketBridge(pds, s).run();
+ } catch (UnknownHostException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ } catch (IOException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }finally {
+ if(s!=null)
+ try {
+ s.close();
+ } catch (IOException e1) {
+ // TODO 自动生成的 catch 块
+ e1.printStackTrace();
+ }
+ try {
+ soc.close();
+ } catch (IOException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }
+ }
+ }
+ });
+ tcpl.open();
+ }
+ private static void openLocalPort(String bip) throws IOException {
+ if(tcpl2!=null)
+ tcpl2.close();
+ MultipurposeSocketAddress mpsa=new MultipurposeSocketAddress(bip);
+ tcpl2=new TCPListener(new MultipurposeSocketAddress("KLALB", "::0", 23333));
+ tcpl2.setCon((soc)->{
+ Socket s=null;
+ try {
+ s=mpsa.connectSocket();
+ new SocketBridge(soc, s).run();
+ } catch (UnknownHostException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ } catch (IOException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }finally {
+ if(s!=null)
+ try {
+ s.close();
+ } catch (IOException e1) {
+ // TODO 自动生成的 catch 块
+ e1.printStackTrace();
+ }
+ try {
+ soc.close();
+ } catch (IOException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }
+ }
+
+ });
+ tcpl2.open();
+ }
+ private static String fileRead(String filePath){
+ //1.定义一个BufferedReader对象,将文件内容读取到缓存
+ BufferedReader bufferedReader =null;
+ String returnInfo="";
+ try{
+ // 2.定义一个file对象
+ File file = new File(filePath);//定义一个file对象,用来初始化FileReader
+ // 3.定义一个fileReader对象
+ FileReader reader = new FileReader(file);
+ // 4.定义一个BufferedReader对象,将文件内容读取到缓存
+ bufferedReader = new BufferedReader(reader);
+ // 5.定义一个字符串缓存,将字符串存放缓存中
+ StringBuilder stringBuilder = new StringBuilder();
+ String str = "";
+ while ((str =bufferedReader.readLine()) != null) {//逐行读取文件内容,不读取换行符和末尾的空格
+ stringBuilder.append(str + "\n");//将读取的字符串添加换行符后累加存放在缓存中
+ }
+ returnInfo = stringBuilder.toString();
+ }catch (Exception e){
+ e.printStackTrace();
+ }finally {
+ try{
+ if(bufferedReader!=null){
+ bufferedReader.close();
+ }
+ }catch (Exception e){
+ e.printStackTrace();
+ }
+ }
+ return returnInfo;
+ }
+
+}
diff --git a/src/org/kne/cloud/network/klalb/VADDRPacket.java b/src/org/kne/cloud/network/klalb/VADDRPacket.java
index 1fa903d..2c861a2 100644
--- a/src/org/kne/cloud/network/klalb/VADDRPacket.java
+++ b/src/org/kne/cloud/network/klalb/VADDRPacket.java
@@ -17,12 +17,12 @@ public class VADDRPacket extends KLALBPacket {
}
public VADDRPacket() {
- super();
+ super(VADDR);
}
@Override
public String toString() {
- return "VADDR"+vaddr;
+ return "VADDR "+vaddr.getHostAddress();
}
@Override
diff --git a/src/org/kne/cloud/network/mport/DefaultServerSocketFactory.java b/src/org/kne/cloud/network/mport/DefaultServerSocketFactory.java
index 093fd60..7c94200 100644
--- a/src/org/kne/cloud/network/mport/DefaultServerSocketFactory.java
+++ b/src/org/kne/cloud/network/mport/DefaultServerSocketFactory.java
@@ -10,19 +10,16 @@ public class DefaultServerSocketFactory extends ServerSocketFactory {
@Override
public ServerSocket createServerSocket(int port) throws IOException {
- // TODO 自动生成的方法存根
return new ServerSocket(port);
}
@Override
public ServerSocket createServerSocket(int port, int backlog) throws IOException {
- // TODO 自动生成的方法存根
return new ServerSocket(port, backlog);
}
@Override
public ServerSocket createServerSocket(int port, int backlog, InetAddress ifAddress) throws IOException {
- // TODO 自动生成的方法存根
return new ServerSocket(port, backlog, ifAddress);
}
diff --git a/src/org/kne/cloud/network/mport/DefaultSocketFactory.java b/src/org/kne/cloud/network/mport/DefaultSocketFactory.java
new file mode 100644
index 0000000..e22f23a
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/DefaultSocketFactory.java
@@ -0,0 +1,31 @@
+package org.kne.cloud.network.mport;
+
+import java.io.IOException;
+import java.net.InetAddress;
+import java.net.Socket;
+import java.net.UnknownHostException;
+import javax.net.SocketFactory;
+
+public class DefaultSocketFactory extends SocketFactory {
+ public Socket createSocket() {
+ return new Socket();
+ }
+
+ public Socket createSocket(String paramString, int paramInt) throws IOException, UnknownHostException {
+ return new Socket(paramString, paramInt);
+ }
+
+ public Socket createSocket(InetAddress paramInetAddress, int paramInt) throws IOException {
+ return new Socket(paramInetAddress, paramInt);
+ }
+
+ public Socket createSocket(String paramString, int paramInt1, InetAddress paramInetAddress, int paramInt2)
+ throws IOException, UnknownHostException {
+ return new Socket(paramString, paramInt1, paramInetAddress, paramInt2);
+ }
+
+ public Socket createSocket(InetAddress paramInetAddress1, int paramInt1, InetAddress paramInetAddress2,
+ int paramInt2) throws IOException {
+ return new Socket(paramInetAddress1, paramInt1, paramInetAddress2, paramInt2);
+ }
+}
diff --git a/src/org/kne/cloud/network/mport/HostPort.java b/src/org/kne/cloud/network/mport/HostPort.java
deleted file mode 100644
index 1671d63..0000000
--- a/src/org/kne/cloud/network/mport/HostPort.java
+++ /dev/null
@@ -1,63 +0,0 @@
-package org.kne.cloud.network.mport;
-
-import java.io.Serializable;
-import java.net.Inet6Address;
-import java.net.InetAddress;
-import java.net.InetSocketAddress;
-import java.net.SocketAddress;
-import java.net.UnknownHostException;
-import java.util.Objects;
-
-public class HostPort implements Serializable{
- /**
- *
- */
- private static final long serialVersionUID = 1L;
- private String host;
- private int port;
- public HostPort(String hostport) {
- int index =hostport.lastIndexOf(":");
- host=hostport.substring(0,index);
- port=Integer.parseInt(hostport.substring(index+1));
-
- }
- @Override
- public int hashCode() {
- return Objects.hash(host, port);
- }
- @Override
- public boolean equals(Object obj) {
- if (this == obj)
- return true;
- if (obj == null)
- return false;
- if (getClass() != obj.getClass())
- return false;
- HostPort other = (HostPort) obj;
- return Objects.equals(host, other.host) && port == other.port;
- }
- public HostPort(String host2, int port2) {
- host=host2;
- port=port2;
- }
- public HostPort(InetSocketAddress remoteSocketAddress) {
- host=remoteSocketAddress.getAddress().getHostAddress();
- port=remoteSocketAddress.getPort();
- }
- public String getHost() {
- return host;
- }
- public int getPort() {
- return port;
- }
- @Override
- public String toString() {
- if(host.contains(":")) {
- return "["+host+"]:"+port;
- }
- return host+":"+port;
- }
- public SocketAddress getSocketAddress() {
- return new InetSocketAddress(host, port);
- }
-}
diff --git a/src/org/kne/cloud/network/mport/HostPortMap.java b/src/org/kne/cloud/network/mport/HostPortMap.java
index e6c17a0..5bb3530 100644
--- a/src/org/kne/cloud/network/mport/HostPortMap.java
+++ b/src/org/kne/cloud/network/mport/HostPortMap.java
@@ -1,8 +1,9 @@
package org.kne.cloud.network.mport;
import java.util.HashMap;
+import java.util.LinkedHashMap;
-public class HostPortMap extends HashMap {
+public class HostPortMap extends LinkedHashMap {
/**
*
@@ -10,9 +11,9 @@ public class HostPortMap extends HashMap {
private static final long serialVersionUID = 1L;
- public HostPort putN$H(String str) {
+ public MultipurposeSocketAddress putN$H(String str) {
String[] strx=str.split("\\$");
- return super.put(strx[0], new HostPort( strx[1]));
+ return super.put(strx[0], new MultipurposeSocketAddress( strx[1]));
}
}
diff --git a/src/org/kne/cloud/network/mport/MultipurposeSocketAddress.java b/src/org/kne/cloud/network/mport/MultipurposeSocketAddress.java
new file mode 100644
index 0000000..9eb24b8
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/MultipurposeSocketAddress.java
@@ -0,0 +1,159 @@
+package org.kne.cloud.network.mport;
+
+import java.io.IOException;
+import java.io.Serializable;
+import java.net.Inet6Address;
+import java.net.InetAddress;
+import java.net.InetSocketAddress;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.net.SocketAddress;
+import java.net.UnknownHostException;
+import java.util.*;
+import java.util.Objects;
+
+import javax.net.ServerSocketFactory;
+import javax.net.SocketFactory;
+
+public class MultipurposeSocketAddress implements Serializable{
+ /**
+ *
+ */
+ private static Map socketFactoryRegister=new HashMap<>();
+ private static Map serverSocketFactoryRegister=new HashMap<>();
+ static {
+ socketFactoryRegister.put("TCP", new DefaultSocketFactory());
+ serverSocketFactoryRegister.put("TCP", new DefaultServerSocketFactory());
+ }
+
+ public static Map getSocketFactoryRegister() {
+ return socketFactoryRegister;
+ }
+ public static Map getServerSocketFactoryRegister() {
+ return serverSocketFactoryRegister;
+ }
+
+ private static final long serialVersionUID = 1L;
+ private String type;
+ private String host;
+ private ProtocolStack ps=new ProtocolStack();
+
+
+ private int port;
+
+ public MultipurposeSocketAddress(String hostport) {
+ this(hostport,"TCP");
+ }
+ public MultipurposeSocketAddress(String hostport,String defaulttype) {
+ hostport=hostport.trim();
+ if (hostport.startsWith("{")) {
+ int i=hostport.indexOf("}");
+ if(i==-1) {
+ throw new IllegalArgumentException("missing }");
+ }
+ type=hostport.substring(1, i);
+ hostport=hostport.substring(i+1);
+ }else {
+ type=defaulttype;
+ }
+ int index =hostport.lastIndexOf(":");
+ host=hostport.substring(0,index);
+ host=host.replace("[", "");
+ host=host.replace("]", "");
+ String tps=hostport.substring(index+1);
+ int v=tps.indexOf("(");
+ if(v!=-1) {
+ ps.add(new Protocol(tps.substring(v+1,tps.lastIndexOf(")") )));
+ tps=tps.substring(0, v);
+ }
+ port=Integer.parseInt(tps);
+
+ }
+ @Override
+ public int hashCode() {
+ return Objects.hash(host, port, ps, type);
+ }
+ @Override
+ public boolean equals(Object obj) {
+ if (this == obj)
+ return true;
+ if (obj == null)
+ return false;
+ if (getClass() != obj.getClass())
+ return false;
+ MultipurposeSocketAddress other = (MultipurposeSocketAddress) obj;
+ return Objects.equals(host, other.host) && port == other.port && Objects.equals(ps, other.ps)
+ && Objects.equals(type, other.type);
+ }
+ public MultipurposeSocketAddress(String host2, int port2) {
+ this("TCP", host2, port2);
+ }
+
+ public MultipurposeSocketAddress(String type,String host2, int port2) {
+ host=host2;
+ port=port2;
+ this.type=type;
+ }
+ public MultipurposeSocketAddress(InetSocketAddress remoteSocketAddress) {
+ this("TCP",remoteSocketAddress);
+ }
+ public MultipurposeSocketAddress(String type,InetSocketAddress remoteSocketAddress) {
+ host=remoteSocketAddress.getAddress().getHostAddress();
+ port=remoteSocketAddress.getPort();
+ this.type=type;
+ }
+ public String getHost() {
+ return host;
+ }
+ public int getPort() {
+ return port;
+ }
+ @Override
+ public String toString() {
+ StringBuilder sb=new StringBuilder();
+ if(!"TCP".equals(type)) {
+ sb.append('{').append(type).append('}');
+ }
+ if(host.contains(":")) {
+ sb.append('[').append( host).append( "]:").append( port);
+ }else {
+ sb.append( host).append( ':').append( port);
+ }
+ if(!ps.isEmpty()) {
+ sb.append('(').append(ps.pop().getName()).append(')');
+ }
+ return sb.toString();
+ }
+ public Socket connectSocket(InetAddress bindip,int bindport,int timeout) throws UnknownHostException, IOException {
+ Socket s=socketFactoryRegister.get(type).createSocket();
+ s.bind(new InetSocketAddress(bindip, bindport));
+ s.connect(new InetSocketAddress(host, port),timeout);
+ return s;
+ }
+ public Socket connectSocket(InetAddress bindip,int bindport) throws UnknownHostException, IOException {
+ Socket s=socketFactoryRegister.get(type).createSocket();
+ s.bind(new InetSocketAddress(bindip, bindport));
+ s.connect(new InetSocketAddress(host, port));
+ return s;
+ }
+ public Socket connectSocket() throws UnknownHostException, IOException {
+ Socket s=socketFactoryRegister.get(type).createSocket();
+ s.connect(new InetSocketAddress(host, port));
+ return s;
+ }
+ public Socket connectSocket(int timeout) throws UnknownHostException, IOException {
+ Socket s=socketFactoryRegister.get(type).createSocket();
+ s.connect(new InetSocketAddress(host, port),timeout);
+ return s;
+ }
+
+ public ServerSocket listenServerSocket() throws UnknownHostException, IOException {
+ ServerSocket sk=serverSocketFactoryRegister.get(type).createServerSocket(port, 50, InetAddress.getByName(host));
+ return sk;
+ }
+ public ServerSocket listenServerSocket(int backlog) throws UnknownHostException, IOException {
+ ServerSocket sk=serverSocketFactoryRegister.get(type).createServerSocket(port, backlog, InetAddress.getByName(host));
+ return sk;
+ }
+
+}
diff --git a/src/org/kne/cloud/network/mport/NetworkService.java b/src/org/kne/cloud/network/mport/NetworkService.java
new file mode 100644
index 0000000..4ffd5b1
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/NetworkService.java
@@ -0,0 +1,9 @@
+package org.kne.cloud.network.mport;
+
+public interface NetworkService {
+ public void listen(MultipurposeSocketAddress msa);
+ public void unlisten(MultipurposeSocketAddress msc);
+
+ public void connect(MultipurposeSocketAddress msa);
+ public void unconnect(MultipurposeSocketAddress msc);
+}
diff --git a/src/org/kne/cloud/network/mport/PortMultiUse.java b/src/org/kne/cloud/network/mport/PortMultiUse.java
index c672e55..56d292a 100644
--- a/src/org/kne/cloud/network/mport/PortMultiUse.java
+++ b/src/org/kne/cloud/network/mport/PortMultiUse.java
@@ -21,7 +21,7 @@ public class PortMultiUse {
}
Scanner scn=new Scanner(f);
int remp=scn.nextInt();
- System.out.println("端口复用程序V0.2");
+ System.out.println("端口复用程序V2.0");
System.out.println("开放端口:"+remp);
while(scn.hasNext() ) {
String s=scn.next();
diff --git a/src/org/kne/cloud/network/mport/PortRelay.java b/src/org/kne/cloud/network/mport/PortRelay.java
index 7452281..5d2df50 100644
--- a/src/org/kne/cloud/network/mport/PortRelay.java
+++ b/src/org/kne/cloud/network/mport/PortRelay.java
@@ -16,13 +16,13 @@ public class PortRelay {
public PortRelay(int port, HostPortMap services) throws IOException {
this.services = services;
this.port = port;
- ssc=new TCPListener(port);
- ssc.setServerSocketFactory(new ProtocolDetectorServerSocketFactory());
+ MultipurposeSocketAddress.getServerSocketFactoryRegister().put("ProtocolDetectorServerSocket", new ProtocolDetectorServerSocketFactory());
+ ssc=new TCPListener(new MultipurposeSocketAddress("ProtocolDetectorServerSocket", "0.0.0.0", port));
}
public void start() throws IOException {
ssc.setCon((s)->{
- Socket sl = new Socket();
+ Socket sl = null;
ProtocolDetectorSocket pds=(ProtocolDetectorSocket) s;
String nx;
if(pds.getProtocolStack().isEmpty()) {
@@ -30,18 +30,19 @@ public class PortRelay {
}else {
nx=pds.getProtocolStack().pop().getName();
}
- HostPort hp=services.get(nx);
+ MultipurposeSocketAddress hp=services.get(nx);
InetSocketAddress sa=(InetSocketAddress) s.getRemoteSocketAddress();
System.out.println(sa.getAddress().getHostAddress() + ":" + port
+ "-" + "(" + nx + ")->" + hp.getHost() + ":" + hp.getPort());
try {
- sl.connect(hp.getSocketAddress());
+ sl=hp.connectSocket();
new SocketBridge(s, sl).run();
} catch (IOException e) {
e.printStackTrace();
}finally {
try {
+ if(sl!=null)
sl.close();
} catch (IOException e) {
e.printStackTrace();
diff --git a/src/org/kne/cloud/network/mport/Protocol.java b/src/org/kne/cloud/network/mport/Protocol.java
index 94da63e..09098c0 100644
--- a/src/org/kne/cloud/network/mport/Protocol.java
+++ b/src/org/kne/cloud/network/mport/Protocol.java
@@ -33,7 +33,7 @@ public boolean equals(Object obj) {
@Override
public String toString() {
- return "Protocol [name=" + name + "]";
+ return name;
}
diff --git a/src/org/kne/cloud/network/mport/ProxyProfileAnalyser.java b/src/org/kne/cloud/network/mport/ProxyProfileAnalyser.java
new file mode 100644
index 0000000..ff55fb3
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/ProxyProfileAnalyser.java
@@ -0,0 +1,63 @@
+package org.kne.cloud.network.mport;
+
+import java.io.File;
+import java.io.FileNotFoundException;
+import java.io.InputStream;
+import java.util.*;
+
+public class ProxyProfileAnalyser {
+ private Listentrys=new ArrayList();
+ public List getEntrys() {
+ return entrys;
+ }
+ public ProxyProfileAnalyser(InputStream f) {
+ Scanner sr=new Scanner(f);
+ analyse(sr);
+ }
+ public ProxyProfileAnalyser(File f) throws FileNotFoundException {
+ Scanner sr=new Scanner(f);
+ analyse(sr);
+ }
+ public ProxyProfileAnalyser(String s) {
+ Scanner sr=new Scanner(s);
+ analyse(sr);
+ }
+ private void analyse(Scanner sr) {
+ int n=0;
+ while(sr.hasNext()) {
+ String sl=sr.nextLine().trim();
+ int x=sl.indexOf("//");
+ if(x!=-1)
+ sl=sl.substring(0, x);
+ n++;
+ if(sl.isEmpty()) {
+ continue;
+ }
+ String[] sp=sl.split("->");
+ if(sp.length!=2) {
+ throw new IllegalArgumentException("Format error at line "+n);
+ }
+ String l=sp[0].trim();
+ String r=sp[1].trim();
+ //System.out.println(l+" "+r);
+ if(l.contains(":")) {
+ if(r.contains(":")){
+ SocketToSocketProxyProfileEntry sts=new SocketToSocketProxyProfileEntry(l,r);
+ entrys.add(sts);
+
+ }else {
+ SocketToServiceProxyProfileEntry sts=new SocketToServiceProxyProfileEntry(l,r);
+ entrys.add(sts);
+ }
+ }else {
+ if(r.contains(":")) {
+ ServiceToSocketProxyProfileEntry sts=new ServiceToSocketProxyProfileEntry(l,r);
+ entrys.add(sts);
+ }else {
+
+ }
+ }
+ }
+ }
+
+}
diff --git a/src/org/kne/cloud/network/mport/ProxyProfileEntry.java b/src/org/kne/cloud/network/mport/ProxyProfileEntry.java
new file mode 100644
index 0000000..134871a
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/ProxyProfileEntry.java
@@ -0,0 +1,16 @@
+package org.kne.cloud.network.mport;
+
+import java.net.Socket;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.function.Consumer;
+
+import javax.net.ServerSocketFactory;
+import javax.net.SocketFactory;
+
+public class ProxyProfileEntry {
+ private static Map register=new HashMap<>();
+ public static Map getRegister() {
+ return register;
+ }
+}
diff --git a/src/org/kne/cloud/network/mport/ProxyProfileExecutor.java b/src/org/kne/cloud/network/mport/ProxyProfileExecutor.java
new file mode 100644
index 0000000..66d66c0
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/ProxyProfileExecutor.java
@@ -0,0 +1,57 @@
+package org.kne.cloud.network.mport;
+
+import java.io.File;
+import java.io.FileNotFoundException;
+import java.io.InputStream;
+import java.util.ArrayList;
+import java.util.Iterator;
+import java.util.List;
+
+public class ProxyProfileExecutor {
+ private Listcurrent=new ArrayList<>();
+
+ public List getCurrent() {
+ return current;
+ }
+ public void applyNewProxyList(Listl) {
+ for (Iterator iterator = l.iterator(); iterator.hasNext();) {
+ ProxyProfileEntry proxyProfileEntry = (ProxyProfileEntry) iterator.next();
+ if(!current.contains(proxyProfileEntry)) {
+ if(proxyProfileEntry instanceof SocketToSocketProxyProfileEntry) {
+ SocketToSocketProxyProfileEntry stsppe=(SocketToSocketProxyProfileEntry) proxyProfileEntry;
+ }else if(proxyProfileEntry instanceof SocketToServiceProxyProfileEntry){
+ SocketToServiceProxyProfileEntry stsppe=(SocketToServiceProxyProfileEntry) proxyProfileEntry;
+ stsppe.getDes().listen(stsppe.getSrc());
+ }else if(proxyProfileEntry instanceof ServiceToSocketProxyProfileEntry) {
+ ServiceToSocketProxyProfileEntry stsppe=(ServiceToSocketProxyProfileEntry) proxyProfileEntry;
+ stsppe.getSrc().connect(stsppe.getDes());
+ }
+ }
+ }
+ for (Iterator iterator = current.iterator(); iterator.hasNext();) {
+ ProxyProfileEntry proxyProfileEntry = (ProxyProfileEntry) iterator.next();
+ if(!l.contains(proxyProfileEntry)) {
+ if(proxyProfileEntry instanceof SocketToSocketProxyProfileEntry) {
+ SocketToSocketProxyProfileEntry stsppe=(SocketToSocketProxyProfileEntry) proxyProfileEntry;
+ }else if(proxyProfileEntry instanceof SocketToServiceProxyProfileEntry){
+ SocketToServiceProxyProfileEntry stsppe=(SocketToServiceProxyProfileEntry) proxyProfileEntry;
+ stsppe.getDes().unlisten(stsppe.getSrc());
+ }else if(proxyProfileEntry instanceof ServiceToSocketProxyProfileEntry) {
+ ServiceToSocketProxyProfileEntry stsppe=(ServiceToSocketProxyProfileEntry) proxyProfileEntry;
+ stsppe.getSrc().unconnect(stsppe.getDes());
+ }
+ }
+ }
+ current.clear();
+ current.addAll(l);
+ }
+ public void load(File f) throws FileNotFoundException {
+ applyNewProxyList(new ProxyProfileAnalyser(f).getEntrys());
+ }
+ public void load(String s) {
+ applyNewProxyList(new ProxyProfileAnalyser(s).getEntrys());
+ }
+ public void load(InputStream in) {
+ applyNewProxyList(new ProxyProfileAnalyser(in).getEntrys());
+ }
+}
diff --git a/src/org/kne/cloud/network/mport/ServiceElement.java b/src/org/kne/cloud/network/mport/ServiceElement.java
index c5c3731..3bb5cff 100644
--- a/src/org/kne/cloud/network/mport/ServiceElement.java
+++ b/src/org/kne/cloud/network/mport/ServiceElement.java
@@ -9,15 +9,15 @@ import java.util.Map;
public class ServiceElement {
public String name;
public String protocol;
- public HostPort ipport;
+ public MultipurposeSocketAddress ipport;
public ServiceElement(String ini) throws UnknownHostException {
String[]t=ini.split("\\$");
if(t.length==1) {
protocol="";
- ipport=new HostPort(t[0]);
+ ipport=new MultipurposeSocketAddress(t[0]);
}else {
protocol=t[0];
- ipport=new HostPort(t[1]);
+ ipport=new MultipurposeSocketAddress(t[1]);
}
}
public ServiceElement() {
diff --git a/src/org/kne/cloud/network/mport/ServiceToSocketProxyProfileEntry.java b/src/org/kne/cloud/network/mport/ServiceToSocketProxyProfileEntry.java
new file mode 100644
index 0000000..acf96d5
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/ServiceToSocketProxyProfileEntry.java
@@ -0,0 +1,18 @@
+package org.kne.cloud.network.mport;
+
+import java.util.function.Consumer;
+
+public class ServiceToSocketProxyProfileEntry extends ProxyProfileEntry{
+ public NetworkService getSrc() {
+ return src;
+ }
+ public MultipurposeSocketAddress getDes() {
+ return des;
+ }
+ public ServiceToSocketProxyProfileEntry(String l, String r) {
+ src=getRegister().get(l.substring(1, l.length()-1));
+ des=new MultipurposeSocketAddress(r);
+ }
+ private NetworkService src;
+ private MultipurposeSocketAddress des;
+}
diff --git a/src/org/kne/cloud/network/mport/SocketToServiceProxyProfileEntry.java b/src/org/kne/cloud/network/mport/SocketToServiceProxyProfileEntry.java
new file mode 100644
index 0000000..a0cc947
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/SocketToServiceProxyProfileEntry.java
@@ -0,0 +1,19 @@
+package org.kne.cloud.network.mport;
+
+import java.net.Socket;
+import java.util.function.Consumer;
+
+public class SocketToServiceProxyProfileEntry extends ProxyProfileEntry {
+ public SocketToServiceProxyProfileEntry(String l, String r) {
+ src=new MultipurposeSocketAddress(l);
+ des=getRegister().get(r.substring(1, r.length()-1)) ;
+ }
+ public MultipurposeSocketAddress getSrc() {
+ return src;
+ }
+ public NetworkService getDes() {
+ return des;
+ }
+ private MultipurposeSocketAddress src;
+ private NetworkService des;
+}
diff --git a/src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java b/src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java
new file mode 100644
index 0000000..2da4c94
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java
@@ -0,0 +1,16 @@
+package org.kne.cloud.network.mport;
+
+public class SocketToSocketProxyProfileEntry extends ProxyProfileEntry {
+ public SocketToSocketProxyProfileEntry(String l, String r) {
+ src=new MultipurposeSocketAddress(l);
+ des=new MultipurposeSocketAddress(r);
+ }
+ private MultipurposeSocketAddress src;
+ private MultipurposeSocketAddress des;
+ public MultipurposeSocketAddress getSrc() {
+ return src;
+ }
+ public MultipurposeSocketAddress getDes() {
+ return des;
+ }
+}
diff --git a/src/org/kne/cloud/network/mport/TCPListener.java b/src/org/kne/cloud/network/mport/TCPListener.java
index db869b0..4a9da4d 100644
--- a/src/org/kne/cloud/network/mport/TCPListener.java
+++ b/src/org/kne/cloud/network/mport/TCPListener.java
@@ -1,6 +1,7 @@
package org.kne.cloud.network.mport;
import java.io.IOException;
+import java.net.InetAddress;
import java.net.ServerSocket;
import java.net.Socket;
import java.util.function.Consumer;
@@ -8,9 +9,11 @@ import java.util.function.Consumer;
import javax.net.ServerSocketFactory;
public class TCPListener {
- private int port;
- private ServerSocket servers;
- private ServerSocketFactory serverSocketFactory=new DefaultServerSocketFactory();
+ private ServerSocket serverSocket;
+ public ServerSocket getServerSocket() {
+ return serverSocket;
+ }
+
private volatile boolean flag=false;
private Consumercon;
@@ -19,7 +22,7 @@ public class TCPListener {
public void run() {
while(flag){
try {
- Socket soce=servers.accept();
+ Socket soce=serverSocket.accept();
ThreadTool.makeVThreadIfSupport("端口监听线程",()->{
con.accept(soce);
}).start();
@@ -29,41 +32,34 @@ public class TCPListener {
}
}
};
- public TCPListener(int port) throws IOException {
- this.port=port;
- }
+ private MultipurposeSocketAddress multipurposeSocketAddress;
- public ServerSocketFactory getServerSocketFactory() {
- return serverSocketFactory;
- }
-
- public void setServerSocketFactory(ServerSocketFactory serverSocketFactory) {
- this.serverSocketFactory = serverSocketFactory;
+ public TCPListener(MultipurposeSocketAddress multipurposeSocketAddress) {
+ this.multipurposeSocketAddress=multipurposeSocketAddress;
}
public void open() throws IOException {
flag=true;
- servers=serverSocketFactory.createServerSocket(port);
+ serverSocket=multipurposeSocketAddress.listenServerSocket();
//servers=new ServerSocket(port);
new Thread(r).start();
}
public void close() {
flag=false;
- if(servers!=null) {
+ if(serverSocket!=null) {
try {
- servers.close();
+ serverSocket.close();
} catch (IOException e) {
e.printStackTrace();
}
+ serverSocket=null;
}
}
- public int getPort() {
- return port;
- }
- public void setPort(int port) {
- this.port = port;
+
+ public MultipurposeSocketAddress getMultipurposeSocketAddress() {
+ return multipurposeSocketAddress;
}
public Consumer getCon() {
diff --git a/src/org/kne/cloud/network/nathole/TCPNatHoleTeat.java b/src/org/kne/cloud/network/nathole/TCPNatHoleTeat.java
deleted file mode 100644
index 741725f..0000000
--- a/src/org/kne/cloud/network/nathole/TCPNatHoleTeat.java
+++ /dev/null
@@ -1,101 +0,0 @@
-package org.kne.cloud.network.nathole;
-
-import java.io.IOException;
-import java.lang.reflect.Field;
-import java.lang.reflect.InvocationTargetException;
-import java.lang.reflect.Method;
-import java.net.InetAddress;
-import java.net.InetSocketAddress;
-import java.net.ServerSocket;
-import java.net.Socket;
-import java.net.SocketAddress;
-import java.net.SocketImpl;
-import java.net.SocketTimeoutException;
-import java.net.UnknownHostException;
-import java.util.Arrays;
-
-public class TCPNatHoleTeat {
- public static void main(String[] args) throws IOException, NoSuchFieldException, SecurityException, IllegalArgumentException, IllegalAccessException, NoSuchMethodException, InvocationTargetException {
- ServerSocket ssk=new ServerSocket(8089);
- Field si= ssk.getClass().getDeclaredField("impl");
- si.setAccessible(true);
- SocketImpl sii= (SocketImpl) si.get(ssk);
- Method sim= sii.getClass().getDeclaredMethod("connect", SocketAddress.class,int.class);
- sim.setAccessible(true);
- new Thread(()->{
- try {
- Thread.sleep(1000);
- } catch (InterruptedException e1) {
- e1.printStackTrace();
- }
- while(true) {
- try {
- sim.invoke(sii, new InetSocketAddress(InetAddress.getByName("3.3.3.3"), 56) ,1);
- } catch (IllegalAccessException e) {
- // TODO 自动生成的 catch 块
- e.printStackTrace();
- } catch (IllegalArgumentException e) {
- // TODO 自动生成的 catch 块
- e.printStackTrace();
- } catch (InvocationTargetException e) {
- } catch (UnknownHostException e) {
- // TODO 自动生成的 catch 块
- e.printStackTrace();
- }
- }
- }).start();
- ssk.accept();
- //x("183.198.152.240", 12187);
-
- }
-
- private static void x(String nat, int port) throws IOException {
- for (int i = 10000; i < 20000; i++) {
- System.out.println("try" + i);
- Socket s = null;
- ServerSocket srs = null;
- try {
- s = new Socket();
- s.bind(new InetSocketAddress(InetAddress.getByName("0.0.0.0"), i));
- try {
- s.connect(new InetSocketAddress(InetAddress.getByName("1.1.1.1"), 443), 1);
- }catch(SocketTimeoutException e) {
-
- }
- s.close();
-
- srs = new ServerSocket();
- srs.setReuseAddress(true);
- srs.bind(new InetSocketAddress("0.0.0.0", i));
- ServerSocket srs2 = srs;
- new Thread(() -> {
- try {
- while (true) {
- Socket sa = srs2.accept();
- System.out.println(sa.toString());
- //sa.close();
- // System.exit(0);
- }
- } catch (IOException e) {
- // e.printStackTrace();
- }
- }).start();
-
- /*Socket ste = new Socket();
- try {
- ste.connect(new InetSocketAddress(nat, port), 1);
- } catch (SocketTimeoutException e) {
-
- }
- ste.close();*/
- } catch (IOException e) {
- e.printStackTrace();
- } finally {
- if (s != null)
- s.close();
- /*if (srs != null)
- srs.close();*/
- }
- }
- }
-}