diff --git a/.classpath b/.classpath index 0039a3b..407aa40 100644 --- a/.classpath +++ b/.classpath @@ -2,8 +2,13 @@ + + + + + - + diff --git a/KLALB协议规范V2.0.docx b/KLALB协议规范V2.0.docx index 0886dc0..e36c7b1 100644 Binary files a/KLALB协议规范V2.0.docx and b/KLALB协议规范V2.0.docx differ diff --git a/kserver.ini b/kserver.ini index 944768a..cb1a88c 100644 --- a/kserver.ini +++ b/kserver.ini @@ -1,3 +1,3 @@ -virtualip=67c:72ce:765e:4db1:a02e:92fa:1959:29eb +virtualip=cefe:49b8:f837:4b00:b9cd:dbdd:c299:349 bind=0.0.0.0:4569 -local=127.0.0.1:36555 +local=127.0.0.1:5212 diff --git a/linetable - 副本.txt b/linetable - 副本.txt index 320126c..b1682b2 100644 --- a/linetable - 副本.txt +++ b/linetable - 副本.txt @@ -15,5 +15,4 @@ frp.freefrp.net:49965 frp1.freefrp.net:49965 frp2.freefrp.net:49965 frp4.freefrp.net:49965 -192.168.0.233:4569 -192.168.1.233:4569 \ No newline at end of file +cn-he-plc-2.openfrp.top:4569 \ No newline at end of file diff --git a/linetable.txt b/linetable.txt index cb2c257..d11d069 100644 --- a/linetable.txt +++ b/linetable.txt @@ -1,20 +1,2 @@ -43.248.189.107:65529 -cn-bj-bgp-3.openfrp.top:65529 -180.76.147.250:65529 -cn-ah-dx-1.natfrp.cloud:65529 -cn-nn-dx-1.natfrp.cloud:65529 -cn-wh-dx-1.natfrp.cloud:65529 -cn-zz-bgp-10.natfrp.cloud:23330 -cn-zz-bgp-7.natfrp.cloud:33336 -43.143.109.64:49965 -frp.104300.xyz:49965 -us.afrps.cn:49966 -hk.afrps.cn:49966 -la.afrps.cn:49966 -frp.freefrp.net:49965 -frp1.freefrp.net:49965 -frp2.freefrp.net:49965 -frp4.freefrp.net:49965 -192.168.0.233:4569 -192.168.1.233:4569 -cn-he-plc-2.openfrp.top:4569 \ No newline at end of file +{UDP}192.168.1.233:4569 +{UDP}192.168.0.233:4569 \ No newline at end of file diff --git a/src/org/kne/cloud/network/mport/DefaultServerSocketFactory.java b/src/org/kne/cloud/network/DefaultServerSocketFactory.java similarity index 78% rename from src/org/kne/cloud/network/mport/DefaultServerSocketFactory.java rename to src/org/kne/cloud/network/DefaultServerSocketFactory.java index 7c94200..83b29b5 100644 --- a/src/org/kne/cloud/network/mport/DefaultServerSocketFactory.java +++ b/src/org/kne/cloud/network/DefaultServerSocketFactory.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.net.InetAddress; @@ -7,6 +7,12 @@ import java.net.ServerSocket; import javax.net.ServerSocketFactory; public class DefaultServerSocketFactory extends ServerSocketFactory { + + + @Override + public ServerSocket createServerSocket() throws IOException { + return new ServerSocket(); + } @Override public ServerSocket createServerSocket(int port) throws IOException { diff --git a/src/org/kne/cloud/network/mport/DefaultSocketFactory.java b/src/org/kne/cloud/network/DefaultSocketFactory.java similarity index 88% rename from src/org/kne/cloud/network/mport/DefaultSocketFactory.java rename to src/org/kne/cloud/network/DefaultSocketFactory.java index e22f23a..9e2a4ac 100644 --- a/src/org/kne/cloud/network/mport/DefaultSocketFactory.java +++ b/src/org/kne/cloud/network/DefaultSocketFactory.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.net.InetAddress; @@ -7,7 +7,9 @@ import java.net.UnknownHostException; import javax.net.SocketFactory; public class DefaultSocketFactory extends SocketFactory { - public Socket createSocket() { + + @Override + public Socket createSocket() throws IOException { return new Socket(); } diff --git a/src/org/kne/cloud/network/mport/FilterServerSocket.java b/src/org/kne/cloud/network/FilterServerSocket.java similarity index 95% rename from src/org/kne/cloud/network/mport/FilterServerSocket.java rename to src/org/kne/cloud/network/FilterServerSocket.java index 68d0826..759b7a6 100644 --- a/src/org/kne/cloud/network/mport/FilterServerSocket.java +++ b/src/org/kne/cloud/network/FilterServerSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.net.InetAddress; diff --git a/src/org/kne/cloud/network/mport/FilterSocket.java b/src/org/kne/cloud/network/FilterSocket.java similarity index 95% rename from src/org/kne/cloud/network/mport/FilterSocket.java rename to src/org/kne/cloud/network/FilterSocket.java index 663f7bc..678ceb6 100644 --- a/src/org/kne/cloud/network/mport/FilterSocket.java +++ b/src/org/kne/cloud/network/FilterSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.FilterInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/HTTPDetectorItem.java b/src/org/kne/cloud/network/HTTPDetectorItem.java similarity index 91% rename from src/org/kne/cloud/network/mport/HTTPDetectorItem.java rename to src/org/kne/cloud/network/HTTPDetectorItem.java index 6a2da0b..ab7e3c2 100644 --- a/src/org/kne/cloud/network/mport/HTTPDetectorItem.java +++ b/src/org/kne/cloud/network/HTTPDetectorItem.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.DataInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/HTTPSDetectorItem.java b/src/org/kne/cloud/network/HTTPSDetectorItem.java similarity index 86% rename from src/org/kne/cloud/network/mport/HTTPSDetectorItem.java rename to src/org/kne/cloud/network/HTTPSDetectorItem.java index 051b6e8..77929a7 100644 --- a/src/org/kne/cloud/network/mport/HTTPSDetectorItem.java +++ b/src/org/kne/cloud/network/HTTPSDetectorItem.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.DataInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/HostPortMap.java b/src/org/kne/cloud/network/HostPortMap.java similarity index 87% rename from src/org/kne/cloud/network/mport/HostPortMap.java rename to src/org/kne/cloud/network/HostPortMap.java index 5bb3530..f18308f 100644 --- a/src/org/kne/cloud/network/mport/HostPortMap.java +++ b/src/org/kne/cloud/network/HostPortMap.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.util.HashMap; import java.util.LinkedHashMap; diff --git a/src/org/kne/cloud/network/mport/KLALBDetectorItem.java b/src/org/kne/cloud/network/KLALBDetectorItem.java similarity index 87% rename from src/org/kne/cloud/network/mport/KLALBDetectorItem.java rename to src/org/kne/cloud/network/KLALBDetectorItem.java index 69c4e88..ee1bf6f 100644 --- a/src/org/kne/cloud/network/mport/KLALBDetectorItem.java +++ b/src/org/kne/cloud/network/KLALBDetectorItem.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.DataInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/MultipurposeSocketAddress.java b/src/org/kne/cloud/network/MultipurposeSocketAddress.java similarity index 50% rename from src/org/kne/cloud/network/mport/MultipurposeSocketAddress.java rename to src/org/kne/cloud/network/MultipurposeSocketAddress.java index 9eb24b8..e3c3f17 100644 --- a/src/org/kne/cloud/network/mport/MultipurposeSocketAddress.java +++ b/src/org/kne/cloud/network/MultipurposeSocketAddress.java @@ -1,7 +1,8 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.io.Serializable; +import java.net.DatagramSocket; import java.net.Inet6Address; import java.net.InetAddress; import java.net.InetSocketAddress; @@ -19,20 +20,19 @@ public class MultipurposeSocketAddress implements Serializable{ /** * */ - private static Map socketFactoryRegister=new HashMap<>(); - private static Map serverSocketFactoryRegister=new HashMap<>(); + private static Map socketTypeRegister=new HashMap<>(); + + static { - socketFactoryRegister.put("TCP", new DefaultSocketFactory()); - serverSocketFactoryRegister.put("TCP", new DefaultServerSocketFactory()); + socketTypeRegister.put("TCP", new SocketType(new DefaultSocketFactory(), new DefaultServerSocketFactory())); + socketTypeRegister.put("UDP", new SocketType(new DefaultDatagramSocketFactory(),new DefaultDatagramServerSocketFactory())); } - public static Map getSocketFactoryRegister() { - return socketFactoryRegister; - } - public static Map getServerSocketFactoryRegister() { - return serverSocketFactoryRegister; - } + + public static Map getSocketTypeRegister() { + return socketTypeRegister; + } private static final long serialVersionUID = 1L; private String type; private String host; @@ -108,6 +108,9 @@ public class MultipurposeSocketAddress implements Serializable{ public int getPort() { return port; } + public String getType() { + return type; + } @Override public String toString() { StringBuilder sb=new StringBuilder(); @@ -125,35 +128,99 @@ public class MultipurposeSocketAddress implements Serializable{ return sb.toString(); } public Socket connectSocket(InetAddress bindip,int bindport,int timeout) throws UnknownHostException, IOException { - Socket s=socketFactoryRegister.get(type).createSocket(); + SocketFactory sf=socketTypeRegister.get(type).getSocketFactory(); + if(sf==null) { + throw new UnsupportedOperationException("Socket Unsupported"); + } + Socket s=sf.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(); + SocketFactory sf=socketTypeRegister.get(type).getSocketFactory(); + if(sf==null) { + throw new UnsupportedOperationException("Socket Unsupported"); + } + Socket s=sf.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(); + SocketFactory sf=socketTypeRegister.get(type).getSocketFactory(); + if(sf==null) { + throw new UnsupportedOperationException("Socket Unsupported"); + } + Socket s=sf.createSocket(); s.connect(new InetSocketAddress(host, port)); return s; } public Socket connectSocket(int timeout) throws UnknownHostException, IOException { - Socket s=socketFactoryRegister.get(type).createSocket(); + SocketFactory sf=socketTypeRegister.get(type).getSocketFactory(); + if(sf==null) { + throw new UnsupportedOperationException("Socket Unsupported"); + } + Socket s=sf.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)); + ServerSocketFactory srf=socketTypeRegister.get(type).getServerSocketFactory(); + if(srf==null) { + throw new UnsupportedOperationException("ServerSocket Unsupported"); + } + ServerSocket sk=srf.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)); + ServerSocketFactory srf=socketTypeRegister.get(type).getServerSocketFactory(); + if(srf==null) { + throw new UnsupportedOperationException("ServerSocket Unsupported"); + } + ServerSocket sk=srf.createServerSocket(port, backlog, InetAddress.getByName(host)); + return sk; + } + public boolean isStream() { + return checkIsStream(type); + } + public static boolean checkIsStream(String type2) { + return socketTypeRegister.get(type2).isStream(); + } + public DatagramSocket connectDatagramSocket(InetAddress bindip,int bindport) throws IOException { + DatagramSocketFactory dgs=socketTypeRegister.get(type).getDatagramSocketFactory(); + if(dgs==null) { + throw new UnsupportedOperationException("DatagramSocket Unsupported"); + } + DatagramSocket dgd=dgs.createSocket(); + dgd.bind(new InetSocketAddress(bindip, bindport)); + dgd.connect(new InetSocketAddress(host, port)); + return dgd; + } + public DatagramSocket connectDatagramSocket() throws IOException { + DatagramSocketFactory dgs=socketTypeRegister.get(type).getDatagramSocketFactory(); + if(dgs==null) { + throw new UnsupportedOperationException("DatagramSocket Unsupported"); + } + DatagramSocket dgd=dgs.createSocket(); + dgd.connect(new InetSocketAddress(host, port)); + return dgd; + } + public DatagramServerSocket listenDatagramServerSocket() throws UnknownHostException, IOException { + DatagramServerSocketFactory srf=socketTypeRegister.get(type).getDatagramServerSocketFactory(); + if(srf==null) { + throw new UnsupportedOperationException("DatagramServerSocket Unsupported"); + } + DatagramServerSocket sk=srf.createDatagramServerSocket(port, 50, InetAddress.getByName(host)); + return sk; + } + public DatagramServerSocket listenDatagramServerSocket(int backlog) throws UnknownHostException, IOException { + DatagramServerSocketFactory srf=socketTypeRegister.get(type).getDatagramServerSocketFactory(); + if(srf==null) { + throw new UnsupportedOperationException("DatagramServerSocket Unsupported"); + } + DatagramServerSocket sk=srf.createDatagramServerSocket(port, backlog, InetAddress.getByName(host)); return sk; } - } diff --git a/src/org/kne/cloud/network/mport/NetworkService.java b/src/org/kne/cloud/network/NetworkService.java similarity index 84% rename from src/org/kne/cloud/network/mport/NetworkService.java rename to src/org/kne/cloud/network/NetworkService.java index 4ffd5b1..9e4a6b5 100644 --- a/src/org/kne/cloud/network/mport/NetworkService.java +++ b/src/org/kne/cloud/network/NetworkService.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; public interface NetworkService { public void listen(MultipurposeSocketAddress msa); diff --git a/src/org/kne/cloud/network/mport/PortMultiUse.java b/src/org/kne/cloud/network/PortMultiUse.java similarity index 94% rename from src/org/kne/cloud/network/mport/PortMultiUse.java rename to src/org/kne/cloud/network/PortMultiUse.java index 56d292a..5f327da 100644 --- a/src/org/kne/cloud/network/mport/PortMultiUse.java +++ b/src/org/kne/cloud/network/PortMultiUse.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.File; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/PortRelay.java b/src/org/kne/cloud/network/PortRelay.java similarity index 78% rename from src/org/kne/cloud/network/mport/PortRelay.java rename to src/org/kne/cloud/network/PortRelay.java index 5d2df50..78123b0 100644 --- a/src/org/kne/cloud/network/mport/PortRelay.java +++ b/src/org/kne/cloud/network/PortRelay.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.DataInputStream; import java.io.DataOutputStream; @@ -12,12 +12,12 @@ import java.util.*; public class PortRelay { private int port; private HostPortMap services; - private TCPListener ssc; + private SocketListener ssc; public PortRelay(int port, HostPortMap services) throws IOException { this.services = services; this.port = port; - MultipurposeSocketAddress.getServerSocketFactoryRegister().put("ProtocolDetectorServerSocket", new ProtocolDetectorServerSocketFactory()); - ssc=new TCPListener(new MultipurposeSocketAddress("ProtocolDetectorServerSocket", "0.0.0.0", port)); + MultipurposeSocketAddress.getSocketTypeRegister().put("ProtocolDetectorServerSocket",new SocketType(null, new ProtocolDetectorServerSocketFactory()) ); + ssc=new SocketListener(new MultipurposeSocketAddress("ProtocolDetectorServerSocket", "0.0.0.0", port)); } public void start() throws IOException { ssc.setCon((s)->{ @@ -55,7 +55,6 @@ public class PortRelay { } }); - ssc.open(); System.out.println("已打开端口:" + port); } diff --git a/src/org/kne/cloud/network/mport/Protocol.java b/src/org/kne/cloud/network/Protocol.java similarity index 87% rename from src/org/kne/cloud/network/mport/Protocol.java rename to src/org/kne/cloud/network/Protocol.java index 09098c0..198c136 100644 --- a/src/org/kne/cloud/network/mport/Protocol.java +++ b/src/org/kne/cloud/network/Protocol.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.util.Objects; diff --git a/src/org/kne/cloud/network/mport/ProtocolDetector.java b/src/org/kne/cloud/network/ProtocolDetector.java similarity index 94% rename from src/org/kne/cloud/network/mport/ProtocolDetector.java rename to src/org/kne/cloud/network/ProtocolDetector.java index e369d66..d076551 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetector.java +++ b/src/org/kne/cloud/network/ProtocolDetector.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.io.InputStream; diff --git a/src/org/kne/cloud/network/mport/ProtocolDetectorItem.java b/src/org/kne/cloud/network/ProtocolDetectorItem.java similarity index 76% rename from src/org/kne/cloud/network/mport/ProtocolDetectorItem.java rename to src/org/kne/cloud/network/ProtocolDetectorItem.java index 5c3ca7b..d704f9b 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetectorItem.java +++ b/src/org/kne/cloud/network/ProtocolDetectorItem.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.InputStream; import java.util.function.Function; diff --git a/src/org/kne/cloud/network/mport/ProtocolDetectorServerSocket.java b/src/org/kne/cloud/network/ProtocolDetectorServerSocket.java similarity index 94% rename from src/org/kne/cloud/network/mport/ProtocolDetectorServerSocket.java rename to src/org/kne/cloud/network/ProtocolDetectorServerSocket.java index 659f91a..4e0b766 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetectorServerSocket.java +++ b/src/org/kne/cloud/network/ProtocolDetectorServerSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.io.InputStream; diff --git a/src/org/kne/cloud/network/mport/ProtocolDetectorServerSocketFactory.java b/src/org/kne/cloud/network/ProtocolDetectorServerSocketFactory.java similarity index 82% rename from src/org/kne/cloud/network/mport/ProtocolDetectorServerSocketFactory.java rename to src/org/kne/cloud/network/ProtocolDetectorServerSocketFactory.java index fe5b9e4..62571e9 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetectorServerSocketFactory.java +++ b/src/org/kne/cloud/network/ProtocolDetectorServerSocketFactory.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.net.InetAddress; @@ -17,9 +17,13 @@ public class ProtocolDetectorServerSocketFactory extends DefaultServerSocketFact this.pd = pd; } + @Override + public ServerSocket createServerSocket() throws IOException { + return new ProtocolDetectorServerSocket(super.createServerSocket()); + } + @Override public ServerSocket createServerSocket(int port) throws IOException { - // TODO 自动生成的方法存根 return new ProtocolDetectorServerSocket( super.createServerSocket(port),pd); } @@ -33,13 +37,11 @@ public class ProtocolDetectorServerSocketFactory extends DefaultServerSocketFact @Override public ServerSocket createServerSocket(int port, int backlog) throws IOException { - // TODO 自动生成的方法存根 return new ProtocolDetectorServerSocket(super.createServerSocket(port, backlog),pd); } @Override public ServerSocket createServerSocket(int port, int backlog, InetAddress ifAddress) throws IOException { - // TODO 自动生成的方法存根 return new ProtocolDetectorServerSocket(super.createServerSocket(port, backlog, ifAddress),pd); } diff --git a/src/org/kne/cloud/network/mport/ProtocolDetectorSocket.java b/src/org/kne/cloud/network/ProtocolDetectorSocket.java similarity index 90% rename from src/org/kne/cloud/network/mport/ProtocolDetectorSocket.java rename to src/org/kne/cloud/network/ProtocolDetectorSocket.java index 858e368..48ccf30 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetectorSocket.java +++ b/src/org/kne/cloud/network/ProtocolDetectorSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.BufferedInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/ProtocolStack.java b/src/org/kne/cloud/network/ProtocolStack.java similarity index 65% rename from src/org/kne/cloud/network/mport/ProtocolStack.java rename to src/org/kne/cloud/network/ProtocolStack.java index 13d793a..58b1935 100644 --- a/src/org/kne/cloud/network/mport/ProtocolStack.java +++ b/src/org/kne/cloud/network/ProtocolStack.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.util.Stack; diff --git a/src/org/kne/cloud/network/mport/ProxyProfileEntry.java b/src/org/kne/cloud/network/Proxy.java similarity index 65% rename from src/org/kne/cloud/network/mport/ProxyProfileEntry.java rename to src/org/kne/cloud/network/Proxy.java index 134871a..e0e09fb 100644 --- a/src/org/kne/cloud/network/mport/ProxyProfileEntry.java +++ b/src/org/kne/cloud/network/Proxy.java @@ -1,5 +1,7 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; +import java.io.Closeable; +import java.net.InetSocketAddress; import java.net.Socket; import java.util.HashMap; import java.util.Map; @@ -8,9 +10,13 @@ import java.util.function.Consumer; import javax.net.ServerSocketFactory; import javax.net.SocketFactory; -public class ProxyProfileEntry { +public abstract class Proxy implements Closeable{ private static Map register=new HashMap<>(); public static Map getRegister() { return register; } + + + + } diff --git a/src/org/kne/cloud/network/mport/ProxyProfileAnalyser.java b/src/org/kne/cloud/network/ProxyProfileAnalyser.java similarity index 100% rename from src/org/kne/cloud/network/mport/ProxyProfileAnalyser.java rename to src/org/kne/cloud/network/ProxyProfileAnalyser.java diff --git a/src/org/kne/cloud/network/mport/ProxyProfileExecutor.java b/src/org/kne/cloud/network/ProxyProfileExecutor.java similarity index 100% rename from src/org/kne/cloud/network/mport/ProxyProfileExecutor.java rename to src/org/kne/cloud/network/ProxyProfileExecutor.java diff --git a/src/org/kne/cloud/network/mport/RDPDetectorItem.java b/src/org/kne/cloud/network/RDPDetectorItem.java similarity index 86% rename from src/org/kne/cloud/network/mport/RDPDetectorItem.java rename to src/org/kne/cloud/network/RDPDetectorItem.java index 5d0d38c..dbaa497 100644 --- a/src/org/kne/cloud/network/mport/RDPDetectorItem.java +++ b/src/org/kne/cloud/network/RDPDetectorItem.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.DataInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/SSHDetectorItem.java b/src/org/kne/cloud/network/SSHDetectorItem.java similarity index 87% rename from src/org/kne/cloud/network/mport/SSHDetectorItem.java rename to src/org/kne/cloud/network/SSHDetectorItem.java index 280c071..ebe068c 100644 --- a/src/org/kne/cloud/network/mport/SSHDetectorItem.java +++ b/src/org/kne/cloud/network/SSHDetectorItem.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.DataInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/ServiceElement.java b/src/org/kne/cloud/network/ServiceElement.java similarity index 91% rename from src/org/kne/cloud/network/mport/ServiceElement.java rename to src/org/kne/cloud/network/ServiceElement.java index 3bb5cff..6fd86b7 100644 --- a/src/org/kne/cloud/network/mport/ServiceElement.java +++ b/src/org/kne/cloud/network/ServiceElement.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.net.UnknownHostException; import java.util.HashMap; diff --git a/src/org/kne/cloud/network/mport/ServiceToSocketProxyProfileEntry.java b/src/org/kne/cloud/network/ServiceToSocketProxy.java similarity index 53% rename from src/org/kne/cloud/network/mport/ServiceToSocketProxyProfileEntry.java rename to src/org/kne/cloud/network/ServiceToSocketProxy.java index acf96d5..6299a34 100644 --- a/src/org/kne/cloud/network/mport/ServiceToSocketProxyProfileEntry.java +++ b/src/org/kne/cloud/network/ServiceToSocketProxy.java @@ -1,18 +1,24 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; +import java.io.IOException; import java.util.function.Consumer; -public class ServiceToSocketProxyProfileEntry extends ProxyProfileEntry{ +public class ServiceToSocketProxy extends Proxy{ public NetworkService getSrc() { return src; } public MultipurposeSocketAddress getDes() { return des; } - public ServiceToSocketProxyProfileEntry(String l, String r) { + public ServiceToSocketProxy(String l, String r) { src=getRegister().get(l.substring(1, l.length()-1)); des=new MultipurposeSocketAddress(r); } private NetworkService src; private MultipurposeSocketAddress des; + @Override + public void close() throws IOException { + // TODO 自动生成的方法存根 + + } } diff --git a/src/org/kne/cloud/network/mport/SocketBridge.java b/src/org/kne/cloud/network/SocketBridge.java similarity index 55% rename from src/org/kne/cloud/network/mport/SocketBridge.java rename to src/org/kne/cloud/network/SocketBridge.java index e4cedc2..27c0469 100644 --- a/src/org/kne/cloud/network/mport/SocketBridge.java +++ b/src/org/kne/cloud/network/SocketBridge.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; @@ -6,33 +6,37 @@ import java.net.*; import org.kne.io.Task; public class SocketBridge extends Task{ - private Socket a; - private Socket b; - public StreamBridge getSab() { - return sab; + protected Socket a; + protected Socket b; + protected StreamBridge bridgeAB; + protected StreamBridge bridgeBA; + public StreamBridge getBridgeAB() { + return bridgeAB; } - public StreamBridge getSba() { - return sba; + public StreamBridge getBridgeBA() { + return bridgeBA; } - private StreamBridge sab; - private StreamBridge sba; public SocketBridge(Socket a, Socket b) throws IOException { super(); this.a = a; this.b = b; - sab = new StreamBridge(getAIN(), getBOUT()); - sba = new StreamBridge(getBIN(), getAOUT()); + createStreamBridge(); + } + protected void createStreamBridge() throws IOException { + bridgeAB = new StreamBridge(a.getInputStream(), b.getOutputStream()); + bridgeBA = new StreamBridge(b.getInputStream(), a.getOutputStream()); } @Override protected void runTask() { try { - sab.runAtNewThread("SocketBridge A->B thread"); - sba.runAtNewThread("SocketBridge B->A thread"); - sab.waitfortask(); - sba.waitfortask(); + bridgeAB.runAtNewThread("SocketBridge A->B thread"); + bridgeBA.runAtNewThread("SocketBridge B->A thread"); + bridgeAB.waitfortask(); + bridgeBA.waitfortask(); } catch (Exception e) { e.printStackTrace(); }finally { + //new Exception().printStackTrace(); try { a.close(); } catch (IOException e) { @@ -53,7 +57,7 @@ public class SocketBridge extends Task{ public Socket getB() { return b; } - protected OutputStream getBOUT() throws IOException { + /*protected OutputStream getBOUT() throws IOException { return b.getOutputStream(); } protected InputStream getBIN() throws IOException { @@ -64,6 +68,6 @@ public class SocketBridge extends Task{ } protected InputStream getAIN() throws IOException { return a.getInputStream(); - } + }*/ } diff --git a/src/org/kne/cloud/network/mport/TCPListener.java b/src/org/kne/cloud/network/SocketListener.java similarity index 52% rename from src/org/kne/cloud/network/mport/TCPListener.java rename to src/org/kne/cloud/network/SocketListener.java index 4a9da4d..b222819 100644 --- a/src/org/kne/cloud/network/mport/TCPListener.java +++ b/src/org/kne/cloud/network/SocketListener.java @@ -1,22 +1,25 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; +import java.io.Closeable; import java.io.IOException; import java.net.InetAddress; import java.net.ServerSocket; import java.net.Socket; +import java.net.SocketException; +import java.net.UnknownHostException; import java.util.function.Consumer; import javax.net.ServerSocketFactory; -public class TCPListener { - private ServerSocket serverSocket; +public class SocketListener implements Closeable,AutoCloseable{ + protected ServerSocket serverSocket; public ServerSocket getServerSocket() { return serverSocket; } - private volatile boolean flag=false; + private volatile boolean flag=true; - private Consumercon; + private volatile Consumercon; private Runnable r=new Runnable() { @Override public void run() { @@ -24,8 +27,26 @@ public class TCPListener { try { Socket soce=serverSocket.accept(); ThreadTool.makeVThreadIfSupport("端口监听线程",()->{ + if(con!=null) { + try { con.accept(soce); + }catch(Exception e) { + e.printStackTrace(); + try { + soce.close(); + } catch (IOException er) { + er.printStackTrace(); + } + } + }else { + try { + soce.close(); + } catch (IOException e) { + e.printStackTrace(); + } + } }).start(); + }catch(SocketException e) { } catch (IOException e) { e.printStackTrace(); } @@ -34,14 +55,17 @@ public class TCPListener { }; private MultipurposeSocketAddress multipurposeSocketAddress; - public TCPListener(MultipurposeSocketAddress multipurposeSocketAddress) { + public SocketListener(MultipurposeSocketAddress multipurposeSocketAddress) throws IOException { this.multipurposeSocketAddress=multipurposeSocketAddress; + open(); } - - public void open() throws IOException { - flag=true; + public SocketListener(ServerSocket tserverSocket) throws IOException { + this.serverSocket=tserverSocket; + open(); + } + protected void open() throws UnknownHostException, IOException { + if(serverSocket==null) serverSocket=multipurposeSocketAddress.listenServerSocket(); - //servers=new ServerSocket(port); new Thread(r).start(); } public void close() { diff --git a/src/org/kne/cloud/network/mport/SocketToServiceProxyProfileEntry.java b/src/org/kne/cloud/network/SocketToServiceProxy.java similarity index 55% rename from src/org/kne/cloud/network/mport/SocketToServiceProxyProfileEntry.java rename to src/org/kne/cloud/network/SocketToServiceProxy.java index a0cc947..dc407ac 100644 --- a/src/org/kne/cloud/network/mport/SocketToServiceProxyProfileEntry.java +++ b/src/org/kne/cloud/network/SocketToServiceProxy.java @@ -1,10 +1,11 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; +import java.io.IOException; import java.net.Socket; import java.util.function.Consumer; -public class SocketToServiceProxyProfileEntry extends ProxyProfileEntry { - public SocketToServiceProxyProfileEntry(String l, String r) { +public class SocketToServiceProxy extends Proxy { + public SocketToServiceProxy(String l, String r) { src=new MultipurposeSocketAddress(l); des=getRegister().get(r.substring(1, r.length()-1)) ; } @@ -16,4 +17,9 @@ public class SocketToServiceProxyProfileEntry extends ProxyProfileEntry { } private MultipurposeSocketAddress src; private NetworkService des; + @Override + public void close() throws IOException { + // TODO 自动生成的方法存根 + + } } diff --git a/src/org/kne/cloud/network/SocketToSocketProxy.java b/src/org/kne/cloud/network/SocketToSocketProxy.java new file mode 100644 index 0000000..836020b --- /dev/null +++ b/src/org/kne/cloud/network/SocketToSocketProxy.java @@ -0,0 +1,130 @@ +package org.kne.cloud.network; + +import java.io.IOException; +import java.net.InetAddress; +import java.net.InetSocketAddress; +import java.net.Socket; +import java.net.UnknownHostException; +import java.util.Map; + +public class SocketToSocketProxy extends Proxy { + private SocketListener sl; + private DatagramSocketListener dsl; + private MultipurposeSocketAddress listen, cbind, defaultConnect; + private Map detectedConnect; + + public SocketToSocketProxy(MultipurposeSocketAddress listen, MultipurposeSocketAddress connect) throws IOException { + this(listen, new MultipurposeSocketAddress(new InetSocketAddress(0)), connect); + } + + public SocketToSocketProxy(String l, String r) throws IOException { + this(new MultipurposeSocketAddress(l), new MultipurposeSocketAddress(r)); + } + + public SocketToSocketProxy(MultipurposeSocketAddress listen, MultipurposeSocketAddress cbind, + MultipurposeSocketAddress defaultConnect, Map detectedConnect) + throws IOException { + this.listen = listen; + this.cbind = cbind; + this.defaultConnect = defaultConnect; + this.detectedConnect = detectedConnect; + open(); + } + + public SocketToSocketProxy(MultipurposeSocketAddress listen, MultipurposeSocketAddress cbind, + MultipurposeSocketAddress defaultConnect) throws IOException { + this(listen, cbind, defaultConnect,null); + } + + private void open() throws IOException { + if (listen.isStream()) { + if (detectedConnect == null || detectedConnect.isEmpty()) { + sl = new SocketListener(listen); + } else { + sl = new SocketListener(new ProtocolDetectorServerSocket(listen.listenServerSocket())); + } + sl.setCon((sox) -> { + Socket sk = null; + try { + if (sox instanceof ProtocolDetectorSocket) { + ProtocolStack ps = ((ProtocolDetectorSocket) sox).getProtocolStack(); + if (!ps.isEmpty()) { + MultipurposeSocketAddress pmsa = detectedConnect.get(ps.pop().getName()); + if(pmsa!=null) { + sk = pmsa.connectSocket(InetAddress.getByName(cbind.getHost()), cbind.getPort()); + }else { + sk = defaultConnect.connectSocket(InetAddress.getByName(cbind.getHost()), cbind.getPort()); + + } + } else { + sk = defaultConnect.connectSocket(InetAddress.getByName(cbind.getHost()), cbind.getPort()); + } + } else { + sk = defaultConnect.connectSocket(InetAddress.getByName(cbind.getHost()), cbind.getPort()); + } + runBridge(sox,sk); + } catch (IOException e) { + e.printStackTrace(); + } finally { + if (sk != null) + try { + sk.close(); + } catch (IOException e) { + e.printStackTrace(); + } + try { + sox.close(); + } catch (IOException e) { + e.printStackTrace(); + } + } + }); + } else { + dsl = new DatagramSocketListener(listen); + dsl.setCon((dox) -> { + + }); + } + } + + protected void runBridge(Socket sk, Socket sox) throws IOException { + SocketBridge sb = new SocketBridge(sk, sox); + sb.run(); + } + + public MultipurposeSocketAddress getListen() { + return listen; + } + + public MultipurposeSocketAddress getCbind() { + return cbind; + } + + public MultipurposeSocketAddress getDefaultConnect() { + return defaultConnect; + } + + @Override + public void close() throws IOException { + if (dsl != null) + dsl.close(); + if (sl != null) + sl.close(); + } + + public InetAddress getListenAddress() { + if (sl != null) { + return sl.getServerSocket().getInetAddress(); + } else { + return dsl.getDatagramServerSocket().getLocalAddress(); + } + } + + public int getListenPort() { + if (sl != null) { + return sl.getServerSocket().getLocalPort(); + } else { + return dsl.getDatagramServerSocket().getLocalPort(); + } + } +} diff --git a/src/org/kne/cloud/network/mport/StreamBridge.java b/src/org/kne/cloud/network/StreamBridge.java similarity index 71% rename from src/org/kne/cloud/network/mport/StreamBridge.java rename to src/org/kne/cloud/network/StreamBridge.java index 73f1c7c..ea5d97d 100644 --- a/src/org/kne/cloud/network/mport/StreamBridge.java +++ b/src/org/kne/cloud/network/StreamBridge.java @@ -1,15 +1,16 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.*; import java.util.Objects; +import java.util.UUID; import org.kne.io.Task; public class StreamBridge extends Task{ - private InputStream in; - private OutputStream out; - private int blocksize=65535; - private long delay=0; + protected InputStream in; + protected OutputStream out; + protected int blocksize=65535; + protected long delay=0; public StreamBridge(InputStream in, OutputStream out) { super(); Objects.requireNonNull(in); @@ -20,12 +21,16 @@ public class StreamBridge extends Task{ @Override public void runTask() { + //FileOutputStream fos=null; try { + //fos=new FileOutputStream(UUID.randomUUID()+".txt"); int v=-1; byte[]b=new byte[blocksize]; while ((v=in.read(b))!=-1) { + // fos.write(b, 0, v); out.write(b,0,v); out.flush(); + if(delay>0) try { Thread.sleep(delay); } catch (InterruptedException e) { @@ -36,13 +41,21 @@ public class StreamBridge extends Task{ // TODO 自动生成的 catch 块 e.printStackTrace(); }finally { + /*if(fos!=null) + try { + fos.close(); + } catch (IOException e1) { + e1.printStackTrace(); + }*/ try { + if(in!=null) in.close(); } catch (IOException e) { // TODO 自动生成的 catch 块 e.printStackTrace(); } try { + if(in!=null) out.close(); } catch (IOException e) { // TODO 自动生成的 catch 块 diff --git a/src/org/kne/cloud/network/mport/ThreadTool.java b/src/org/kne/cloud/network/ThreadTool.java similarity index 72% rename from src/org/kne/cloud/network/mport/ThreadTool.java rename to src/org/kne/cloud/network/ThreadTool.java index 2982542..62238e1 100644 --- a/src/org/kne/cloud/network/mport/ThreadTool.java +++ b/src/org/kne/cloud/network/ThreadTool.java @@ -1,10 +1,10 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.lang.reflect.Method; public class ThreadTool { public static boolean first=true; - public static boolean forceD=true; + public static boolean forceD=false; public static Thread makeVThreadIfSupport(String name,Runnable r) { if(forceD) return new Thread(r, name); @@ -21,7 +21,7 @@ public class ThreadTool { }catch(Throwable e) { //e.printStackTrace(); if(first) { - System.out.println("请使用java19以上版本并开启--enable-preview选项以提高性能!"); + System.out.println("提示:请使用java21以上版本以提高本软件的数据转发性能!"); first=false; } return new Thread(r, name); @@ -46,7 +46,7 @@ public class ThreadTool { }catch(Throwable e) { //e.printStackTrace(); if(first) { - System.out.println("请使用java19以上版本以提高性能!"); + System.out.println("提示:请使用java21以上版本以提高本软件的数据转发性能!"); first=false; } Thread rt=new Thread(r, name); @@ -54,4 +54,14 @@ public class ThreadTool { return rt; } } + + public static Thread makePThreadIfSupport(String name, Runnable r) { + Thread rt=new Thread(r, name); + return rt; + } + public static Thread makePDaemonThreadIfSupport(String name, Runnable r) { + Thread rt=new Thread(r, name); + rt.setDaemon(true); + return rt; + } } diff --git a/src/org/kne/cloud/network/mport/VirtualServerSocket.java b/src/org/kne/cloud/network/VirtualServerSocket.java similarity index 82% rename from src/org/kne/cloud/network/mport/VirtualServerSocket.java rename to src/org/kne/cloud/network/VirtualServerSocket.java index 0fab1ce..80137ed 100644 --- a/src/org/kne/cloud/network/mport/VirtualServerSocket.java +++ b/src/org/kne/cloud/network/VirtualServerSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.lang.reflect.Constructor; @@ -13,8 +13,11 @@ import javax.sql.rowset.RowSetMetaDataImpl; public abstract class VirtualServerSocket extends ServerSocket { - public VirtualServerSocket(SocketImpl si) throws IOException { + private VirtualSocketImpl virtualImpl; + + public VirtualServerSocket(VirtualSocketImpl si) throws IOException { super(si); + this.virtualImpl=si; /*try { Class c=ServerSocket.class; Field con =c.getDeclaredField("impl"); @@ -45,6 +48,10 @@ public abstract class VirtualServerSocket extends ServerSocket { }*/ } + public VirtualSocketImpl getVirtualImpl() { + return virtualImpl; + } + @Override public abstract VirtualSocket accept() throws IOException ; diff --git a/src/org/kne/cloud/network/mport/VirtualSocket.java b/src/org/kne/cloud/network/VirtualSocket.java similarity index 68% rename from src/org/kne/cloud/network/mport/VirtualSocket.java rename to src/org/kne/cloud/network/VirtualSocket.java index 25ef25c..1665fdc 100644 --- a/src/org/kne/cloud/network/mport/VirtualSocket.java +++ b/src/org/kne/cloud/network/VirtualSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.FilterInputStream; import java.io.IOException; @@ -13,21 +13,25 @@ import java.nio.channels.SocketChannel; public abstract class VirtualSocket extends Socket { - private VirtualSocketImpl si; + private VirtualSocketImpl virtualImpl; + + public VirtualSocketImpl getVirtualImpl() { + return virtualImpl; + } public VirtualSocket(VirtualSocketImpl si) throws SocketException { super(si); - this.si=si; + this.virtualImpl=si; } @Override public InputStream getInputStream() throws IOException { - return si.getInputStream(); + return virtualImpl.getInputStream(); } @Override public OutputStream getOutputStream() throws IOException { - return si.getOutputStream(); + return virtualImpl.getOutputStream(); } diff --git a/src/org/kne/cloud/network/mport/VirtualSocketImpl.java b/src/org/kne/cloud/network/VirtualSocketImpl.java similarity index 88% rename from src/org/kne/cloud/network/mport/VirtualSocketImpl.java rename to src/org/kne/cloud/network/VirtualSocketImpl.java index 44a0097..9160906 100644 --- a/src/org/kne/cloud/network/mport/VirtualSocketImpl.java +++ b/src/org/kne/cloud/network/VirtualSocketImpl.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.io.InputStream; diff --git a/src/org/kne/cloud/network/klalb/ACKTPacket.java b/src/org/kne/cloud/network/klalb/ACKTPacket.java index b920bc8..d2d1ad3 100644 --- a/src/org/kne/cloud/network/klalb/ACKTPacket.java +++ b/src/org/kne/cloud/network/klalb/ACKTPacket.java @@ -1,19 +1,25 @@ package org.kne.cloud.network.klalb; import java.io.DataInput; +import java.io.DataInputStream; import java.io.DataOutput; +import java.io.DataOutputStream; import java.io.IOException; +import java.util.concurrent.atomic.AtomicInteger; -public class ACKTPacket extends KLALBPacket { +public class ACKTPacket extends KLALBPacket implements PortPacket { private int sport,dport; private long number; private boolean avaliable; - public ACKTPacket(int sport,int dport,long number,boolean avaliable) { + private int sendcount; + + public ACKTPacket(int sport,int dport,long number,boolean avaliable,int sendcount) { super(ACKT); this.sport=sport; this.dport=dport; this.number=number; this.avaliable=avaliable; + this.sendcount=sendcount; } public ACKTPacket() { @@ -31,7 +37,7 @@ public class ACKTPacket extends KLALBPacket { @Override public long getLength() { - return super.getLength()+17; + return super.getLength()+18; } public int getDport() { @@ -45,22 +51,68 @@ public class ACKTPacket extends KLALBPacket { public boolean isAvaliable() { return avaliable; } + + public int getSendcount() { + return sendcount; + } + + private final byte[] writeBuffer = new byte[18]; @Override protected void writeToStream(DataOutput dto) throws IOException { super.writeToStream(dto); - dto.writeInt(sport); + /*dto.writeInt(sport); dto.writeInt(dport); dto.writeLong(number); - dto.writeBoolean(avaliable); - } + dto.writeByte(sendcount); + dto.writeBoolean(avaliable);*/ + writeBuffer[0] = (byte)(sport >>> 24); + writeBuffer[1] = (byte)(sport >>> 16); + writeBuffer[2] = (byte)(sport >>> 8); + writeBuffer[3] = (byte)(sport >>> 0); + + writeBuffer[4] = (byte)(dport >>> 24); + writeBuffer[5] = (byte)(dport >>> 16); + writeBuffer[6] = (byte)(dport >>> 8); + writeBuffer[7] = (byte)(dport >>> 0); + writeBuffer[8] = (byte)(number >>> 56); + writeBuffer[9] = (byte)(number >>> 48); + writeBuffer[10] = (byte)(number >>> 40); + writeBuffer[11] = (byte)(number >>> 32); + writeBuffer[12] = (byte)(number >>> 24); + writeBuffer[13] = (byte)(number >>> 16); + writeBuffer[14] = (byte)(number >>> 8); + writeBuffer[15] = (byte)(number >>> 0); + + writeBuffer[16]=(byte) getSendRecord().size(); + + writeBuffer[17]=(byte) (avaliable ? 1 : 0); + dto.write(writeBuffer); + } + private final byte[] readBuffer = new byte[18]; @Override protected void readFromStream(DataInput din) throws IOException { super.readFromStream(din); - sport=din.readInt(); + /*sport=din.readInt(); dport=din.readInt(); number=din.readLong(); - avaliable=din.readBoolean(); + sendcount=din.readUnsignedByte(); + avaliable=din.readBoolean();*/ + + din.readFully(readBuffer); + sport=((readBuffer[0] << 24) + (readBuffer[1] << 16) + (readBuffer[2] << 8) + (readBuffer[3] << 0)); + dport=((readBuffer[4] << 24) + (readBuffer[5] << 16) + (readBuffer[6] << 8) + (readBuffer[7] << 0)); + number=(((long)readBuffer[8] << 56) + + ((long)(readBuffer[9] & 255) << 48) + + ((long)(readBuffer[10] & 255) << 40) + + ((long)(readBuffer[11] & 255) << 32) + + ((long)(readBuffer[12] & 255) << 24) + + ((readBuffer[13] & 255) << 16) + + ((readBuffer[14] & 255) << 8) + + ((readBuffer[15] & 255) << 0)); + sendcount=readBuffer[16]&0xff; + avaliable=(readBuffer[17] != 0); } -} + +} \ No newline at end of file diff --git a/src/org/kne/cloud/network/klalb/LINESPacket.java b/src/org/kne/cloud/network/klalb/ADDLINESPacket.java similarity index 81% rename from src/org/kne/cloud/network/klalb/LINESPacket.java rename to src/org/kne/cloud/network/klalb/ADDLINESPacket.java index 8e586fe..d941d12 100644 --- a/src/org/kne/cloud/network/klalb/LINESPacket.java +++ b/src/org/kne/cloud/network/klalb/ADDLINESPacket.java @@ -5,7 +5,7 @@ import java.io.DataOutput; import java.io.IOException; import java.net.Inet6Address; -public class LINESPacket extends KLALBPacket { +public class ADDLINESPacket extends KLALBPacket { private String lines; @@ -13,13 +13,13 @@ public class LINESPacket extends KLALBPacket { return lines; } - public LINESPacket(String lines) { - super(LINES); + public ADDLINESPacket(String lines) { + super(ADDLINES); this.lines=lines; } - public LINESPacket() { - super(LINES); + public ADDLINESPacket() { + super(ADDLINES); } private int UTFlength(String str) { int strlen = str.length(); @@ -38,7 +38,7 @@ public class LINESPacket extends KLALBPacket { } @Override public String toString() { - return "LINES\n"+lines; + return "ADDLINES\n"+lines; } @Override diff --git a/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java b/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java index c4490f6..0516448 100644 --- a/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java +++ b/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java @@ -11,14 +11,14 @@ public class ByteArrayRecycle { this.capacity = capacity; this.length = length; } - public synchronized void recycle(byte[]b) { + public void recycle(byte[]b) { if(b.length!=length) throw new IllegalArgumentException("wrong length"); if(rec.size(),PortPacket{ + public static final ByteArrayRecycle arrayRecycle=new ByteArrayRecycle(5000,65535); @Override public long getLength() { - return super.getLength()+18+size; + return super.getLength()+19+size; } private int sport,dport; private long number; private int size; private byte[]data; + private int sendcount; + + volatile long resendtimer=System.nanoTime(); + public DATATPacket(int sport,int dport,long number,byte[]data,int size) { super(DATAT); this.sport=sport; @@ -24,6 +34,10 @@ public class DATATPacket extends KLALBPacket { this.size=size; } + public int getSendcount() { + return sendcount; + } + public DATATPacket() { super(DATAT); } @@ -48,28 +62,107 @@ public class DATATPacket extends KLALBPacket { public String toString() { return "DATAT "+sport+"->"+dport+" "+number+"["+size+"]"; } - + private final byte[] writeBuffer = new byte[19]; @Override protected void writeToStream(DataOutput dto) throws IOException { super.writeToStream(dto); - dto.writeInt(sport); + /*dto.writeInt(sport); dto.writeInt(dport); dto.writeLong(number); - dto.writeChar(size); + dto.writeByte(getSendRecord().size()); + dto.writeChar(size);*/ + + writeBuffer[0] = (byte)(sport >>> 24); + writeBuffer[1] = (byte)(sport >>> 16); + writeBuffer[2] = (byte)(sport >>> 8); + writeBuffer[3] = (byte)(sport >>> 0); + + writeBuffer[4] = (byte)(dport >>> 24); + writeBuffer[5] = (byte)(dport >>> 16); + writeBuffer[6] = (byte)(dport >>> 8); + writeBuffer[7] = (byte)(dport >>> 0); + + writeBuffer[8] = (byte)(number >>> 56); + writeBuffer[9] = (byte)(number >>> 48); + writeBuffer[10] = (byte)(number >>> 40); + writeBuffer[11] = (byte)(number >>> 32); + writeBuffer[12] = (byte)(number >>> 24); + writeBuffer[13] = (byte)(number >>> 16); + writeBuffer[14] = (byte)(number >>> 8); + writeBuffer[15] = (byte)(number >>> 0); + + writeBuffer[16]=(byte) getSendRecord().size(); + + writeBuffer[17]=(byte) (size>>>8); + writeBuffer[18]=(byte) (size>>>0); + dto.write(writeBuffer); + //CRC32 crc=new CRC32(); + // crc.update(data, 0, size); dto.write(data,0,size); + //dto.writeLong(crc.getValue()); } + private final byte[] readBuffer = new byte[19]; @Override protected void readFromStream(DataInput din) throws IOException { super.readFromStream(din); - sport=din.readInt(); + /*sport=din.readInt(); dport=din.readInt(); number=din.readLong(); - size=din.readChar(); + sendcount=din.readUnsignedByte(); + size=din.readChar();*/ + + din.readFully(readBuffer); + sport=((readBuffer[0] << 24) + (readBuffer[1] << 16) + (readBuffer[2] << 8) + (readBuffer[3] << 0)); + dport=((readBuffer[4] << 24) + (readBuffer[5] << 16) + (readBuffer[6] << 8) + (readBuffer[7] << 0)); + number=(((long)readBuffer[8] << 56) + + ((long)(readBuffer[9] & 255) << 48) + + ((long)(readBuffer[10] & 255) << 40) + + ((long)(readBuffer[11] & 255) << 32) + + ((long)(readBuffer[12] & 255) << 24) + + ((readBuffer[13] & 255) << 16) + + ((readBuffer[14] & 255) << 8) + + ((readBuffer[15] & 255) << 0)); + sendcount=readBuffer[16]&0xff; + size=(((readBuffer[17]&0xff) << 8) + ((readBuffer[18]&0xff) << 0)); + data=arrayRecycle.create(); din.readFully(data,0,size); + /*CRC32 crc32=new CRC32(); + crc32.update(data, 0, size); + if(din.readLong()!=crc32.getValue()) { + throw new StreamCorruptedException("CRC32 error!"); + }*/ } + @Override + public int hashCode() { + return Objects.hash(number); + } + + @Override + public boolean equals(Object obj) { + if (this == obj) + return true; + if (obj == null) + return false; + if (getClass() != obj.getClass()) + return false; + DATATPacket other = (DATATPacket) obj; + return number == other.number; + } + public int getSize() { return size; } + + @Override + public int compareTo(DATATPacket o) { + if(number>o.number) { + return 1; + }else if(number selflineTableSupplier=()->{return null;}; - - public Supplier getSelflineTableSupplier() { return selflineTableSupplier; @@ -53,149 +75,156 @@ public class KLALBController { this.selflineTableSupplier = selflineTableSupplier; } - public LineManager getLineManager() { - if(lineManager==null) - lineManager=createLineManager(); - return lineManager; - } - - protected LineManager createLineManager() { - return new LineManager(this); - } - private Inet6Address self; public Inet6Address getSelf() { return self; } - - private Map> routes = new ConcurrentHashMap<>(); - - protected Map> getRoutes() { - return routes; + private List lines=new ArrayList<>(); + //private ReadWriteLock lineslock=new ReentrantReadWriteLock(); + + public List getLines() { + return lines; } - private void setLine(Inet6Address vaddr, KLALBRemoteSocket krs) { - synchronized (routes) { - List al = routes.computeIfAbsent(vaddr, (vaddr2) -> { - return new ArrayList(); - }); - al.add(krs); - } + private PortBinder streamPortBinder=new PortBinder(this); + + protected PortBinder getStreamPortBinder() { + return streamPortBinder; } - private void removeLine(KLALBRemoteSocket krs) { - synchronized (routes) { - Iterator>> iter = routes.entrySet().iterator(); - while (iter.hasNext()) { - Map.Entry> entry = (Map.Entry>) iter - .next(); - entry.getValue().remove(krs); - if (entry.getValue().isEmpty()) { - iter.remove(); - } + + public void reconnectImmediately() { + synchronized (lines) { + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + klalbRemoteLine.reconnectImmediately(); } } } + + private PacketReceiver prc=new PacketReceiver(); + private class PacketReceiver implements KLALBPacketConsumer{ - public void addRemoteSocket(KLALBRemoteSocket krs) { - krs.setController(this); - CountDownLatch cdl = new CountDownLatch(1); - krs.setPacketReceiver((rec) -> { + @Override + public void accept(KLALBRemoteLine krs, KLALBPacket rec) { try { - if (rec instanceof SYNTPacket) { - SYNTPacket synt = (SYNTPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(synt.getDport()); - if (kvi != null) { - if (kvi.isListening()) { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), synt); - } - } else { - - sendPacketToAddress(krs.getRemoteVaddr(), new RSTPacket(synt.getDport(), synt.getSport()), - 65537); + if(rec instanceof PortPacket&&krs.getRemoteVaddr()!=null) { + PortPacket pt=(PortPacket) rec; + if(!streamPortBinder.distributePacketToConsumer(krs, pt)) { + if(!(pt instanceof RSTPacket)) + sendPacketToAddress(krs.getRemoteVaddr(), new RSTPacket(pt.getDport(), pt.getSport()), + 0,2); } - } else if (rec instanceof SACKTPacket) { - SACKTPacket sackt = (SACKTPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(sackt.getDport()); - if (kvi != null) { - if (!kvi.isListening()) { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), sackt); - } - } - } else if (rec instanceof RSTPacket) { - RSTPacket rst = (RSTPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(rst.getDport()); - if (kvi != null) { - if (kvi.isListening()) { - KLALBVirtualSocketImpl kvi2 = kvi.getAccepts() - .get(new InetSocketAddress(krs.getRemoteVaddr(), rst.getSport())); - if (kvi2 != null) { - kvi2.getPackReceiver().accept(krs.getRemoteVaddr(), rst); - } - } else { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), rst); - } - } - } else if (rec instanceof DATATPacket) { - DATATPacket datat = (DATATPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(datat.getDport()); - if (kvi != null) { - if (kvi.isListening()) { - KLALBVirtualSocketImpl kvi2 = kvi.getAccepts() - .get(new InetSocketAddress(krs.getRemoteVaddr(), datat.getSport())); - if (kvi2 != null) { - kvi2.getPackReceiver().accept(krs.getRemoteVaddr(), datat); - } - } else { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), datat); - } - } - } else if (rec instanceof ACKTPacket) { - ACKTPacket ackt = (ACKTPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(ackt.getDport()); - if (kvi != null) { - if (kvi.isListening()) { - KLALBVirtualSocketImpl kvi2 = kvi.getAccepts() - .get(new InetSocketAddress(krs.getRemoteVaddr(), ackt.getSport())); - if (kvi2 != null) { - kvi2.getPackReceiver().accept(krs.getRemoteVaddr(), ackt); - } - } else { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), ackt); - } - } - } else if (rec instanceof VADDRPacket) { - VADDRPacket var = (VADDRPacket) rec; - setLine(var.getVaddr(), krs); - cdl.countDown(); - }else if(rec instanceof LINESPacket) { - LINESPacket lpt=(LINESPacket) rec; + }else { + switch (rec.getType()) { + case KLALBPacket.ADDLINES: + ADDLINESPacket lpt=(ADDLINESPacket) rec; String s=lpt.getLines(); Scanner scn=new Scanner(s); while(scn.hasNext()) { String sn=scn.nextLine(); - getLineManager().addHostPort(new MultipurposeSocketAddress(sn)); + MultipurposeSocketAddress msa= new MultipurposeSocketAddress(sn); + addRemoteLines(msa); + } + break; + default: + break; } - } catch (NoRouteToHostException e) { + } + } catch (IOException e) { + // TODO 自动生成的 catch 块 e.printStackTrace(); } - }); - 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) { - e.printStackTrace(); + } + + + } + public void addRemoteLines(MultipurposeSocketAddress target) throws SocketTimeoutException, SocketException { + synchronized (lines) { + + try { + Enumerationeu= NetworkInterface.getNetworkInterfaces(); + while (eu.hasMoreElements()) { + NetworkInterface networkInterface = (NetworkInterface) eu.nextElement(); + if(networkInterface.isUp()) { + //System.out.println(networkInterface+" "+networkInterface.isUp()); + Enumerationei= networkInterface.getInetAddresses(); + while (ei.hasMoreElements()) { + InetAddress inetAddress = (InetAddress) ei.nextElement(); + MultipurposeSocketAddress bind=new MultipurposeSocketAddress(inetAddress.getHostAddress(),0); + if(!checkContainsTargetAndBind(target,bind)) { + //System.out.println(target+" "+bind); + addRemoteLine( new KLALBRemoteLine(target,bind)); + } + } + } + } + } catch (SocketException e) { + if(!checkContainsTarget(target)) + addRemoteLine( new KLALBRemoteLine(target)); + throw e; + } + } } + public Inet6Address getRemoteVaddrBySocketAddress(MultipurposeSocketAddress target) throws SocketTimeoutException { + KLALBRemoteLine kr=null; + synchronized(lines) { + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + if(target.equals(klalbRemoteLine.getSocketAddress())) { + kr=klalbRemoteLine; + } + } + } + if(kr==null) { + kr=new KLALBRemoteLine(target); + addRemoteLine(kr); + kr.waitForRemoteVaddrAvaliable(20000); + }else { + kr.reconnectImmediately(); + kr.waitForRemoteVaddrAvaliable(20000); + } + return kr.getRemoteVaddr(); + } + private boolean checkContainsTargetAndBind(MultipurposeSocketAddress target,MultipurposeSocketAddress bind) { + boolean b=false; + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + if(bind.equals(klalbRemoteLine.getBindAddress())&&target.equals(klalbRemoteLine.getSocketAddress())) { + b=true; + break; + } + } + return b; + } + + private boolean checkContainsTarget(MultipurposeSocketAddress target) { + boolean b=false; + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + if(target.equals(klalbRemoteLine.getSocketAddress())) { + b=true; + break; + } + } + return b; + } + + public void addRemoteLine(KLALBRemoteLine krs) throws SocketTimeoutException { + krs.setPacketReceiver(prc); + krs.setLocalVaddrSupplier(()->{return self;}); + krs.startIO(); + String selflineTable=selflineTableSupplier.get(); + if(selflineTable!=null) + krs.sendPacket(new ADDLINESPacket(selflineTable), 0); + synchronized (lines) { + lines.add(krs); + } + } + public KLALBController(Inet6Address self) { this.self = self; @@ -205,135 +234,149 @@ public class KLALBController { this.self=KLALBUtils.uuidToIP(UUID.randomUUID()); } - private Map bindmap = new ConcurrentHashMap<>(); - - protected Map getBindmap() { - return bindmap; - } - protected KLALBVirtualSocketImpl createVirtualImpl() { return new KLALBVirtualSocketImpl(this); } - protected int bind(KLALBVirtualSocketImpl klalbVirtualSocketImpl, int port) throws BindException { - synchronized (bindmap) { - if (port == 0) { - port = allocPort(); - } - if (bindmap.putIfAbsent(port, klalbVirtualSocketImpl) != null) { - throw new BindException("port " + port + " is already bind!"); - } - return port; - } - } - - protected void unbind(KLALBVirtualSocketImpl klalbVirtualSocketImpl) { - synchronized (bindmap) { - Set> s = bindmap.entrySet(); - Iterator> it = s.iterator(); - while (it.hasNext()) { - Entry object = it.next(); - if (klalbVirtualSocketImpl.equals(object.getValue())) { - it.remove(); - return; - } - } - } - } - - protected int allocPort() throws BindException { - int i = 1; - while (bindmap.containsKey(i)) { - if (i == 65535) { - throw new BindException("can't alloc port"); - } - i++; - } - return i; - } + protected void sendPacketToAddress(Inet6Address addr, KLALBPacket syntPacket, int priority) - throws NoRouteToHostException { + throws IOException { sendPacketToAddress(addr, syntPacket, priority, 1); } - - protected void sendPacketToAddress(Inet6Address addr, KLALBPacket packet, int priority, int count) - throws NoRouteToHostException { - synchronized (routes) { - List l = routes.get(addr); - if (l == null || l.isEmpty()) { - throw new NoRouteToHostException("address unreachable: " + addr); + + private void updateLines2(Inet6Address addr) throws SocketTimeoutException { + List l=new ArrayList(); + synchronized (lines) { + for (int i = 0; i < lines.size(); i++) { + KLALBRemoteLine klalbRemoteLine = lines.get(i); + if(klalbRemoteLine.isClosed()) { + lines.remove(i); + i--; + continue; + } + if(addr.equals(klalbRemoteLine.getRemoteVaddr())&&klalbRemoteLine.getMonitor().getState()==Monitor.ONLINE) { + l.add(klalbRemoteLine); + } + } + } + if(l.isEmpty()) { + lines2.remove(addr); + }else { + lines2.put(addr, l); + } } - List l2 = (List) ((ArrayList) l).clone(); + private Map> lines2 = new ConcurrentHashMap<>(); + private volatile long itm=System.nanoTime(); + protected void sendPacketToAddress(Inet6Address addr, KLALBPacket packet, int priority, int count) + throws IOException { + /*if(packet instanceof RSTPacket) { + new Exception("-RST-").printStackTrace(); + }*/ + //TimeDebugger tdb=new TimeDebugger(); + //tdb.putTime("start"); + + List lines2x; + //loop:while(true) { + long cur=System.nanoTime(); + if(cur-itm>10000000L) { + itm=cur; + lines2.clear(); + } + lines2x=lines2.get(addr); + if (lines2x == null || lines2x.isEmpty()) { + updateLines2(addr); + lines2x=lines2.get(addr); + } + if (lines2x == null || lines2x.isEmpty()) { + throw new NoRouteToHostException("address unreachable: " + addr); + } + //tdb.putTime("selectLines"); + /* for (Iterator iterator = lines2x.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + if(klalbRemoteLine.statLengthBefore(priority)<=65536*10) { + break loop; + } + } + try { + //System.out.println("slp"); + Thread.sleep(1); + } catch (InterruptedException e) { + e.printStackTrace(); + } + }*/ + List l2 = (List) ((ArrayList) lines2x).clone(); l2.removeAll(packet.getSendRecord()); if(l2.isEmpty()) { - l2 = (List) ((ArrayList) l).clone(); + l2 = (List) ((ArrayList) lines2x).clone(); } + + //tdb.putTime("findAvaliable"); + int count0 = Math.min(count, l2.size()); - LineDecitionComparator ldc = new LineDecitionComparator(l2,packet, priority); - Collections.sort(l2,ldc); + l2.forEach((r)->{ + r.runPredict(packet,priority); + }); + //Collections.shuffle(l2); + Collections.sort(l2); + //tdb.putTime("makeDecision"); //System.out.println(l2); - for (Iterator iterator = l2.iterator(); iterator.hasNext();) { - KLALBRemoteSocket krst = (KLALBRemoteSocket) iterator.next(); + for (int i = 0; i < l2.size(); i++) { + KLALBRemoteLine krst =l2.get(i); krst.sendPacket(packet, priority); packet.getSendRecord().add(krst); count0--; - Thread.yield(); if (count0 <= 0) break; } - } - + //tdb.putTime("sendPacket"); + //tdb.print(); + /*if(packet instanceof RSTPacket) + new Exception().printStackTrace();*/ } 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); - }); + synchronized (lines) { + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + klalbRemoteLine.remoeFromSendQueue(klalbPacket); } } } - protected boolean checkIsBind(KLALBVirtualSocketImpl klalbVirtualSocketImpl) { - return bindmap.containsValue(klalbVirtualSocketImpl); - } private Timer t=new Timer("数据包发送计时器", true); public Timer getTimer() { 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{ + + public void registerToProxyTypeAs(String proxyname) { + MultipurposeSocketAddress.getSocketTypeRegister().put(proxyname+"_Stream",socketType); + //ProxyProfileEntry.getRegister().put(proxyname, this); + } - @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); - } + /*@Override + public void listen(MultipurposeSocketAddress msa) { + // TODO 自动生成的方法存根 } + @Override + public void unlisten(MultipurposeSocketAddress msc) { + // TODO 自动生成的方法存根 + + } + + @Override + public void connect(MultipurposeSocketAddress msa) { + // TODO 自动生成的方法存根 + + } + + @Override + public void unconnect(MultipurposeSocketAddress msc) { + // TODO 自动生成的方法存根 + + }*/ + } diff --git a/src/org/kne/cloud/network/klalb/KLALBInputStream.java b/src/org/kne/cloud/network/klalb/KLALBInputStream.java index af6489b..915a38e 100644 --- a/src/org/kne/cloud/network/klalb/KLALBInputStream.java +++ b/src/org/kne/cloud/network/klalb/KLALBInputStream.java @@ -1,5 +1,7 @@ package org.kne.cloud.network.klalb; import static org.kne.cloud.network.klalb.KLALBPacket.*; + +import java.io.DataInput; import java.io.DataInputStream; import java.io.EOFException; import java.io.IOException; @@ -21,49 +23,68 @@ public class KLALBInputStream extends DataInputStream { if(bv!=2) throw new StreamCorruptedException("remote version is V"+bv+"."+sv+",not V2.0"); } - public synchronized KLALBPacket readPacket() throws IOException { - int type=read(); + public KLALBPacket readPacket() throws IOException { + return readKLALBPacketFromStream(this); + } + public static KLALBPacket readKLALBPacketFromStream(DataInputStream in) throws IOException { +int type=in.read(); if(type==-1) { - throw new EOFException(); + return null; } KLALBPacket klp; switch(type) { case PING:klp=new PINGPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case PONG: klp=new PONGPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case SYNT: klp=new SYNTPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case SACKT: klp=new SACKTPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case RST: klp=new RSTPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case DATAT: klp=new DATATPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case ACKT: klp=new ACKTPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case VADDR: klp=new VADDRPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; - case LINES: - klp=new LINESPacket(); - klp.readFromStream(this); + case ADDLINES: + klp=new ADDLINESPacket(); + klp.readFromStream(in); return klp; + case NACKT: + klp=new NACKTPacket(); + klp.readFromStream(in); + return klp; + case TEST: + klp=new TESTPacket(); + klp.readFromStream(in); + return klp; + case VADDRACK: + klp=new VADDRACKPacket(); + klp.readFromStream(in); + return klp; + case VADDRREQ: + klp=new VADDRREQPacket(); + klp.readFromStream(in); + 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 index b9440ab..7624060 100644 --- a/src/org/kne/cloud/network/klalb/KLALBMain.java +++ b/src/org/kne/cloud/network/klalb/KLALBMain.java @@ -1,72 +1,58 @@ package org.kne.cloud.network.klalb; +import java.io.BufferedReader; import java.io.File; +import java.io.FileNotFoundException; +import java.io.FileReader; 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; +import org.kne.cloud.network.KLALBProxyConfigJsonExecuter; +import org.kne.cloud.network.SocketToServiceProxy; +import org.kne.cloud.network.SocketToSocketProxy; 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); - + System.out.println(CONST.klalb+" V"+CONST.klalbver); Scanner scn=new Scanner(System.in); + + File configJson=new File("klalbconfig.json"); + + KLALBProxyConfigJsonExecuter kpcje=new KLALBProxyConfigJsonExecuter(); + kpcje.loadConfigJson(configJson); while(true) { String s=scn.next(); String[]sc=s.split(" "); switch(sc[0]) { case "help": + System.out.print("help:查看命令使用说明"); System.out.print("state:查看线路状态"); - System.out.println("reload:重新加载线路配置"); + //System.out.println("reload:重新加载线路配置文件"); + System.out.println("reconnect:所有离线线路跳过重连等待时间立即尝试重连"); + System.out.println("stop:退出程序"); + 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()); + System.out.println("状态\t上传流量\t下载流量\t上传速度\t下载速度\t延迟\t抖动\t下一次重试"); + for (Iterator iterator = kpcje.getKlalbController().getLines().iterator(); iterator.hasNext();) { + KLALBRemoteLine hostPort = iterator.next(); + System.out.println(hostPort .toString2()); + //System.out.println(); } break; + case "stop": + System.exit(0); + break; + case "reconnect": + kpcje.getKlalbController().reconnectImmediately(); + 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/KLALBOutputStream.java b/src/org/kne/cloud/network/klalb/KLALBOutputStream.java index 3de1e53..457b344 100644 --- a/src/org/kne/cloud/network/klalb/KLALBOutputStream.java +++ b/src/org/kne/cloud/network/klalb/KLALBOutputStream.java @@ -12,11 +12,15 @@ public class KLALBOutputStream extends DataOutputStream { byte[]b=new byte[] {'K','L','A','L','B'}; write(b); writeInt(2); - writeInt(0); + writeInt(1); flush(); } - public synchronized void writePacket(KLALBPacket klb) throws IOException { - write(klb.getType()); - klb.writeToStream(this); + public void writePacket(KLALBPacket klb) throws IOException { + writeKLALBPacketToStream(this, klb); + } + public static void writeKLALBPacketToStream(DataOutputStream out,KLALBPacket klb) throws IOException { + out.write(klb.getType()); + klb.writeToStream(out); + out.flush(); } } diff --git a/src/org/kne/cloud/network/klalb/KLALBPacket.java b/src/org/kne/cloud/network/klalb/KLALBPacket.java index 8bc5ac0..5f65c2d 100644 --- a/src/org/kne/cloud/network/klalb/KLALBPacket.java +++ b/src/org/kne/cloud/network/klalb/KLALBPacket.java @@ -8,7 +8,7 @@ import java.io.ObjectInput; import java.io.ObjectOutput; import java.util.Vector; -public abstract class KLALBPacket{ +public abstract class KLALBPacket implements Sumable{ public static final int PING=0; public static final int PONG=1; public static final int SYNT=2; @@ -18,7 +18,11 @@ 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; + public static final int ADDLINES=9; + public static final int NACKT=10; + public static final int TEST=11; + public static final int VADDRACK=12; + public static final int VADDRREQ=13; private int type; public KLALBPacket(int type) { @@ -32,18 +36,32 @@ public abstract class KLALBPacket{ public int getType() { return type; } + + private long sndtime,rcvtime; protected void writeToStream(DataOutput dto) throws IOException { - + sndtime=System.nanoTime(); } protected void readFromStream(DataInput din) throws IOException { - + rcvtime=System.nanoTime(); + } + public long getSndtime() { + return sndtime; + } + public long getRcvtime() { + return rcvtime; } public long getLength() { return 1; } - private Vector sendRecord=new Vector<>(); + @Override + public long getValue() { + return getLength(); + } - public Vector getSendRecord() { + + private Vector sendRecord=new Vector<>(); + + public Vector getSendRecord() { return sendRecord; } diff --git a/src/org/kne/cloud/network/klalb/KLALBUtils.java b/src/org/kne/cloud/network/klalb/KLALBUtils.java index 3e4a5bb..f0e314f 100644 --- a/src/org/kne/cloud/network/klalb/KLALBUtils.java +++ b/src/org/kne/cloud/network/klalb/KLALBUtils.java @@ -5,6 +5,7 @@ import java.io.DataOutputStream; import java.io.IOException; import java.net.Inet6Address; import java.net.UnknownHostException; +import java.util.List; import java.util.UUID; public class KLALBUtils { @@ -54,4 +55,24 @@ public class KLALBUtils { } return fmt; } + public static int binarySearchDATATPacketNumber(Listlst,long number) { + int low =0; + int high=lst.size()-1; + int middle=0; + if(low>high||numberlst.get(high).getNumber()) { + return -1; + } + while(low<=high) { + middle=(low+high)/2; + long xn=lst.get(middle).getNumber(); + if(xn>number) { + high=middle-1; + }else if(xn backlogQueue; + //private CountDownLatch reseted=new CountDownLatch(1); + private BlockingQueue backlogQueue; + private Thread sendDequeLock; private Queue sendDeque = new ConcurrentLinkedQueue<>(); - private List sendlist=new Vector<>(); + private Map sendmap=new HashMap<>(); + private volatile Thread sendthread; + //private List sendlist=new ArrayList<>(); private TimerTask sendCheckTask=new SendCheckTask(); private class SendCheckTask extends TimerTask{ - + public void run() { - synchronized (sendlist) { - for (int i = 0; i < sendlist.size(); i++) { + synchronized (sendmap) { + Collection cdp=sendmap.values(); + for (Iterator iterator = cdp.iterator(); iterator.hasNext();) { + DATATPacket dtp = (DATATPacket) iterator.next(); try { - sendlist.get(i).check(i); + long x=System.nanoTime(); + long dt=x-dtp.resendtimer; + long limit= dtp.getSendRecord().size()*RTTAvg*6+10000000; + if(dt>limit) { + if(dtp.getSendRecord().size()>=20) { + throw new IOException("send error!"); + } + controller.sendPacketToAddress((Inet6Address) remoteaddr,dtp, 5-1); + System.out.println("第"+(dtp.getSendRecord().size()-1)+"次重传:"+dtp); + dtp.resendtimer=x; + + } + } catch (IOException e) { + e.printStackTrace(); + try { + close0(true); + } catch (IOException e1) { + e1.printStackTrace(); + } + break; + } + } + + } + /*for (int i = 0; i < sendlist.size(); i++) { + try { + DATATPacket dtp= sendlist.get(i); + long x=System.nanoTime(); + long dt=x-dtp.resendtimer; + long limit= dtp.getSendRecord().size()*(1000000000L*i+RTTAvg); + if(dt>limit) { + if(dtp.getSendRecord().size()>=30) { + throw new IOException("send error!"); + } + controller.sendPacketToAddress((Inet6Address) address,dtp, 5-1); + System.out.println("第"+(dtp.getSendRecord().size()-1)+"次重传:"+dtp); + dtp.resendtimer=x; + + } } catch (IOException e) { e.printStackTrace(); try { @@ -101,111 +173,56 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } break; } - } - } + }*/ + + //System.out.println(isClosed()); } } - private volatile boolean succeed, refused; + + private TimerTask flowControlTask=new FlowControlTask(); + private class FlowControlTask extends TimerTask{ + @Override + public void run() { + if(sendDeque.isEmpty()) { + try {if(getLocalPort()!=0&&getPort()!=0) + if(remoteaddr instanceof Inet6Address&&(!remoteaddr.isAnyLocalAddress())) + controller.sendPacketToAddress((Inet6Address) remoteaddr, new ACKTPacket(getLocalPort(),getPort(), -1, + true,0), 0); + } catch (IOException e) { + try { + close0(true); + } catch (IOException e1) { + e1.printStackTrace(); + } + e.printStackTrace(); + } + + } + } + + } + + protected boolean isListening() { return backlogQueue != null; } - private List inputchache = new ArrayList<>(); + //private int inputcross = 0; + private List inputchache = new ArrayList<>(); private long inputcount = 0; private volatile boolean avaliable = true; - private BiConsumer packReceiver = new BiConsumer() { + + + private long RTTAvg=1000000000L; + private volatile boolean ignoreBindCheck; + private boolean connected; - @Override - public void accept(Inet6Address from, KLALBPacket u) { - try { - if (u instanceof SYNTPacket) { - if (isListening()) { - if (backlogQueue.offer(new InetSocketAddress(from, ((SYNTPacket) u).getSport()))) { - controller.sendPacketToAddress(from, - new SACKTPacket(localport, ((SYNTPacket) u).getSport()), 65537,2); - - } else { - controller.sendPacketToAddress(from, new RSTPacket(localport, ((SYNTPacket) u).getSport()), - 65537,2); - } - } else { - controller.sendPacketToAddress(from, new RSTPacket(localport, ((SYNTPacket) u).getSport()), - 65537,2); - } - } else if (u instanceof SACKTPacket) { - if (connecting) { - connecting = false; - succeed = true; - cdl.countDown(); - } - } else if (u instanceof RSTPacket) { - if (connecting) { - connecting = false; - refused = true; - cdl.countDown(); - close(); - } else { - close(); - } - } else if (u instanceof DATATPacket) { - DATATPacket dtp = (DATATPacket) u; - synchronized (inputchache) { - if (dtp.getNumber() >= inputcount) { - inputchache.add(dtp); - while (true) { - DATATPacket kkb = null; - for (int i = 0; i < inputchache.size(); i++) { - DATATPacket klalbBlock = inputchache.get(i); - if (klalbBlock.getNumber() == inputcount) { - inputchache.remove(i); - i--; - kkb = klalbBlock; - break; - } - } - if (kkb == null) - break; - sendDeque.add(kkb); - synchronized (sendDeque) { - sendDeque.notifyAll(); - - } - inputcount++; - } - } - } - controller.sendPacketToAddress(from, new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp.getNumber(), - countInputBytes() < inputchachesize), 32768,1); - } else if (u instanceof ACKTPacket) { - ACKTPacket ackt = (ACKTPacket) u; - avaliable = ackt.isAvaliable(); - AtomicReferencekl=new AtomicReference<>(); - synchronized (sendlist) { - - sendlist.removeIf((tsk)->{ - boolean b=tsk.getPacket().getNumber()==ackt.getNumber(); - if(b) - kl.set(tsk.getPacket()); - return b; - }); - sendlist.notifyAll(); - } - if(kl.get()!=null) { - controller.removeFromSend(from,kl.get()); - DATATPacket.arrayRecycle.recycle(kl.get().getData()); - } - } - } catch (IOException e) { - e.printStackTrace(); - } - } - - - - }; + + + private int countInputBytes() { AtomicInteger i = new AtomicInteger(0); sendDeque.forEach((c) -> { @@ -215,23 +232,17 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } private int countOutputBytes() { AtomicInteger i = new AtomicInteger(0); - sendlist.forEach((c) -> { - i.addAndGet(c.getPacket().getSize()); + sendmap.values().forEach((c) -> { + i.addAndGet(c.getSize()); }); return i.get(); } public KLALBController getController() { return controller; } - - public BiConsumer getPackReceiver() { - return packReceiver; - } - public KLALBVirtualSocketImpl(KLALBController kc) { super(); this.controller = kc; - kc.getTimer().schedule(sendCheckTask, 50, 50); } @Override @@ -259,7 +270,7 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { case SocketOptions.SO_SNDBUF: return outputchachesize; case SocketOptions.SO_BINDADDR: - return bindaddr; + return localaddr; default: return null; } @@ -283,43 +294,58 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { connect(new InetSocketAddress(address, port), 10000); } - private CountDownLatch cdl = new CountDownLatch(1); - @Override protected void connect(SocketAddress address, int timeout) throws IOException { - if (!controller.checkIsBind(this)) { + if(!ignoreBindCheck) + if (!controller.getStreamPortBinder().checkIsBind(this)) { bind(Inet6Address.getByName("::0"), 0); } - connecting = true; port = ((InetSocketAddress) address).getPort(); - this.address = ((InetSocketAddress) address).getAddress(); - controller.sendPacketToAddress((Inet6Address) this.address, new SYNTPacket(localport, port), 65537); + this.address=this.remoteaddr = (Inet6Address) ((InetSocketAddress) address).getAddress(); + controller.getStreamPortBinder().connect(this); + + + if(connected) + throw new SocketException("already connected"); + connected=true; + + controller.getTimer().schedule(sendCheckTask, 50, 50); + //controller.sendPacketToAddress((Inet6Address) this.remoteaddr, new SYNTPacket(localport, port), 0,1); try { - if (timeout == 0) { - cdl.await(); - } else { - cdl.await(timeout, TimeUnit.MILLISECONDS); + getKVSIOutputStream().write(compress); + getKVSIOutputStream().forceFlush(); + getKVSIOutputStream().waitForAllAcknowledged(timeout); + if(compress==1) { + vout=new DeflaterOutputStream(vout, new Deflater(Deflater.BEST_COMPRESSION, true), LIMIT, true); + }else { + vout = getKVSIOutputStream(); } - } catch (InterruptedException e) { - e.printStackTrace(); + int comp=getKVSIInputStream().read(); + if(comp==1) { + vin=new InflaterInputStream(getKVSIInputStream(),new Inflater(true),LIMIT); + }else { + vin = getKVSIInputStream(); + } - connecting = false; - if (succeed) { - } else if (refused) { - throw new ConnectException("connect refused"); - } else { + }catch(SocketTimeoutException e) { throw new SocketTimeoutException("connect time out"); + }catch(SocketException e) {//e.printStackTrace(); + throw new ConnectException("connect refused"); } + + controller.getTimer().schedule(flowControlTask, 5000, 5000); } - private volatile boolean connecting = false; + //private volatile boolean connecting = false; - public boolean isConnecting() { + /*public boolean isConnecting() { return connecting; - } - + }*/ @Override protected void bind(InetAddress host, int port) throws IOException { + bind(host,port,false); + } + protected void bind(InetAddress host, int port,boolean ignoreBindCheck) throws IOException { if (host.equals(Inet4Address.getByName("0.0.0.0"))) { host = Inet6Address.getByName("::0"); } @@ -329,46 +355,58 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { if ((!host.isAnyLocalAddress()) && (!host.equals(controller.getSelf()))) { throw new BindException("must bind to self"); } - bindaddr = (Inet6Address) host; - this.localport = controller.bind(this, port); + localaddr = (Inet6Address) host; + localport=port; + this.ignoreBindCheck=ignoreBindCheck; + if(!ignoreBindCheck) + controller.getStreamPortBinder().bind(this); + } + + @Override + public int getPort() { + return super.getPort(); } @Override protected void listen(int backlog) throws IOException { backlogQueue = new ArrayBlockingQueue<>(backlog); - address=bindaddr; - } - - private Map accepts = new ConcurrentHashMap<>(); - - public Map getAccepts() { - return accepts; + controller.getStreamPortBinder().listen(this); + address=localaddr; } @Override protected void accept(SocketImpl s) throws IOException { - try { KLALBVirtualSocketImpl kvsi = (KLALBVirtualSocketImpl) s; - InetSocketAddress isa = backlogQueue.take(); + + Object[] p=null; + while(true) { + if (isClosed()) + throw new SocketException("Socket is closed"); + p=backlogQueue.peek(); + if(p!=null) { + break; + } + try { + Thread.sleep(1); + } catch (InterruptedException e) { + e.printStackTrace(); + } + } + InetSocketAddress isa = (InetSocketAddress) p[0]; kvsi.inputchachesize=inputchachesize; kvsi.outputchachesize=outputchachesize; - kvsi.port = isa.getPort(); - kvsi.address = isa.getAddress(); + + kvsi.address = isa.getAddress(); + kvsi.remoteaddr=(Inet6Address) isa.getAddress(); + + kvsi.bind(localaddr, localport,true); + kvsi.accept((KLALBRemoteLine)p[1],(KLALBPacket) p[2]); + kvsi.connect(isa, 5000); + backlogQueue.poll(); + /*kvsi.port = isa.getPort(); kvsi.localport = localport; - kvsi.bindaddr = bindaddr; - accepts.put(isa, kvsi); - kvsi.setCloseListener((x) -> { - accepts.remove(isa); - }); - } catch (InterruptedException e) { - e.printStackTrace(); - } - } - - private Consumer acceptedSocketCloseListener; - - private void setCloseListener(Consumer lsr) { - this.acceptedSocketCloseListener = lsr; + kvsi.localaddr = localaddr;*/ + } private InputStream vin; @@ -382,27 +420,7 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { @Override public int read() throws IOException { if (dtp == null ) { - count = 0; - while (true) { - if (isClosed()) - throw new SocketException("Socket is closed"); - DATATPacket dtp2 = sendDeque.poll(); - if (dtp2 != null) { - dtp = dtp2; - if(countInputBytes() < inputchachesize) { - controller.sendPacketToAddress((Inet6Address) address, new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp2.getNumber(), - true), 32768); - } - break; - } - try { - synchronized (sendDeque) { - sendDeque.wait(100); - } - } catch (InterruptedException e) { - e.printStackTrace(); - } - } + nextPacket(); } if (dtp.getSize() == 0) { return -1; @@ -416,6 +434,29 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } } + private DATATPacket nextPacket() throws IOException { + count = 0; + while (true) { + if (isClosed()) + throw new SocketException("Socket is closed"); + DATATPacket dtp2 = sendDeque.poll(); + if (dtp2 != null) { + dtp = dtp2; + checkFlowControl(dtp2); + break; + } + sendDequeLock=Thread.currentThread(); + LockSupport.parkNanos(1000000); + } + return dtp; + } + private void checkFlowControl(DATATPacket dtp2) throws IOException { + if(countInputBytes() >= inputchachesize-LIMIT*4) { + controller.sendPacketToAddress(remoteaddr, new ACKTPacket(dtp2.getDport(), dtp2.getSport(), dtp2.getNumber(), + true,dtp2.getSendcount()), 0); + } + } + @Override public int read(byte[] b, int off, int len) throws IOException { if (b == null) { @@ -430,29 +471,59 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { if (dtp == null ) { - count = 0; - while (true) { - if (isClosed()) - throw new SocketException("Socket is closed"); - DATATPacket dtp2 = sendDeque.poll(); - if (dtp2 != null) { - dtp = dtp2; - if(countInputBytes() < inputchachesize) { - controller.sendPacketToAddress((Inet6Address) address, new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp2.getNumber(), - true), 32768); - } - break; - } - try { - synchronized (sendDeque) { - sendDeque.wait(100); - - } - } catch (InterruptedException e) { - e.printStackTrace(); - } + nextPacket(); + } + if (dtp.getSize() == 0) { + return -1; + } else { + b[off]= dtp.getData()[count++] ; + if(count==dtp.getSize()) { + DATATPacket.arrayRecycle.recycle(dtp.getData()); + dtp=null; } } + int i = 1; + try { + while (i < len) { + + if (dtp == null ) { + nextPacket(); + } + if (dtp.getSize() == 0) { + break; + } + int remainD=len-i; + int min=Math.min(dtp.getSize()-count, remainD); + System.arraycopy(dtp.getData(), count, b, off + i, min); + count+=min; + i+=min; + //b[off + i]= dtp.getData()[count++] ; + if(count>=dtp.getSize()) { + DATATPacket.arrayRecycle.recycle(dtp.getData()); + dtp=null; + } + + } + } catch (IOException ee) { + } + return i; + } + /*@Override + public int read(byte[] b, int off, int len) throws IOException { + if (b == null) { + throw new NullPointerException(); + } else if (off < 0 || len < 0 || len > b.length - off) { + throw new IndexOutOfBoundsException(); + } else if (len == 0) { + return 0; + } + + len = Math.min(len, available()); + + + if (dtp == null ) { + nextPacket(); + } if (dtp.getSize() == 0) { return -1; } else { @@ -467,28 +538,7 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { for (; i < len; i++) { if (dtp == null ) { - count = 0; - while (true) { - if (isClosed()) - throw new SocketException("Socket is closed"); - DATATPacket dtp2 = sendDeque.poll(); - if (dtp2 != null) { - dtp = dtp2; - if(countInputBytes() < inputchachesize) { - controller.sendPacketToAddress((Inet6Address) address, new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp2.getNumber(), - true), 32768); - } - break; - } - try { - synchronized (sendDeque) { - sendDeque.wait(100); - - } - } catch (InterruptedException e) { - e.printStackTrace(); - } - } + nextPacket(); } if (dtp.getSize() == 0) { break; @@ -503,13 +553,18 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } catch (IOException ee) { } return i; - } - + }*/ @Override public void close() throws IOException { + try { + close0(); + }finally { + getKVSIOutputStream().close0(); + } + } + private void close0() throws IOException{ } - @Override public int available() throws IOException { AtomicInteger i = new AtomicInteger(0); @@ -524,138 +579,256 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } private volatile long outputcount = 0; + private int compress=0; - private static final int LIMIT=65535; + public int getCompress() { + return compress; + } + + public void setCompress(int compress) { + this.compress = compress; + } + + private static final int LIMIT=8669; private class KVSIOutputStream extends OutputStream { + private byte[] cache=DATATPacket.arrayRecycle.create(); private int count=0; - private Object lock=new Object(); + private Lock olock=new ReentrantLock(); @Override public void write(int b) throws IOException { if (isClosed()) throw new SocketException("Socket is closed"); - synchronized (lock) { - - cache[count++]=(byte) b; + olock.lock(); + try { + cache[count++]=(byte) b;//1429?5773?8669 if (count >= LIMIT) {//1429?5773?8669 flush0(); }else { flush(); } + }finally { + olock.unlock(); } } + + @Override public void write(byte[] b, int off, int len) throws IOException { if (isClosed()) throw new SocketException("Socket is closed"); - int ol=off+len; - synchronized (lock) { + + olock.lock(); + try { + while(len>0) { + int remainD=LIMIT-count; + int min=Math.min(len, remainD); + System.arraycopy(b, off, cache, count, min); + count+=min; + off+=min; + len-=min; + if (count >= LIMIT) {//1429?5773?8669 + flush0(); + } + } + /*int ol=off+len; for (int i = off; i < ol; i++) { cache[count++]=b[i]; if (count >= LIMIT) {//1429?5773?8669 flush0(); } - } + }*/ flush(); + }finally { + olock.unlock(); + } + } + public void waitForAllAcknowledged(int timeout) throws IOException { + long start=System.nanoTime(); + while(true){ + if (isClosed()) + throw new SocketException("Socket is closed"); + //System.out.println(sendmap.size()); + if(sendmap.isEmpty()) + break; + if(timeout!=0&&(System.nanoTime()-start>timeout*1000000)) + throw new SocketTimeoutException("wait for acknowledged timout"); + sendthread=Thread.currentThread(); + LockSupport.parkNanos(1000000L); } } - private volatile TimerTask tt; @Override public void flush() throws IOException { - if(nodelay) { - synchronized (lock) { - flush0(); - } - }else { - if(tt==null) { - tt=new TimerTask() { - - @Override - public void run() { - if(isClosed()) - cancel(); - try { - synchronized (lock) { - flush0(); + if(nodelay) { + olock.lock(); + try { + flush0(); + }finally { + olock.unlock(); + } + }else { + AtomicReference ioe=new AtomicReference<>(); + if(tt==null) { + tt=new TimerTask() { + + @Override + public void run() { + if(isClosed()) + cancel(); + try { + olock.lock(); + try { + flush0(); + }finally { + olock.unlock(); + } + } catch (IOException e) { + ioe.set(e); + } } - } catch (IOException e) { - e.printStackTrace(); + }; + new Timer("粘包计时线程").scheduleAtFixedRate(tt, delaytime, delaytime); + } + IOException ioex=ioe.get(); + if(ioex!=null) { + ioex.fillInStackTrace(); + throw ioex; } } - }; - new Timer("粘包计时线程").scheduleAtFixedRate(tt, 0, delaytime); - } + } + public void forceFlush() throws IOException{ + olock.lock(); + try { + flush0(); + }finally { + olock.unlock(); } } - private void flush0() throws IOException { if (count > 0) { + //TimeDebugger td=new TimeDebugger(); + //td.putTime("start"); + while (!avaliable) { try { - while (!avaliable) { Thread.sleep(1); - } } catch (InterruptedException e) { e.printStackTrace(); } - try { - while(countOutputBytes()>outputchachesize) { - synchronized (sendlist) { - sendlist.wait(10); } + //td.putTime("waitForAvaliable"); + while(true){ + if (isClosed()) + throw new SocketException("Socket is closed"); + //System.out.println(sendmap.size()); + boolean b=sendmap.size()<=outputchachesize/LIMIT; + if(b) + break; + sendthread=Thread.currentThread(); + LockSupport.parkNanos(1000000L); } - } catch (InterruptedException e) { - // TODO 自动生成的 catch 块 - e.printStackTrace(); - } + //td.putTime("waitForCache"); byte[] ba=cache; cache=DATATPacket.arrayRecycle.create(); - - SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, ba,count), 5,20); + //td.putTime("flushBuffer"); + DATATPacket pack=new DATATPacket(localport, port, outputcount++, ba,count); + //td.putTime("createPacket"); + controller.sendPacketToAddress(remoteaddr,pack,1{ + try { + DATATPacket dp; + while((dp=kis.nextPacket()).getSize()!=0) { + los.write(dp.getData(), 0, dp.getSize()); + los.flush(); + } + }catch(IOException e) { + e.printStackTrace(); + }finally { + try { + los.close(); + } catch (IOException e) { + e.printStackTrace(); + } + try { + kis.close(); + } catch (IOException e) { + e.printStackTrace(); + } + } + }); + Thread t2=ThreadTool.makeVThreadIfSupport("本地接收线程", ()->{ + try { + int len=-1; + while(true) { + if((len=lis.read(kos.cache,0,LIMIT))==-1) { + break; + } + kos.olock.lock(); + try { + kos.count=len; + kos.flush0(); + }finally { + kos.olock.unlock(); + + } + } + }catch(IOException e) { + e.printStackTrace(); + }finally { + try { + lis.close(); + } catch (IOException e) { + e.printStackTrace(); + } + try { + kos.close(); + } catch (IOException e) { + e.printStackTrace(); + } + } + }); + t1.start(); + t2.start(); + try { + t1.join(); + t2.join(); + } catch (InterruptedException e) { + e.printStackTrace(); + } + close(); + } + + @Override + public void accept(KLALBRemoteLine from, KLALBPacket u) { + try { + //System.out.println(this+" "+u); + switch (u.getType()) { + case KLALBPacket.RST: + //reseted.countDown(); + if(!isListening()) + close0(false); + break; + case KLALBPacket.DATAT: + DATATPacket dtp = (DATATPacket) u; + if(isListening()) { + if(dtp.getNumber()==0) { + synchronized (backlogQueue) { + + //controller.getStreamPortBinder().checkIsConnected(new Pair); + InetSocketAddress is=new InetSocketAddress(from.getRemoteVaddr(), dtp.getSport()); + AtomicBoolean ab=new AtomicBoolean(true); + for (Iterator iterator = backlogQueue.iterator(); iterator.hasNext();) { + Object[] objects = (Object[]) iterator.next(); + if(is.equals(objects[0])) { + ab.set(false); + break; + } + } + if(ab.get()) { + if(controller.getStreamPortBinder().checkIsConnect(this,new InetSocketAddress(from.getRemoteVaddr(), dtp.getSport()))) { + ab.set(false); + } + } + + + if(ab.get()) + if(backlogQueue.offer(new Object[] { is,from,dtp})) { + /* controller.sendPacketToAddress(from.getRemoteVaddr(), new ACKTPacket(dtp.getDport(), dtp.getSport(),dtp.getNumber(),true,0), + 0,2);*/ + }else { + controller.sendPacketToAddress(from.getRemoteVaddr(), new RSTPacket(dtp.getDport(), dtp.getSport()), + 0,2); + } + + + } + }else { + controller.sendPacketToAddress(from.getRemoteVaddr(), new RSTPacket(dtp.getDport(), dtp.getSport()), + 0,2); + } + }else { + controller.sendPacketToAddress(from.getRemoteVaddr(), new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp.getNumber(), + countInputBytes() < inputchachesize,dtp.getSendcount()), 0,2); + synchronized (inputchache) { + if (dtp.getNumber() >= inputcount) { + int currindex=(int) (dtp.getNumber()-inputcount); + int reqsize=1+currindex; + while(reqsize>inputchache.size()) { + inputchache.add(new AtomicInteger()); + } + for (int i = 0; i < currindex; i++) { + Object o=inputchache.get(i); + if(o instanceof AtomicInteger) { + ((AtomicInteger) o).incrementAndGet(); + if(((AtomicInteger) o).get()==50) { + controller.sendPacketToAddress(from.getRemoteVaddr(), new NACKTPacket(dtp.getDport(), dtp.getSport(), inputcount+i), 0,1); + System.out.println("请求快速重传:"+(inputcount+i)); + } + } + } + inputchache.set(currindex, dtp); + + Iteratoritr=inputchache.iterator(); + while (itr.hasNext()) { + Object datatPacket = itr.next(); + if(datatPacket instanceof DATATPacket) { + itr.remove(); + sendDeque.add((DATATPacket) datatPacket); + LockSupport.unpark(sendDequeLock); + inputcount++; + }else { + break; + } + + } + + /* int sze=inputchache.size(); + for (int i = 0; i < sze+1; i++) { + if(i < sze) { + DATATPacket klalbBlock = inputchache.get(i); + if(dtp.getNumber()itr=inputchache.iterator(); + while (itr.hasNext()) { + DATATPacket datatPacket = (DATATPacket) itr.next(); + if(datatPacket.getNumber()==inputcount) { + itr.remove(); + synchronized (sendDeque) { + sendDeque.add(datatPacket); + sendDeque.notifyAll(); + } + inputcount++; + inputcross=0; + }else { + inputcross++; + if(inputcross>5) { + inputcross=0; + controller.sendPacketToAddress(from.getRemoteVaddr(), new NACKTPacket(dtp.getDport(), dtp.getSport(), inputcount), 32768); + System.out.println("请求快速重传:"+inputcount); + } + break; + } + } + + */ + /* + if (dtp.getNumber() >= inputcount) { + inputchache.add(dtp); + while (true) { + DATATPacket kkb = null; + for (int i = 0; i < inputchache.size(); i++) { + DATATPacket klalbBlock = inputchache.get(i); + if (klalbBlock.getNumber() == inputcount) { + inputchache.remove(i); + i--; + kkb = klalbBlock; + break; + } + } + if (kkb == null) + break; + sendDeque.add(kkb); + synchronized (sendDeque) { + sendDeque.notifyAll(); + + } + inputcount++; + } + } + */ + } + } + //dbg.println(from.getMonitor()+","+dtp.getNumber()); + } + break; + case KLALBPacket.ACKT: + ACKTPacket ackt = (ACKTPacket) u; + if(isListening()) { + controller.sendPacketToAddress(from.getRemoteVaddr(), new RSTPacket(ackt.getDport(), ackt.getSport()), + 0,2); + }else { + avaliable = ackt.isAvaliable(); + DATATPacket kl=null; + synchronized (sendmap) { + + kl=sendmap.remove(ackt.getNumber()); + + } + /*synchronized (sendlist) { + int val=KLALBUtils.binarySearchDATATPacketNumber(sendlist, ackt.getNumber()); + System.out.println(val); + if(val!=-1) { + kl=sendlist.remove(val); + } + }*/ + /*synchronized (sendlist) { + int val=0; + for (Iterator iterator = sendlist.iterator(); iterator.hasNext();val++) { + DATATPacket inetSocketAddress = (DATATPacket) iterator.next(); + if(inetSocketAddress.getNumber()==ackt.getNumber()) { + iterator.remove(); + kl=inetSocketAddress; + break; + } + } + //System.out.println(val); + }*/ + if(kl!=null) { + if(sendthread!=null) + LockSupport.unpark(sendthread); + controller.removeFromSend(from.getRemoteVaddr(),kl); + DATATPacket.arrayRecycle.recycle(kl.getData()); + if(kl.getSendRecord().size()==1) { + long RTTC=ackt.getRcvtime()- kl.getSndtime(); + if(RTTC>RTTAvg) { + RTTAvg=(RTTAvg+RTTC)/2; + }else { + RTTAvg= (RTTAvg*99+RTTC)/100; + } + //System.out.println(RTTAvg); + } + } + } + break; + case KLALBPacket.NACKT: + NACKTPacket nackt=(NACKTPacket) u; + DATATPacket st=null; + synchronized (sendmap) { + st=sendmap.get(nackt.getNumber()); + } + /*for (Iterator iterator = sendlist.iterator(); iterator.hasNext();) { + DATATPacket st = (DATATPacket) iterator.next(); + if(st.getNumber()==nackt.getNumber()) { + + } + }*/ + if(st!=null) { + controller.sendPacketToAddress(remoteaddr,st, 5-1); + } + break; + default: + break; + } + } catch (IOException e) { + e.printStackTrace(); + } + } + + @Override + public InetAddress getRemoteInetAddress() { + return remoteaddr; + } + + @Override + public InetAddress getLocalInetAddress() { + return localaddr; + } + + @Override + public void setLocalPort(int i) { + localport=i; + } + } diff --git a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java index 2999198..b7f3dfb 100644 --- a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java +++ b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java @@ -1,5 +1,6 @@ package org.kne.cloud.network.klalb; +import java.net.SocketException; import java.util.Comparator; import java.util.HashMap; import java.util.Iterator; @@ -7,28 +8,23 @@ import java.util.List; import java.util.Map; import java.util.concurrent.atomic.AtomicLong; -public class LineDecitionComparator implements Comparator { - private Map predictedlatencys=new HashMap<>(); - public LineDecitionComparator (List krs,KLALBPacket curr,int priority) { - for (Iterator iterator = krs.iterator(); iterator.hasNext();) { - KLALBRemoteSocket klalbRemoteSocket = (KLALBRemoteSocket) iterator.next(); - double x=klalbRemoteSocket.getMonitor().getLatency()/2; - x+=klalbRemoteSocket.getMonitor().getJitter()/2; - AtomicLong al=new AtomicLong(0); - klalbRemoteSocket.getSendQueue().forEach((v)->{ - if(v.getPriority()>=priority) { - al.addAndGet(v.getPacket().getLength()); - } - }); - al.addAndGet(curr.getLength()); +public class LineDecitionComparator implements Comparator { + private Map predictedlatencys=new HashMap<>(); + public LineDecitionComparator (List krs,KLALBPacket curr,int priority) { + for (Iterator iterator = krs.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteSocket = (KLALBRemoteLine) iterator.next(); + long x=klalbRemoteSocket.getSndDelayFactor(); + //long x=klalbRemoteSocket.getMonitor().getLatency()>>1; + /*long al=0; + al=klalbRemoteSocket.statLengthBefore(priority)+curr.getLength(); long speed=klalbRemoteSocket.getMonitor(). getOutSpeed(); if(speed==0) { - if(al.get()>0) { + if(al>0) { x=Long.MAX_VALUE; } }else { - x+=al.get()*1000000000.0/speed; - } + x+=al*1000000000L/speed; + }*/ /*System.out.println(klalbRemoteSocket); System.out.println(klalbRemoteSocket.getSendQueue().size()); @@ -36,14 +32,14 @@ public class LineDecitionComparator implements Comparator { predictedlatencys.put(klalbRemoteSocket, x); } } - public Map getPredictedlatencys() { + public Map getPredictedlatencys() { return predictedlatencys; } @Override - public int compare(KLALBRemoteSocket o1, KLALBRemoteSocket o2) { - double t1=predictedlatencys.get(o1); - double t2=predictedlatencys.get(o2); + public int compare(KLALBRemoteLine o1, KLALBRemoteLine o2) { + long t1=predictedlatencys.get(o1); + long t2=predictedlatencys.get(o2); if(t1>t2) { return 1; }else if(t1>1; + this.latencyAvg=(latencyAvg*9+ newlatency)/10; //} + } + + if(latencyMin==Long.MIN_VALUE) { + latencyMin=newlatency; + }else { + /*if(newlatencyinSpeedMax) { + /*if(inSpeed>inSpeedMax) { inSpeedMax=inSpeed; }else { inSpeedMax=(inSpeedMax*9999+inSpeed)/10000; @@ -128,11 +187,22 @@ public class Monitor { outSpeedMax=outSpeed; }else { outSpeedMax=(outSpeedMax*9999+outSpeed)/10000; - } + }*/ + + this.inSpeedAvg=(inSpeedAvg*9+ inSpeed)/10; + this.outSpeedAvg=(outSpeedAvg*9+ outSpeed)/10; if(changeListener!=null) { changeListener.accept(this); } + + } + public void updateOutSpeedMax() { + if(outSpeed>outSpeedMax) { + outSpeedMax=outSpeed; + }else { + outSpeedMax=(outSpeedMax*9999+outSpeed)/10000; + } } private ConsumerchangeListener; @@ -142,22 +212,26 @@ public class Monitor { public void setChangeListener(Consumer changeListener) { this.changeListener = changeListener; } - public long getLatency() { - return latency; + public long getLatencyAvg() { + return latencyAvg; } public String toString() { StringBuilder sb=new StringBuilder(); + sb.append(name); + sb.append('\t'); sb.append(getDsc()); sb.append('\t'); - 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"); + 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(latencyAvg/1000000).append("ms"); return sb.toString(); } public String toString2() { StringBuilder sb=new StringBuilder(); + sb.append(name); + sb.append('\n'); sb.append(getDsc()); sb.append('\t'); - 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\t"); + 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(latencyAvg/1000000).append("ms").append("\t").append(jitter/1000000).append("ms\t"); if(state==OFFLINE) { sb.append((coolingTime-(System.currentTimeMillis()-mls))/1000L); diff --git a/src/org/kne/cloud/network/klalb/MonitoredSocket.java b/src/org/kne/cloud/network/klalb/MonitoredSocket.java index 8cf78f4..92e8974 100644 --- a/src/org/kne/cloud/network/klalb/MonitoredSocket.java +++ b/src/org/kne/cloud/network/klalb/MonitoredSocket.java @@ -8,7 +8,7 @@ import java.util.Timer; import java.util.TimerTask; import java.util.concurrent.atomic.AtomicLong; -import org.kne.cloud.network.mport.FilterSocket; +import org.kne.cloud.network.FilterSocket; public class MonitoredSocket extends FilterSocket { public Monitor getMonitor() { diff --git a/src/org/kne/cloud/network/klalb/PONGPacket.java b/src/org/kne/cloud/network/klalb/PONGPacket.java index aa45750..be88e6e 100644 --- a/src/org/kne/cloud/network/klalb/PONGPacket.java +++ b/src/org/kne/cloud/network/klalb/PONGPacket.java @@ -5,39 +5,53 @@ import java.io.DataOutput; import java.io.IOException; public class PONGPacket extends KLALBPacket { - public long getTime() { - return time; + + public long getTimepingsnd() { + return timepingsnd; } - private long time; + + public long getTimepingrcv() { + return timepingrcv; + } + + public long getTimepongsnd() { + return timepongsnd; + } + private long timepingsnd,timepingrcv,timepongsnd; @Override protected void writeToStream(DataOutput dto) throws IOException { super.writeToStream(dto); - dto.writeLong(time); + dto.writeLong(timepingsnd); + dto.writeLong(timepingrcv); + dto.writeLong(timepongsnd); } @Override protected void readFromStream(DataInput din) throws IOException { super.readFromStream(din); - time=din.readLong(); + timepingsnd=din.readLong(); + timepingrcv=din.readLong(); + timepongsnd=din.readLong(); } @Override public long getLength() { - return super.getLength()+8; + return super.getLength()+24; } public PONGPacket() { super(PONG); } - public PONGPacket(long time) { + + public PONGPacket( long timepingsnd, long timepingrcv, long timepongsnd) { super(PONG); - this.time=time; + this.timepingsnd = timepingsnd; + this.timepingrcv = timepingrcv; + this.timepongsnd = timepongsnd; } + @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 6a0ba15..2e42556 100644 --- a/src/org/kne/cloud/network/klalb/RSTPacket.java +++ b/src/org/kne/cloud/network/klalb/RSTPacket.java @@ -4,7 +4,7 @@ import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; -public class RSTPacket extends KLALBPacket { +public class RSTPacket extends KLALBPacket implements PortPacket{ private int sport,dport; public RSTPacket(int sport,int dport) { super(RST); diff --git a/src/org/kne/cloud/network/klalb/SendTask.java b/src/org/kne/cloud/network/klalb/SendTask.java index e738a96..ea54dc6 100644 --- a/src/org/kne/cloud/network/klalb/SendTask.java +++ b/src/org/kne/cloud/network/klalb/SendTask.java @@ -58,8 +58,7 @@ public class SendTask { public void check(int number) throws IOException { - long limit= (1<<(2*Math.min(4,count.get()-1)))*(number<3?200:1000*number)*1000000L; - //long limit=(1<{ - KLALBVirtualSocket kvs=null; - try { - s.setTcpNoDelay(true); - kvs=new KLALBVirtualSocket(kc, kr.getRemoteVaddr(), 23333); - kvs.setTcpNoDelay(true); - SocketBridge dsb=new SocketBridge(s, kvs); - dsb.getSab().setDelay(1); - dsb.run(); - /*CompressedSocketBridge csb=new CompressedSocketBridge(s, kvs); - csb.getSab().setDelay(1); - csb.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(); + new SocketToSocketProxy(new MultipurposeSocketAddress(ap.getProperty("local")), vmsa); System.out.println("提示:输入state并回车可以查看当前线路状态"); while(true) { String s=scn.nextLine(); switch(s) { case "state": System.out.println("状态\t上传流量\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()); + for (Iterator iterator = kc.getLines().iterator(); iterator.hasNext();) { + KLALBRemoteLine hostPort = iterator.next(); + System.out.println(hostPort .toString2()); //System.out.println(); } break; diff --git a/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java index 8027082..df6a136 100644 --- a/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java +++ b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java @@ -4,34 +4,33 @@ 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; +import org.kne.cloud.network.DatagramSocketListener; +import org.kne.cloud.network.MultipurposeSocketAddress; +import org.kne.cloud.network.ProtocolDetectorServerSocketFactory; +import org.kne.cloud.network.ProtocolDetectorSocket; +import org.kne.cloud.network.SocketBridge; +import org.kne.cloud.network.SocketListener; +import org.kne.cloud.network.SocketToSocketProxy; +import org.kne.cloud.network.SocketType; public class SimpleKLALBServer { - public static TCPListener tcpl,tcpl2; + public static SocketListener tcpl; + public static DatagramSocketListener udpl; public static ServerPropties sp; public static KLALBController kc; static { - MultipurposeSocketAddress.getServerSocketFactoryRegister().put("DETTCP", new ProtocolDetectorServerSocketFactory()); + MultipurposeSocketAddress.getSocketTypeRegister().put("DETTCP",new SocketType(null, new ProtocolDetectorServerSocketFactory())); } public static void main(String[] args) throws IOException { + //Debuger dbg=new Debuger(); + //dbg.start(); - System.out.println("KLALB负载均衡V2.0"); + System.out.println(CONST.klalb+" V"+CONST.klalbver); sp=new ServerPropties(); System.out.println("虚拟地址:"+sp.getVirtualIP().getHostAddress()); @@ -58,10 +57,11 @@ public class SimpleKLALBServer { 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()); + System.out.println("状态\t上传流量\t下载流量\t上传速度\t下载速度\t延迟\t抖动\t下一次重试"); + for (Iterator iterator = kc.getLines().iterator(); iterator.hasNext();) { + KLALBRemoteLine hostPort = iterator.next(); + System.out.println(hostPort .toString2()); + //System.out.println(); } break; default: @@ -72,13 +72,21 @@ public class SimpleKLALBServer { private static void openPort(String bip) throws IOException { if(tcpl!=null) tcpl.close(); + if(udpl!=null) { + udpl.close(); + } MultipurposeSocketAddress mpsa=new MultipurposeSocketAddress(bip,"DETTCP"); - tcpl=new TCPListener(mpsa); + tcpl=new SocketListener(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); + KLALBRemoteLine krs=null; + try { + krs = new KLALBRemoteLine(new StreamKLALBPacketLink(pds)); + kc.addRemoteLine(krs); + } catch (IOException e) { + e.printStackTrace(); + } }else { MultipurposeSocketAddress mpsa2=new MultipurposeSocketAddress(sp.getLocal()); Socket s=null; @@ -108,49 +116,21 @@ public class SimpleKLALBServer { } } }); - 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; + udpl=new DatagramSocketListener(new MultipurposeSocketAddress(bip, "UDP")); + udpl.setCon((r)->{ + KLALBRemoteLine krl; try { - s=mpsa.connectSocket(); - s.setTcpNoDelay(true); - soc.setTcpNoDelay(true); - SocketBridge dsb=new SocketBridge(s, soc); - dsb.getSab().setDelay(1); - dsb.run(); - /*CompressedSocketBridge csb=new CompressedSocketBridge(s, soc); - csb.getSab().setDelay(1); - csb.run();*/ - } catch (UnknownHostException e) { - // TODO 自动生成的 catch 块 - e.printStackTrace(); + krl=new KLALBRemoteLine(new DatagramKLALBPacketLink(r)); + kc.addRemoteLine(krl); } 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 void openLocalPort(String bip) throws IOException { + MultipurposeSocketAddress mpsa=new MultipurposeSocketAddress(bip); + new SocketToSocketProxy(new MultipurposeSocketAddress("KLALB_Stream", "::0", 23333),mpsa); } private static String fileRead(String filePath){ //1.定义一个BufferedReader对象,将文件内容读取到缓存 diff --git a/src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java b/src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java deleted file mode 100644 index 2da4c94..0000000 --- a/src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java +++ /dev/null @@ -1,16 +0,0 @@ -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/nathole/NatholeTestS.java b/src/org/kne/cloud/network/nathole/NatholeTestS.java index e08c3bf..488b4c2 100644 --- a/src/org/kne/cloud/network/nathole/NatholeTestS.java +++ b/src/org/kne/cloud/network/nathole/NatholeTestS.java @@ -9,12 +9,12 @@ import java.nio.channels.SocketChannel; public class NatholeTestS { public static void main(String[] args) throws IOException { - for (int i1 = 1024; i1 < 65536; i1++) { + for (int i1 = 1024; i1 < 32768; i1++) { try { SocketChannel sc= SocketChannel.open(); sc.configureBlocking(false) ; sc.bind(new InetSocketAddress("0.0.0.0", i1)); - sc.connect(new InetSocketAddress("10.235.16.1",10300)); + sc.connect(new InetSocketAddress("183.199.49.116",21000)); sc.close(); int i2=i1; new Thread(()->{ diff --git a/src/org/kne/debug/TimeDebugger.java b/src/org/kne/debug/TimeDebugger.java index c9f3a79..c9631d8 100644 --- a/src/org/kne/debug/TimeDebugger.java +++ b/src/org/kne/debug/TimeDebugger.java @@ -3,15 +3,16 @@ package org.kne.debug; import java.util.ArrayList; import java.util.HashMap; import java.util.Hashtable; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; public class TimeDebugger { - Map l=new Hashtable(); + Map l=new LinkedHashMap(); long tmp=-1; public void putTime(String name){ - long v=System.currentTimeMillis(); + long v=System.nanoTime(); if(tmp==-1){ l.put(name, 0l); }else{ @@ -20,6 +21,6 @@ public class TimeDebugger { tmp=v; } public void print(){ - System.err.println(l+"(ms)"); + System.err.println(l+"(ns)"); } }