diff --git a/.classpath b/.classpath
index af6f9dc..7fba2fa 100644
--- a/.classpath
+++ b/.classpath
@@ -2,5 +2,6 @@
+
diff --git a/src/org/kne/cloud/network/klalb/Consts.java b/src/org/kne/cloud/network/klalb/Consts.java
index 01daf2c..a5e57c7 100644
--- a/src/org/kne/cloud/network/klalb/Consts.java
+++ b/src/org/kne/cloud/network/klalb/Consts.java
@@ -1,6 +1,7 @@
package org.kne.cloud.network.klalb;
public class Consts {
- public static final int BLOCKSIZE=32768;
- public static final long PINGTIMENS=1000000000;
+ public static final int BLOCKSIZE=65536;
+ public static final long PINGTIMENS=10000000000L;
+ public static final double A = 0.125;
}
diff --git a/src/org/kne/cloud/network/klalb/IOThreadManager.java b/src/org/kne/cloud/network/klalb/IOThreadManager.java
index 61b6866..0d2d454 100644
--- a/src/org/kne/cloud/network/klalb/IOThreadManager.java
+++ b/src/org/kne/cloud/network/klalb/IOThreadManager.java
@@ -7,8 +7,10 @@ import java.net.SocketException;
import java.util.List;
import java.util.Vector;
+import org.kne.cloud.network.mport.ThreadTool;
+
public class IOThreadManager {
- private KLALBCore klc=new KLALBCore(200);
+ private KLALBCore klc=new KLALBCore(10000);
private Listtcps=new Vector<>();
@@ -48,7 +50,6 @@ public class IOThreadManager {
}catch(Exception e) {
e.printStackTrace();
}finally {
- closeRemote();
}
});
upo.start();
@@ -60,11 +61,16 @@ public class IOThreadManager {
Thread up=ThreadTool.makeVThreadIfSupport("远程发送线程",()->{
try{
while(true) {
+ if(s.getTunnel().getName().contains("Openfrp")) {
+ klc.sendDataBlockControlOnly(s);
+ }else {
klc.sendDataBlock(s);
+ }
}
}catch(InterruptedException s1) {
}catch(Exception e) {
+ e.printStackTrace();
}
});
Thread down=ThreadTool.makeVThreadIfSupport("远程接收线程",()->{
@@ -76,6 +82,7 @@ public class IOThreadManager {
}catch(Exception e) {
up.interrupt();
+ e.printStackTrace();
}
});
diff --git a/src/org/kne/cloud/network/klalb/KLALBBlock.java b/src/org/kne/cloud/network/klalb/KLALBBlock.java
index 98749b8..2124865 100644
--- a/src/org/kne/cloud/network/klalb/KLALBBlock.java
+++ b/src/org/kne/cloud/network/klalb/KLALBBlock.java
@@ -10,7 +10,7 @@ public class KLALBBlock {
volatile long time=-1;
- volatile Thread thread;
+ volatile TCPConnection connect;
public KLALBBlock(byte[] b, int size,long number) {
data=b;
this.size=size;
diff --git a/src/org/kne/cloud/network/klalb/KLALBClient.java b/src/org/kne/cloud/network/klalb/KLALBClient.java
index d52f3ca..b26ed1e 100644
--- a/src/org/kne/cloud/network/klalb/KLALBClient.java
+++ b/src/org/kne/cloud/network/klalb/KLALBClient.java
@@ -16,7 +16,28 @@ import java.util.UUID;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
+import org.kne.cloud.network.mport.IPPort;
+import org.kne.cloud.network.mport.ServiceElement;
+import org.kne.cloud.network.mport.ThreadTool;
+
public class KLALBClient {
+ private static Object olock=new Object();
+ static {
+ new Thread(()->{
+ while(true) {
+ try {
+ Thread.sleep(10000);
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ synchronized(olock) {
+ olock.notifyAll();
+ }
+ }
+ }).start();
+ }
+
+
private List tls=new ArrayList<>();
private TCPListener tcpl;
private ServiceElement sel;
@@ -45,8 +66,14 @@ public class KLALBClient {
tc.getDout().flush();
aig.incrementAndGet();
System.out.println("隧道"+tll+"已连接,可用线路数量:"+ (kcp.getTcps().size()+1));
+ try {
kcp.handleSocket(tc);
- tc.close();
+
+ }catch(IOException e){
+ e.printStackTrace();
+ }finally {
+ tc.close();
+ }
int n=kcp.getTcps().size();
System.out.println("隧道"+tll+"已断开,可用线路数量:"+ n);
if(n<=0) {
@@ -67,10 +94,12 @@ public class KLALBClient {
}
}
- try {
- Thread.sleep(1000);
- } catch (InterruptedException e) {
- e.printStackTrace();
+ synchronized(olock) {
+ try {
+ olock.wait();
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
}
}
}).start();
@@ -123,7 +152,7 @@ public class KLALBClient {
}
int s2=tcpc.getDin().readInt();
System.out.println("隧道列表:");
- for (int i = 0; i < s2; i++) {
+ for (int i = 0; i < s2; i++) {
String str=tcpc.getDin().readUTF();
System.out.println(str);
tls.add(Tunnel.newTunnel(str));
@@ -157,6 +186,7 @@ public class KLALBClient {
kc.setSel(kc.getServices().get(scn.nextInt()));
System.out.println("请输入你要映射的本地端口:");
kc.open(scn.nextInt());
+ System.out.println("服务已开启");
while(true) {
String s=scn.next();
switch(s) {
diff --git a/src/org/kne/cloud/network/klalb/KLALBCore.java b/src/org/kne/cloud/network/klalb/KLALBCore.java
index 7473e5c..d38e6ea 100644
--- a/src/org/kne/cloud/network/klalb/KLALBCore.java
+++ b/src/org/kne/cloud/network/klalb/KLALBCore.java
@@ -5,6 +5,7 @@ import java.io.DataOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
+import java.io.StreamCorruptedException;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
@@ -25,13 +26,15 @@ import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.atomic.AtomicInteger;
+import org.kne.cloud.network.mport.ThreadTool;
+
public class KLALBCore {
private volatile long inputcount = 1;
private volatile long outputcount = 1;
private Set inputcache = Collections.synchronizedSet(new HashSet<>());
- private List outputcache = Collections.synchronizedList(new ArrayList<>());
+ private List outputcache = new Vector<>();
- //private Map>acks=Collections.synchronizedMap(new WeakHashMap<>());
+ private BlockingQueue ackp=new LinkedBlockingQueue<>();
private BlockingQueue ackq=new LinkedBlockingQueue<>();
private volatile boolean inlocal=true;
@@ -47,7 +50,7 @@ public class KLALBCore {
if (closelocal)
throw new InterruptedException();
- while (outputcache.size() > cacheblocks) {
+ while (!outputcache.isEmpty()&&(outputcount-outputcache.get(0).number>cacheblocks)) {
if (closelocal)
Thread.currentThread().interrupt();
Thread.sleep(1);
@@ -76,13 +79,14 @@ public class KLALBCore {
}
}
if (b != null) {
- System.out.println("\tPROCESS:" + inputcache.size());
+ System.out.println("\tSORT:"+inputcount+" " + inputcache.size());
inputcount++;
return b;
}
if (closelocal)
Thread.currentThread().interrupt();
Thread.sleep(1);
+ //System.err.println(inputcount);
}
}finally {
@@ -90,6 +94,17 @@ public class KLALBCore {
}
}
+ private void makeAck(KLALBBlock klalbBlock) {
+ KLALBBlock klb1=new KLALBBlock(null, 0, -klalbBlock.number);
+ klb1.time=klalbBlock.time;
+ ackq.add(klb1);
+ if(ackp.size()<100) {
+ KLALBBlock klb2=new KLALBBlock(null, 0, -klalbBlock.number);
+ klb2.time=klalbBlock.time;
+ ackp.add(klb2);
+ }
+ }
+
public void sendDataBlock(TCPConnection out) throws IOException, InterruptedException {
if (closeremote)
@@ -112,24 +127,24 @@ public class KLALBCore {
for (int i = 0; i < outputcache.size(); i++) {
KLALBBlock kd = outputcache.get(i);
- if (kd.thread == null) {
+ if (kd.connect == null) {
kd.time = System.nanoTime();
- kd.thread = Thread.currentThread();
+ kd.connect = out;
ks = kd;
} else {
- if (kd.thread.isAlive()) {
+ if (kd.connect.isOpen()) {
long timex = (System.nanoTime() - kd.time) / 1000000;
- if (timex > 200000) {
- System.out.println("超时重传:"+kd);
+ if (timex > 100+1000*i) {
kd.time = System.nanoTime();
- kd.thread = Thread.currentThread();
+ kd.connect = out;
ks = kd;
+ System.out.println("超时重传:"+kd);
}
} else {
- System.out.println("掉线重传:"+kd);
kd.time = System.nanoTime();
- kd.thread = Thread.currentThread();
+ kd.connect = out;
ks = kd;
+ System.out.println("掉线重传:"+kd);
}
}
if (ks != null) {
@@ -139,7 +154,7 @@ public class KLALBCore {
}
if (ks != null) {
send0(out, ks);
- // System.out.println(outputcache.size());
+ ks.time = System.nanoTime();
}else {
Thread.sleep(1);
}
@@ -149,19 +164,36 @@ public class KLALBCore {
send0(out, new KLALBBlock(null,0 , 0));
}
}
-
+
+ public void sendDataBlockControlOnly(TCPConnection out) throws IOException, InterruptedException {
+ if (closeremote)
+ throw new InterruptedException();
+ //BlockingQueue bqk=acks.get(out);
+ while (ackp.isEmpty()) {
+ if (closeremote)
+ Thread.currentThread().interrupt();
+ Thread.sleep(1);
+ if(out.checkPingTime()) {
+ send0(out, new KLALBBlock(null,0 , 0));
+ }
+ }
+ KLALBBlock klb=ackp.poll();
+ if(klb!=null) {
+ send0(out, klb);
+ //System.out.println(klb+" "+ackp.size());
+ }
+ if(out.checkPingTime()) {
+ send0(out, new KLALBBlock(null,0 , 0));
+ }
+ }
+
public void receiveDataBlock(TCPConnection in) throws IOException, InterruptedException {
if (closeremote)
throw new InterruptedException();
KLALBBlock x = receive0(in);
long rect=System.nanoTime();
if (x.number > 0) {
-
- KLALBBlock klb=new KLALBBlock(null, 0, -x.number);
- ackq.add(klb);
-
- klb=new KLALBBlock(null, 0, -x.number);
- ackq.add(klb);
+
/*BlockingQueue bq=acks.get(in);
if(bq==null) {
BlockingQueue bqt=new LinkedBlockingQueue<>();
@@ -171,7 +203,6 @@ public class KLALBCore {
acks.get(in).add(klb);
}*/
-
if (x.number >= inputcount) {
while(inputcache.size()>(2*cacheblocks)&&inlocal) {
@@ -181,7 +212,8 @@ public class KLALBCore {
}
inputcache.add(x);
- }
+ }
+ makeAck(x);
} else if(x.number==0) {
ThreadTool.makeVThreadIfSupport("TACK", ()->{
@@ -192,7 +224,6 @@ public class KLALBCore {
klk.time=x.time+tw;
send0(in, klk);
} catch (IOException e) {
- e.printStackTrace();
}finally {
}
}).start();
@@ -211,14 +242,18 @@ public class KLALBCore {
private void send0(TCPConnection tcp, KLALBBlock kd) throws IOException {
DataOutputStream out=tcp.getDout();
synchronized (out) {
+
out.writeLong(kd.number);
if (kd.number > 0) {
+ out.writeLong(System.nanoTime());
out.writeInt(kd.size);
out.write(kd.data, 0, kd.size);
}else if(kd.number==0){
out.writeLong(System.nanoTime());
}else if(kd.number==Long.MIN_VALUE) {
out.writeLong(kd.time);
+ }else {
+ out.writeLong(kd.time);
}
out.flush();
@@ -230,15 +265,23 @@ public class KLALBCore {
KLALBBlock kb = new KLALBBlock();
DataInputStream in=tcp.getDin();
synchronized (in) {
+
kb.number = in.readLong();
if (kb.number > 0) {
+ kb.time=in.readLong();
kb.size = in.readInt();
+ if(kb.size<0||kb.size>65536) {
+ throw new StreamCorruptedException(tcp+"block size error:"+kb.size);
+ }
kb.data = new byte[kb.size];
in.readFully(kb.data);
}else if(kb.number==0){
kb.time=in.readLong();
}else if(kb.number==Long.MIN_VALUE) {
kb.time=in.readLong();
+ }else {
+ kb.time=in.readLong();
+ //tcp.nextRTT(System.nanoTime()-kb.time);
}
}
System.out.println(tcp.getTunnel().getName()+" RECEIVE:" + kb);
diff --git a/src/org/kne/cloud/network/klalb/KLALBSM.java b/src/org/kne/cloud/network/klalb/KLALBSM.java
index a9b9326..374ba21 100644
--- a/src/org/kne/cloud/network/klalb/KLALBSM.java
+++ b/src/org/kne/cloud/network/klalb/KLALBSM.java
@@ -9,6 +9,8 @@ import java.util.Map;
import java.util.Properties;
import java.util.Scanner;
+import org.kne.cloud.network.mport.ServiceElement;
+
public class KLALBSM {
public static Map services=new HashMap();
diff --git a/src/org/kne/cloud/network/klalb/KLALBServer.java b/src/org/kne/cloud/network/klalb/KLALBServer.java
index 7be4091..e4bcd2b 100644
--- a/src/org/kne/cloud/network/klalb/KLALBServer.java
+++ b/src/org/kne/cloud/network/klalb/KLALBServer.java
@@ -14,14 +14,18 @@ import java.util.Map;
import java.util.UUID;
import java.util.WeakHashMap;
+import org.kne.cloud.network.mport.IPPort;
+import org.kne.cloud.network.mport.ServiceElement;
+
public class KLALBServer {
WeakHashMap whm=new WeakHashMap<>();
public KLALBServer(int port, Map services,List tunnels) throws IOException {
TCPListener tcpl=new TCPListener(port);
tcpl.setCon((s)->{
+ TCPConnection tcc=null;
try {
//s.setSoTimeout(10000);
- TCPConnection tcc=new TCPConnection(null,s);
+ tcc=new TCPConnection(null,s);
DataInputStream din=tcc.getDin();
int val=din.readShort()&0xffff;
if(val!=59649) {
@@ -80,6 +84,7 @@ public class KLALBServer {
}
nx.handleSocket(tcc);
int n=nx.getTcps().size();
+ System.out.println(nx.getTcps());
if(n<=0) {
nx.closeLocal();
System.out.println("连接已关闭");
@@ -89,11 +94,8 @@ public class KLALBServer {
}catch(IOException e) {
e.printStackTrace();
}finally {
- try {
- s.close();
- } catch (IOException e) {
- e.printStackTrace();
- }
+ if(tcc!=null)
+ tcc.close();
}
});
tcpl.open();
diff --git a/src/org/kne/cloud/network/klalb/TCPConnection.java b/src/org/kne/cloud/network/klalb/TCPConnection.java
index 948d64a..4c97b5d 100644
--- a/src/org/kne/cloud/network/klalb/TCPConnection.java
+++ b/src/org/kne/cloud/network/klalb/TCPConnection.java
@@ -11,6 +11,11 @@ import java.net.UnknownHostException;
public class TCPConnection {
+ @Override
+ public String toString() {
+ return "TCPConnection [tunnel=" + tunnel + ", connect=" + connect + ", delay=" + delay + "]";
+ }
+
private Tunnel tunnel;
private Socket connect;
private DataInputStream din;
@@ -70,13 +75,13 @@ public class TCPConnection {
if(tunnel!=null) {
tunnel.setDelay(delay);
}
- /*if(connect!=null) {
+ if(connect!=null) {
try {
- connect.setSoTimeout((int) (delay/100000));
+ connect.setSoTimeout(10000);
} catch (SocketException e) {
e.printStackTrace();
}
- }*/
+ }
}
public void setTunnel(Tunnel tunnel) {
@@ -92,9 +97,8 @@ public class TCPConnection {
if(t!=null) {
t.getCCount().incrementAndGet();
}
- // connect.setSoTimeout(10000);
- din = new DataInputStream(new BufferedInputStream(connect.getInputStream(), 65536));
- dout = new DataOutputStream(new BufferedOutputStream(connect.getOutputStream(), 65536));
+ din = new DataInputStream(connect.getInputStream());
+ dout = new DataOutputStream(connect.getOutputStream());
}
public boolean isOpen() {
@@ -115,6 +119,11 @@ public class TCPConnection {
return false;
}
}
+ private volatile long RTT=10000000;
+ public void nextRTT(long NRTT) {
+ RTT=(long) ((1.0-Consts.A)*RTT+Consts.A*NRTT);
+ System.out.println(RTT);
+ }
diff --git a/src/org/kne/cloud/network/klalb/TCPListener.java b/src/org/kne/cloud/network/klalb/TCPListener.java
index 48121b1..d4f666e 100644
--- a/src/org/kne/cloud/network/klalb/TCPListener.java
+++ b/src/org/kne/cloud/network/klalb/TCPListener.java
@@ -5,6 +5,8 @@ import java.net.ServerSocket;
import java.net.Socket;
import java.util.function.Consumer;
+import org.kne.cloud.network.mport.ThreadTool;
+
public class TCPListener {
private int port;
private ServerSocket servers;
diff --git a/src/org/kne/cloud/network/klalb/Tunnel.java b/src/org/kne/cloud/network/klalb/Tunnel.java
index e0aa117..250f1d4 100644
--- a/src/org/kne/cloud/network/klalb/Tunnel.java
+++ b/src/org/kne/cloud/network/klalb/Tunnel.java
@@ -16,6 +16,8 @@ import java.util.WeakHashMap;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
+import org.kne.cloud.network.mport.IPPort;
+
public class Tunnel {
private String name;
private IPPort ipport;
@@ -73,7 +75,9 @@ public class Tunnel {
return name+"$"+ipport;
}
public Socket connectClientSocket() throws UnknownHostException, IOException {
- return new Socket(ipport.getIp(), ipport.getPort());
+ Socket socket=new Socket();
+ socket.connect(ipport.getSocketAddress(), 5000);
+ return socket;
}
@Override
diff --git a/src/org/kne/cloud/network/mport/ByteBufferRecycle.java b/src/org/kne/cloud/network/mport/ByteBufferRecycle.java
new file mode 100644
index 0000000..5f4d451
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/ByteBufferRecycle.java
@@ -0,0 +1,30 @@
+package org.kne.cloud.network.mport;
+
+import java.nio.ByteBuffer;
+import java.util.Stack;
+
+public class ByteBufferRecycle {
+ private int size,count;
+ private Stack s=new Stack();
+
+ public ByteBufferRecycle(int size,int count) {
+ super();
+ this.size = size;
+ this.count=count;
+ }
+
+ public ByteBuffer getBuffer() {
+ if(!s.isEmpty()) {
+ return s.pop();
+ }
+ return ByteBuffer.allocate(size);
+ }
+
+ public void recycle(ByteBuffer bbx) {
+
+ if(s.size() services=new ArrayList();
+ public static void main(String[] args) throws IOException {
+ File f=new File("ports.ini");
+ if(!f.exists()) {
+ f.createNewFile();
+ ept();
+ return;
+ }
+ if(f.length()<=0) {
+ ept();
+ return;//https://101eduyun.com
+ }
+ Scanner scn=new Scanner(f);
+ int remp=scn.nextInt();
+ System.out.println("端口复用程序V0.2");
+ System.out.println("开放端口:"+remp);
+ while(scn.hasNext() ) {
+ String s=scn.next();
+ ServiceElement se=new ServiceElement(s);
+ System.out.println(se);
+ services.add(se);
+ }
+ scn.close();
+ new PortRelay(remp,services).start();
+ System.out.println("提示:运行时输入reload可以重新加载配置文件,已经建立的TCP连接不受影响(暂不支持修改开放端口)");
+ Scanner scn2=new Scanner(System.in);
+ while(true) {
+ String command=scn2.next();
+ switch(command) {
+ case "reload":
+ synchronized (services) {
+ scn=new Scanner(f);
+ services.clear();
+ remp=scn.nextInt();
+ while(scn.hasNext() ) {
+ String s=scn.next();
+ ServiceElement se=new ServiceElement(s);
+ System.out.println(se);
+ services.add(se);
+ }
+ scn.close();
+ }
+ break;
+ }
+ }
+ }
+
+ private static void ept() {
+ System.out.println("请填写配置文件后再启动程序!");
+ System.out.println("格式:");
+ System.out.println("开放端口");
+ System.out.println("服务协议1$服务地址1$服务端口1");
+ System.out.println("服务协议2$服务地址2$服务端口2");
+ System.out.println("......");
+ System.out.println("服务协议n$服务地址n$服务端口n");
+ System.out.println("默认服务地址$默认服务端口");
+ System.out.println("");
+ System.out.println("按照顺序从上向下匹配,若都不是则匹配最后一个");
+ }
+
+}
diff --git a/src/org/kne/cloud/network/mport/PortRelay.java b/src/org/kne/cloud/network/mport/PortRelay.java
new file mode 100644
index 0000000..32bf1e5
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/PortRelay.java
@@ -0,0 +1,297 @@
+package org.kne.cloud.network.mport;
+
+import java.io.DataInputStream;
+import java.io.DataOutputStream;
+import java.io.IOException;
+import java.net.*;
+import java.nio.ByteBuffer;
+import java.nio.channels.*;
+import java.nio.channels.spi.AbstractInterruptibleChannel;
+import java.util.*;
+
+public class PortRelay {
+ private int port;
+ private List services;
+ private TCPListener ssc;
+ public PortRelay(int port, List services) throws IOException {
+ this.services = services;
+ this.port = port;
+ ssc=new TCPListener(port);
+ System.out.println("已打开端口:" + port);
+ }
+ public void start() throws IOException {
+ ssc.setCon((s)->{
+ DataInputStream din=null;
+ DataOutputStream dout=null;
+
+ Socket sl = new Socket();
+ DataInputStream dinl=null;
+ DataOutputStream doutl=null;
+
+ ServiceElement se=null;
+ try {
+ final DataInputStream xdin=new DataInputStream(s.getInputStream());
+ din=xdin;
+ final DataOutputStream xdout =new DataOutputStream(s.getOutputStream());
+ dout=xdout;
+ int n=din.readShort()&0xffff;
+ se = ServiceElement.getbyproc(services, n);
+ InetSocketAddress sa=(InetSocketAddress) s.getRemoteSocketAddress();
+ System.out.println(sa.getAddress().getHostAddress() + ":" + port
+ + "-" + se.proc + "(" + n + ")->" + se.ipport.getIp().getHostAddress() + ":" + se.ipport.getPort());
+ sl.connect(new InetSocketAddress(se.ipport.getIp().getHostAddress(), se.ipport.getPort()));
+ final DataInputStream xdinl=new DataInputStream(sl.getInputStream());
+ dinl=xdinl;
+ final DataOutputStream xdoutl =new DataOutputStream(sl.getOutputStream());
+ doutl=xdoutl;
+ Thread tl=ThreadTool.makeVThreadIfSupport("远程->本地 转发线程", ()->{
+ try {
+ xdoutl.writeShort(n);
+ int len;
+ byte[]b=new byte[8192];
+ while((len=xdin.read(b))!=-1) {
+ xdoutl.write(b,0,len);
+ }
+ } catch (IOException e) {
+ }finally {
+ if(xdin!=null)
+ try {
+ xdin.close();
+ } catch (IOException e) {
+ }
+ if(xdoutl!=null)
+ try {
+ xdoutl.close();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ Thread tr=ThreadTool.makeVThreadIfSupport("本地->远程 转发线程", ()->{
+ try {
+ int len;
+ byte[]b=new byte[8192];
+ while((len=xdinl.read(b))!=-1) {
+ xdout.write(b,0,len);
+ }
+ } catch (IOException e) {
+ }finally {
+ if(xdinl!=null)
+ try {
+ xdinl.close();
+ } catch (IOException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }
+ if(xdout!=null)
+ try {
+ xdout.close();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ tl.start();
+ tr.start();
+ try {
+ tl.join();
+ tr.join();
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ } catch (IOException e) {
+ e.printStackTrace();
+ }finally {
+
+
+ if(sl!=null)
+ try {
+ sl.close();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+
+
+
+ try {
+ s.close();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ ssc.open();
+ }
+ /*public void loop() {
+ Timer timer=new Timer();
+ TimerTask tmk = new TimerTask() {
+ @Override
+ public void run() {
+ try {
+ int v = selector.select();
+ if (v == 0)
+ return;
+ Set sset = selector.selectedKeys();
+ Iterator it = sset.iterator();
+ while (it.hasNext()) {
+ SelectionKey selectionKey = (SelectionKey) it.next();
+ it.remove();
+ if (!selectionKey.isValid()) {
+ selectionKey.cancel();
+ continue;
+ } else if (selectionKey.isConnectable()) {
+ try {
+ SocketChannel sccc=(SocketChannel) selectionKey.channel();
+ if(sccc.isConnectionPending()) {
+ ((SocketChannel) selectionKey.channel()).finishConnect();
+ }else {
+ IOState stt=(IOState) selectionKey.attachment();
+ stt.closeA();
+ selectionKey.cancel();
+ }
+
+ }catch(Exception e) {
+ e.printStackTrace();
+ IOState stt=(IOState) selectionKey.attachment();
+ stt.closeA();
+ }
+ } else if (selectionKey.isAcceptable()) {
+ ServerSocketChannel scc = (ServerSocketChannel) selectionKey.channel();
+ SocketChannel sc = scc.accept();
+ if (sc != null) {
+ sc.configureBlocking(false);
+
+ sc.register(selector, SelectionKey.OP_READ | SelectionKey.OP_WRITE, new IOState(sc) {
+ ServiceElement se = null;
+ ByteBuffer bbf;
+ SocketChannel sssl, sccl;
+ Queue ts = new LinkedList();
+ Queue tc = new LinkedList();
+
+ @Override
+ public void loopRead() throws IOException {
+ if (sssl == null) {
+ sssl = (SocketChannel) getChannel();
+ }
+ if(ts.size()>10)
+ return;
+ if (bbf == null)
+ bbf = rec.getBuffer();
+ try {
+ int ed = sssl.read(bbf);
+ if (ed == -1) {
+
+ sssl.close();
+ if(ts.size()<=0&&sccl!=null)
+ sccl.close();
+ }
+
+ } catch (IOException e) {
+ e.printStackTrace();
+ if(sssl!=null)
+ sssl.close();
+ if(sccl!=null)
+ sccl. close();
+
+ }
+ if (bbf.position() > 0) {
+
+ bbf.flip();
+ if (se == null && bbf.remaining() >= 2) {
+ int n = bbf.getShort()&0xffff;
+ se = ServiceElement.getbyproc(services, n);
+ System.out.println(((InetSocketAddress) ((SocketChannel) getChannel())
+ .getRemoteAddress()).getAddress().getHostAddress() + ":" + port
+ + "-" + se.proc + "(" + n + ")->" + se.ipport.getIp().getHostAddress() + ":" + se.ipport.getPort());
+ sccl = SocketChannel.open();
+ sccl.configureBlocking(false);
+ sccl.register(selector, SelectionKey.OP_CONNECT | SelectionKey.OP_READ
+ | SelectionKey.OP_WRITE, new IOState(sccl) {
+ ByteBuffer bbx;
+
+ @Override
+ public void loopRead() throws IOException {
+ if (bbx == null)
+ bbx = rec.getBuffer();
+ if(tc.size()>10)
+ return;
+ try {
+ int ed = getChannel().read(bbx);
+ if (ed == -1) {
+ sccl.close();
+ if(tc.size()<=0&&sssl!=null)
+ sssl.close();
+ }
+ } catch (IOException e) {
+ e.printStackTrace();
+ if(sssl!=null)
+ sssl.close();
+ if(sccl!=null)
+ sccl. close();
+ }
+ if (bbx.position() > 0) {
+ bbx.flip();
+ tc.add(bbx);
+ bbx = null;
+ }
+ }
+
+ @Override
+ public void loopWrite() throws IOException {
+ ByteBuffer bbf = ts.poll();
+ if (bbf == null)
+ return;
+ getChannel().write(bbf);
+ rec.recycle(bbf);
+ }
+
+ @Override
+ protected void closeA() throws IOException{
+ if(sssl!=null)
+ sssl.close();
+ if(sccl!=null)
+ sccl. close();
+ }
+
+ });
+ sccl.connect(new InetSocketAddress(se.ipport.getIp().getHostAddress(), se.ipport.getPort()));
+ bbf.rewind();
+ }
+ ts.add(bbf);
+ bbf = null;
+ }
+ }
+
+ @Override
+ public void loopWrite() throws IOException {
+ ByteBuffer bbx = tc.poll();
+ if (bbx == null)
+ return;
+ getChannel().write(bbx);
+ rec.recycle(bbx);
+ }
+
+ @Override
+ protected void closeA() throws IOException {
+ if(sssl!=null)
+ sssl.close();
+ if(sccl!=null)
+ sccl. close();
+ }
+
+ });
+ }
+ } else if (selectionKey.isReadable()) {
+ ((IOState) selectionKey.attachment()).loopRead();
+ } else if (selectionKey.isWritable()) {
+ ((IOState) selectionKey.attachment()).loopWrite();
+ }
+ }
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ }
+ };
+ timer.scheduleAtFixedRate(tmk, 0, 1);
+ }*/
+}
diff --git a/src/org/kne/cloud/network/klalb/ServiceElement.java b/src/org/kne/cloud/network/mport/ServiceElement.java
similarity index 93%
rename from src/org/kne/cloud/network/klalb/ServiceElement.java
rename to src/org/kne/cloud/network/mport/ServiceElement.java
index 5972e4e..5096952 100644
--- a/src/org/kne/cloud/network/klalb/ServiceElement.java
+++ b/src/org/kne/cloud/network/mport/ServiceElement.java
@@ -1,4 +1,4 @@
-package org.kne.cloud.network.klalb;
+package org.kne.cloud.network.mport;
import java.net.UnknownHostException;
import java.util.HashMap;
diff --git a/src/org/kne/cloud/network/mport/TCPListener.java b/src/org/kne/cloud/network/mport/TCPListener.java
new file mode 100644
index 0000000..0b84bc6
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/TCPListener.java
@@ -0,0 +1,65 @@
+package org.kne.cloud.network.mport;
+
+import java.io.IOException;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.util.function.Consumer;
+
+public class TCPListener {
+ private int port;
+ private ServerSocket servers;
+ private volatile boolean flag=false;
+
+ private Consumercon;
+ private Runnable r=new Runnable() {
+ @Override
+ public void run() {
+ while(flag){
+ try {
+ Socket soce=servers.accept();
+ ThreadTool.makeVThreadIfSupport("端口监听线程",()->{
+ con.accept(soce);
+ }).start();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ }
+ }
+ };
+ public TCPListener(int port) throws IOException {
+ this.port=port;
+ }
+
+ public void open() throws IOException {
+ flag=true;
+ servers=new ServerSocket(port);
+ new Thread(r).start();
+ }
+ public void close() {
+ flag=false;
+ if(servers!=null) {
+ try {
+ servers.close();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ }
+ }
+
+
+ public int getPort() {
+ return port;
+ }
+ public void setPort(int port) {
+ this.port = port;
+ }
+
+ public Consumer getCon() {
+ return con;
+ }
+
+ public void setCon(Consumer con) {
+ this.con = con;
+ }
+
+}
diff --git a/src/org/kne/cloud/network/klalb/ThreadTool.java b/src/org/kne/cloud/network/mport/ThreadTool.java
similarity index 89%
rename from src/org/kne/cloud/network/klalb/ThreadTool.java
rename to src/org/kne/cloud/network/mport/ThreadTool.java
index d28cf69..2256f20 100644
--- a/src/org/kne/cloud/network/klalb/ThreadTool.java
+++ b/src/org/kne/cloud/network/mport/ThreadTool.java
@@ -1,4 +1,4 @@
-package org.kne.cloud.network.klalb;
+package org.kne.cloud.network.mport;
public class ThreadTool {
public static boolean first=true;