diff --git a/.classpath b/.classpath
index 7a21a4e..0039a3b 100644
--- a/.classpath
+++ b/.classpath
@@ -1,9 +1,9 @@
-
+
-
+
diff --git a/.settings/org.eclipse.jdt.core.prefs b/.settings/org.eclipse.jdt.core.prefs
index c59d0c6..85dd98e 100644
--- a/.settings/org.eclipse.jdt.core.prefs
+++ b/.settings/org.eclipse.jdt.core.prefs
@@ -1,9 +1,9 @@
eclipse.preferences.version=1
org.eclipse.jdt.core.compiler.codegen.inlineJsrBytecode=enabled
org.eclipse.jdt.core.compiler.codegen.methodParameters=do not generate
-org.eclipse.jdt.core.compiler.codegen.targetPlatform=1.8
+org.eclipse.jdt.core.compiler.codegen.targetPlatform=19
org.eclipse.jdt.core.compiler.codegen.unusedLocal=preserve
-org.eclipse.jdt.core.compiler.compliance=1.8
+org.eclipse.jdt.core.compiler.compliance=19
org.eclipse.jdt.core.compiler.debug.lineNumber=generate
org.eclipse.jdt.core.compiler.debug.localVariable=generate
org.eclipse.jdt.core.compiler.debug.sourceFile=generate
@@ -12,4 +12,4 @@ org.eclipse.jdt.core.compiler.problem.enablePreviewFeatures=disabled
org.eclipse.jdt.core.compiler.problem.enumIdentifier=error
org.eclipse.jdt.core.compiler.problem.reportPreviewFeatures=warning
org.eclipse.jdt.core.compiler.release=disabled
-org.eclipse.jdt.core.compiler.source=1.8
+org.eclipse.jdt.core.compiler.source=19
diff --git a/KLALB协议规范V2.0.docx b/KLALB协议规范V2.0.docx
index 0375c53..0886dc0 100644
Binary files a/KLALB协议规范V2.0.docx and b/KLALB协议规范V2.0.docx differ
diff --git a/klalbclient.ini b/klalbclient.ini
new file mode 100644
index 0000000..b6ecda0
--- /dev/null
+++ b/klalbclient.ini
@@ -0,0 +1,3 @@
+#Tue Jul 25 19:03:09 CST 2023
+local=0.0.0.0\:35000
+server=127.0.0.1\:4569
diff --git a/kserver - 副本.ini b/kserver - 副本.ini
new file mode 100644
index 0000000..944768a
--- /dev/null
+++ b/kserver - 副本.ini
@@ -0,0 +1,3 @@
+virtualip=67c:72ce:765e:4db1:a02e:92fa:1959:29eb
+bind=0.0.0.0:4569
+local=127.0.0.1:36555
diff --git a/kserver.ini b/kserver.ini
index 7e44418..944768a 100644
--- a/kserver.ini
+++ b/kserver.ini
@@ -1,3 +1,3 @@
virtualip=67c:72ce:765e:4db1:a02e:92fa:1959:29eb
bind=0.0.0.0:4569
-local=127.0.0.1:5212
+local=127.0.0.1:36555
diff --git a/linetable - 副本.txt b/linetable - 副本.txt
new file mode 100644
index 0000000..320126c
--- /dev/null
+++ b/linetable - 副本.txt
@@ -0,0 +1,19 @@
+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
\ No newline at end of file
diff --git a/linetable.txt b/linetable.txt
new file mode 100644
index 0000000..cb2c257
--- /dev/null
+++ b/linetable.txt
@@ -0,0 +1,20 @@
+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
diff --git a/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java b/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java
new file mode 100644
index 0000000..c4490f6
--- /dev/null
+++ b/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java
@@ -0,0 +1,28 @@
+package org.kne.cloud.network.klalb;
+
+import java.util.concurrent.ConcurrentLinkedQueue;
+
+public class ByteArrayRecycle {
+ private ConcurrentLinkedQueuerec=new ConcurrentLinkedQueue<>();
+ private int capacity;
+ private int length;
+ public ByteArrayRecycle(int capacity, int length) {
+ super();
+ this.capacity = capacity;
+ this.length = length;
+ }
+ public synchronized void recycle(byte[]b) {
+ if(b.length!=length)
+ throw new IllegalArgumentException("wrong length");
+ if(rec.size()"+dport+" "+number+"["+data.length+"]";
+ return "DATAT "+sport+"->"+dport+" "+number+"["+size+"]";
}
@Override
@@ -52,8 +55,8 @@ public class DATATPacket extends KLALBPacket {
dto.writeInt(sport);
dto.writeInt(dport);
dto.writeLong(number);
- dto.writeChar(data.length);
- dto.write(data);
+ dto.writeChar(size);
+ dto.write(data,0,size);
}
@Override
@@ -62,7 +65,11 @@ public class DATATPacket extends KLALBPacket {
sport=din.readInt();
dport=din.readInt();
number=din.readLong();
- data=new byte[din.readChar()];
- din.readFully(data);
+ size=din.readChar();
+ data=arrayRecycle.create();
+ din.readFully(data,0,size);
+ }
+ public int getSize() {
+ return size;
}
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBController.java b/src/org/kne/cloud/network/klalb/KLALBController.java
index 46661d3..0b02fce 100644
--- a/src/org/kne/cloud/network/klalb/KLALBController.java
+++ b/src/org/kne/cloud/network/klalb/KLALBController.java
@@ -264,15 +264,21 @@ public class KLALBController {
if (l == null || l.isEmpty()) {
throw new NoRouteToHostException("address unreachable: " + addr);
}
- int count0 = Math.min(count, l.size());
List l2 = (List) ((ArrayList) l).clone();
- LineDecitionComparator ldc = new LineDecitionComparator(l2, priority);
+ l2.removeAll(packet.getSendRecord());
+ if(l2.isEmpty()) {
+ l2 = (List) ((ArrayList) l).clone();
+ }
+ int count0 = Math.min(count, l2.size());
+ LineDecitionComparator ldc = new LineDecitionComparator(l2,packet, priority);
Collections.sort(l2,ldc);
//System.out.println(l2);
for (Iterator iterator = l2.iterator(); iterator.hasNext();) {
KLALBRemoteSocket krst = (KLALBRemoteSocket) iterator.next();
krst.sendPacket(packet, priority);
+ packet.getSendRecord().add(krst);
count0--;
+ Thread.yield();
if (count0 <= 0)
break;
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBPacket.java b/src/org/kne/cloud/network/klalb/KLALBPacket.java
index 26d5fd7..8bc5ac0 100644
--- a/src/org/kne/cloud/network/klalb/KLALBPacket.java
+++ b/src/org/kne/cloud/network/klalb/KLALBPacket.java
@@ -6,6 +6,7 @@ import java.io.Externalizable;
import java.io.IOException;
import java.io.ObjectInput;
import java.io.ObjectOutput;
+import java.util.Vector;
public abstract class KLALBPacket{
public static final int PING=0;
@@ -40,5 +41,11 @@ public abstract class KLALBPacket{
public long getLength() {
return 1;
}
+ private Vector sendRecord=new Vector<>();
+
+ public Vector getSendRecord() {
+ return sendRecord;
+ }
+
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java b/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java
index 1c0b26d..07baf94 100644
--- a/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java
+++ b/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java
@@ -35,13 +35,12 @@ public KLALBRemoteSocket(Socket socket) {
public KLALBRemoteSocket(Socket socket,Monitor m) {
super(socket,m);
try {
- socket.setSendBufferSize(32768);
- socket.setReceiveBufferSize(65535);
+ //socket.setSendBufferSize(32768);
+ //socket.setReceiveBufferSize(65535);
//System.out.println(socket.getSendBufferSize()+" "+socket.getReceiveBufferSize());
socket.setTcpNoDelay(true);
socket.setSoTimeout(10000);
} catch (SocketException e1) {
- // TODO 自动生成的 catch 块
e1.printStackTrace();
}
ThreadTool.makeVDaemonThreadIfSupport("远程接收线程", () -> {
@@ -49,20 +48,20 @@ public KLALBRemoteSocket(Socket socket) {
while (!isClosed()) {
KLALBPacket kpp = getInputStream().readPacket();
- //System.out.println("RX:" + kpp);
if(kpp instanceof PINGPacket) {
sendPacket(new PONGPacket(((PINGPacket) kpp).getTime()), 65540);
}else if(kpp instanceof PONGPacket) {
PONGPacket png=(PONGPacket) kpp;
getMonitor().updateLatency (System.nanoTime()- png.getTime());
try {
- socket.setSoTimeout(1100);
+ socket.setSoTimeout(600);
} catch (SocketException e1) {
// TODO 自动生成的 catch 块
e1.printStackTrace();
}
}else {
+ System.out.println("RX:" + kpp);
if(kpp instanceof VADDRPacket) {
remoteVaddr=((VADDRPacket) kpp).getVaddr();
}
@@ -94,15 +93,19 @@ public KLALBRemoteSocket(Socket socket) {
PQItem pqi=sendQueue.poll();
if(pqi!=null) {
KLALBPacket kpp=pqi.getPacket();
- //if(!(kpp instanceof PONGPacket))
- //System.out.println("TX:" + kpp);
+ if(!(kpp instanceof PONGPacket))
+ System.out.println("TX:" + kpp);
if(kpp instanceof PONGPacket) {
PONGPacket pp=(PONGPacket) kpp;
pp.redeltaTime(System.nanoTime()-pqi.getAddtime());
}
getOutputStream().writePacket(kpp);
+ //System.out.println(sendQueue.size());
}else {
- Thread.sleep(1);
+ synchronized (sendQueue) {
+ sendQueue.wait(10);
+ }
+ //Thread.sleep(1);
}
}
}
@@ -117,7 +120,6 @@ public KLALBRemoteSocket(Socket socket) {
}
}
}).start();
-
}
@@ -178,6 +180,9 @@ public KLALBRemoteSocket(Socket socket) {
public void sendPacket(KLALBPacket kp,int priority) {
sendQueue.add(new PQItem(kp, priority));
+ synchronized (sendQueue) {
+ sendQueue.notifyAll();
+ }
}
public void setPacketReceiver(Consumer rec) {
@@ -194,11 +199,11 @@ public KLALBRemoteSocket(Socket socket) {
super.close();
sendQueue.forEach((x)->{
try {
+ //System.out.println("断线重发");
controller.sendPacketToAddress(remoteVaddr, x.getPacket(), x.getPriority()+1);
} catch (NoRouteToHostException e) {
}
});
-
}
@Override
diff --git a/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java b/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java
index 093b1c5..3660b3c 100644
--- a/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java
+++ b/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java
@@ -31,12 +31,17 @@ import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.atomic.AtomicReferenceArray;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
+import java.util.zip.Deflater;
+import java.util.zip.DeflaterOutputStream;
+import java.util.zip.Inflater;
+import java.util.zip.InflaterInputStream;
+import org.kne.cloud.network.mport.VirtualSocketImpl;
import org.kne.io.Data;
import java.util.*;
-public class KLALBVirtualSocketImpl extends SocketImpl {
+public class KLALBVirtualSocketImpl extends VirtualSocketImpl {
private KLALBController controller;
protected int getInputchachesize() {
@@ -55,8 +60,19 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
this.outputchachesize = outputchachesize;
}
- private int inputchachesize = 5773 * 100;
- private int outputchachesize = 5773 * 100;
+ private volatile int inputchachesize = LIMIT * 500;
+ private volatile int outputchachesize = LIMIT * 200;
+ private volatile boolean nodelay=false;
+ private volatile long delaytime=1;
+
+ public long getDelaytime() {
+ return delaytime;
+ }
+
+ public void setDelaytime(long delaytime) {
+ this.delaytime = delaytime;
+ }
+
private Inet6Address bindaddr;{
try {
bindaddr=(Inet6Address) Inet6Address.getByName("::0");
@@ -77,15 +93,20 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
try {
sendlist.get(i).check(i);
} catch (IOException e) {
- // TODO 自动生成的 catch 块
e.printStackTrace();
+ try {
+ close();
+ } catch (IOException e1) {
+ e1.printStackTrace();
+ }
+ break;
}
}
}
}
}
- private boolean succeed, refused;
+ private volatile boolean succeed, refused;
protected boolean isListening() {
return backlogQueue != null;
@@ -93,7 +114,7 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
private List inputchache = new ArrayList<>();
private long inputcount = 0;
- private boolean avaliable = true;
+ private volatile boolean avaliable = true;
private BiConsumer packReceiver = new BiConsumer() {
@Override
@@ -104,15 +125,15 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
if (backlogQueue.offer(new InetSocketAddress(from, ((SYNTPacket) u).getSport()))) {
controller.sendPacketToAddress(from,
- new SACKTPacket(localport, ((SYNTPacket) u).getSport()), 65537);
+ new SACKTPacket(localport, ((SYNTPacket) u).getSport()), 65537,2);
} else {
controller.sendPacketToAddress(from, new RSTPacket(localport, ((SYNTPacket) u).getSport()),
- 65537);
+ 65537,2);
}
} else {
controller.sendPacketToAddress(from, new RSTPacket(localport, ((SYNTPacket) u).getSport()),
- 65537);
+ 65537,2);
}
} else if (u instanceof SACKTPacket) {
if (connecting) {
@@ -148,23 +169,34 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
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,2);
+ countInputBytes() < inputchachesize), 32768,1);
} else if (u instanceof ACKTPacket) {
ACKTPacket ackt = (ACKTPacket) u;
avaliable = ackt.isAvaliable();
- AtomicReferencekl=new AtomicReference<>();
+ 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();
@@ -177,14 +209,14 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
private int countInputBytes() {
AtomicInteger i = new AtomicInteger(0);
sendDeque.forEach((c) -> {
- i.addAndGet(c.getData().length);
+ i.addAndGet(c.getSize());
});
return i.get();
}
private int countOutputBytes() {
AtomicInteger i = new AtomicInteger(0);
sendlist.forEach((c) -> {
- i.addAndGet(c.getPacket().getData().length);
+ i.addAndGet(c.getPacket().getSize());
});
return i.get();
}
@@ -204,23 +236,34 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
@Override
public void setOption(int optID, Object value) throws SocketException {
- if(optID==SocketOptions.SO_RCVBUF) {
+ switch(optID) {
+ case SocketOptions.TCP_NODELAY:
+ nodelay=(boolean) value;
+ break;
+ case SocketOptions.SO_RCVBUF:
inputchachesize=(int) value;
- }else if(optID==SocketOptions.SO_SNDBUF) {
+ break;
+ case SocketOptions.SO_SNDBUF:
outputchachesize=(int)value;
+ break;
}
}
@Override
public Object getOption(int optID) throws SocketException {
- if(optID==SocketOptions.SO_BINDADDR) {
- return bindaddr;
- }else if(optID==SocketOptions.SO_RCVBUF) {
+ switch(optID) {
+ case SocketOptions.TCP_NODELAY:
+ return nodelay;
+ case SocketOptions.SO_RCVBUF:
return inputchachesize;
- }else if(optID==SocketOptions.SO_SNDBUF) {
+ case SocketOptions.SO_SNDBUF:
return outputchachesize;
+ case SocketOptions.SO_BINDADDR:
+ return bindaddr;
+ default:
+ return null;
}
- return null;
+
}
@Override
@@ -262,7 +305,6 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
}
connecting = false;
if (succeed) {
-
} else if (refused) {
throw new ConnectException("connect refused");
} else {
@@ -329,8 +371,8 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
this.acceptedSocketCloseListener = lsr;
}
- private KVSIInputStream vin;
- private KVSIOutputStream vout;
+ private InputStream vin;
+ private OutputStream vout;
private class KVSIInputStream extends InputStream {
private DATATPacket dtp = null;
@@ -339,9 +381,7 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
private boolean shutdown=false;
@Override
public int read() throws IOException {
- if(shutdown)
- return -1;
- if (dtp == null || (count >= dtp.getData().length && dtp.getData().length != 0)) {
+ if (dtp == null ) {
count = 0;
while (true) {
if (isClosed())
@@ -356,16 +396,24 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
break;
}
try {
- Thread.sleep(1);
+ synchronized (sendDeque) {
+ sendDeque.wait(100);
+ }
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
- if (dtp.getData().length == 0) {
+ if (dtp.getSize() == 0) {
return -1;
- } else
- return dtp.getData()[count++] & 0xff;
+ } else {
+ int ret= dtp.getData()[count++] & 0xff;
+ if(count==dtp.getSize()) {
+ DATATPacket.arrayRecycle.recycle(dtp.getData());
+ dtp=null;
+ }
+ return ret;
+ }
}
@Override
@@ -379,21 +427,78 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
}
len = Math.min(len, available());
-
- int c = read();
- if (c == -1) {
- return -1;
+
+
+ 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();
+ }
+ }
+ }
+ if (dtp.getSize() == 0) {
+ return -1;
+ } else {
+ b[off]= dtp.getData()[count++] ;
+ if(count==dtp.getSize()) {
+ DATATPacket.arrayRecycle.recycle(dtp.getData());
+ dtp=null;
+ }
}
- b[off] = (byte) c;
-
int i = 1;
try {
for (; i < len; i++) {
- c = read();
- if (c == -1) {
- break;
+
+ 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();
+ }
+ }
+ }
+ if (dtp.getSize() == 0) {
+ break;
+ } else {
+ b[off + i]= dtp.getData()[count++] ;
+ if(count==dtp.getSize()) {
+ DATATPacket.arrayRecycle.recycle(dtp.getData());
+ dtp=null;
+ }
}
- b[off + i] = (byte) c;
}
} catch (IOException ee) {
}
@@ -402,40 +507,92 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
@Override
public void close() throws IOException {
- shutdown=true;
+
}
@Override
public int available() throws IOException {
AtomicInteger i = new AtomicInteger(0);
sendDeque.forEach((V) -> {
- i.addAndGet(V.getData().length);
+ i.addAndGet(V.getSize());
});
if (dtp != null)
- i.addAndGet(dtp.getData().length - count);
+ i.addAndGet(dtp.getSize() - count);
return i.get();
}
}
- private long outputcount = 0;
+ private volatile long outputcount = 0;
+ private static final int LIMIT=65535;
private class KVSIOutputStream extends OutputStream {
-
- private byte[] cache=new byte[65535];
+ private byte[] cache=DATATPacket.arrayRecycle.create();
private int count=0;
+ private Object lock=new Object();
@Override
public void write(int b) throws IOException {
if (isClosed())
throw new SocketException("Socket is closed");
+ synchronized (lock) {
+
cache[count++]=(byte) b;
- if (count >=5773 ) {//1429?
+ if (count >= LIMIT) {//1429?5773?8669
+ flush0();
+ }else {
flush();
+ }
}
}
+
+ @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) {
+ for (int i = off; i < ol; i++) {
+ cache[count++]=b[i];
+ if (count >= LIMIT) {//1429?5773?8669
+ flush0();
+ }
+ }
+ flush();
+ }
+ }
+
+ 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();
+ }
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ }
+ };
+ new Timer("粘包计时线程").scheduleAtFixedRate(tt, 0, delaytime);
+ }
+ }
+ }
+
+
+ private void flush0() throws IOException {
if (count > 0) {
try {
while (!avaliable) {
@@ -446,20 +603,18 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
}
try {
while(countOutputBytes()>outputchachesize) {
- Thread.sleep(1);
+ synchronized (sendlist) {
+ sendlist.wait(10);
+ }
}
} catch (InterruptedException e) {
// TODO 自动生成的 catch 块
e.printStackTrace();
}
- byte[]ba;
- if(count==cache.length) {
- ba=cache;
- cache=new byte[cache.length];
- }else {
- ba=Arrays.copyOf(cache, count);
- }
- SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, ba), 5,10);
+ byte[] ba=cache;
+ cache=DATATPacket.arrayRecycle.create();
+
+ SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, ba,count), 5,20);
x.run();
sendlist.add(x);
/*controller.sendPacketToAddress((Inet6Address) address,
@@ -470,28 +625,36 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
@Override
public void close() throws IOException {
- flush();
- SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, new byte[0]), 5,10);
+ flush0();
+ SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, DATATPacket.arrayRecycle.create(),0), 5,20);
x.run();
sendlist.add(x);
/*controller.sendPacketToAddress((Inet6Address) address,
new DATATPacket(localport, port, outputcount++, new byte[0]), 5);*/
+ tt.cancel();
}
}
@Override
- protected InputStream getInputStream() throws IOException {
+ public InputStream getInputStream() throws IOException {
if (vin == null) {
vin = new KVSIInputStream();
+ int val= vin.read();
+ if(val==1) {
+ vin=new InflaterInputStream(vin,new Inflater(true),LIMIT);
+ }
}
return vin;
}
@Override
- protected OutputStream getOutputStream() throws IOException {
+ public OutputStream getOutputStream() throws IOException {
if (vout == null) {
vout = new KVSIOutputStream();
+ vout.write(1);
+ vout.flush();
+ vout=new DeflaterOutputStream(vout, new Deflater(Deflater.BEST_COMPRESSION, true), LIMIT, true);
}
return vout;
}
@@ -521,17 +684,23 @@ public class KLALBVirtualSocketImpl extends SocketImpl {
protected void close() throws IOException {
if (!isListening() && !isClosed()) {
closed = true;
+ try {
+
controller.sendPacketToAddress((Inet6Address) super.address, new RSTPacket(super.localport, super.port),
- 65537);
+ 65537,2);
+ } catch (NoRouteToHostException e) {
+ }
+ if (acceptedSocketCloseListener != null)
+ acceptedSocketCloseListener.accept(this);
+ else
+ controller.unbind(this);
}
- if (acceptedSocketCloseListener != null)
- acceptedSocketCloseListener.accept(this);
- else
- controller.unbind(this);
sendCheckTask.cancel();
+
+//new Exception().printStackTrace();
}
- private boolean closed;
+ private volatile boolean closed;
private boolean isClosed() {
return closed;
diff --git a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java
index 263a909..2999198 100644
--- a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java
+++ b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java
@@ -8,38 +8,42 @@ import java.util.Map;
import java.util.concurrent.atomic.AtomicLong;
public class LineDecitionComparator implements Comparator {
- private Map predictedlatencys=new HashMap<>();
- public LineDecitionComparator (List krs,int priority) {
+ 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();
- long x=klalbRemoteSocket.getMonitor().getLatency()>>1;
- x+=klalbRemoteSocket.getMonitor().getJitter()>>1;
+ 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());
}
});
- long speed=klalbRemoteSocket.getMonitor(). getOutSpeedMax();
+ al.addAndGet(curr.getLength());
+ long speed=klalbRemoteSocket.getMonitor(). getOutSpeed();
if(speed==0) {
if(al.get()>0) {
x=Long.MAX_VALUE;
}
}else {
- x+=al.get()*1000000000/speed;
+ x+=al.get()*1000000000.0/speed;
}
- //System.out.println(x);
+
+ /*System.out.println(klalbRemoteSocket);
+ System.out.println(klalbRemoteSocket.getSendQueue().size());
+ System.out.println(x);*/
predictedlatencys.put(klalbRemoteSocket, x);
}
}
- public Map getPredictedlatencys() {
+ public Map getPredictedlatencys() {
return predictedlatencys;
}
@Override
public int compare(KLALBRemoteSocket o1, KLALBRemoteSocket o2) {
- long t1=predictedlatencys.get(o1);
- long t2=predictedlatencys.get(o2);
+ double t1=predictedlatencys.get(o1);
+ double t2=predictedlatencys.get(o2);
if(t1>t2) {
return 1;
}else if(t1{
}
}
public synchronized void addHostPort(MultipurposeSocketAddress v) {
-
LineEntry le=new LineEntry(v);
if(putIfAbsent(v,le)==null)
le.open();
@@ -74,10 +73,11 @@ public class LineManager extends Hashtable{
private MultipurposeSocketAddress hp;
private KLALBRemoteSocket krs;
private volatile boolean flag=true;
- private volatile boolean online=false;
private Monitor monitor=new Monitor();
- public boolean isOnline() {
- return online;
+
+
+ public Monitor getMonitor() {
+ return monitor;
}
public MultipurposeSocketAddress getHp() {
@@ -90,7 +90,6 @@ public class LineManager extends Hashtable{
}
public String toString() {
StringBuilder sb=new StringBuilder();
- sb.append(online?"●在线 ":"○离线 ");
sb.append(hp.toString());
sb.append('\t');
sb.append(monitor.toString());
@@ -102,7 +101,6 @@ public class LineManager extends Hashtable{
sb.append(hp.toString());
sb.append('\n');
- sb.append(online?"●在线\t":"○离线\t");
sb.append(monitor.toString2());
return sb.toString();
}
@@ -111,24 +109,34 @@ public class LineManager extends Hashtable{
public void run() {
while(true){
try {
- krs=new KLALBRemoteSocket(hp.connectSocket(),monitor);
+ monitor.setState(Monitor.CONNECTING);
+ krs=new KLALBRemoteSocket(hp.connectSocket(3000),monitor);
//System.out.println("open");
controller.addRemoteSocket(krs);
- online=true;
+ monitor.resetCoolingTime();
+ monitor.setState(Monitor.ONLINE);
krs.waitforlose();
//System.out.println("close");
- online=false;
} catch (UnknownHostException e) {
} catch (IOException e) {
//e.printStackTrace();
+ }finally {
+ monitor.setState(Monitor.OFFLINE);
+ try {
+ if(krs!=null)
+ krs.close();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
}
if(!flag) {
break;
}
try {
synchronized (lock) {
- lock.wait(10000);
+ lock.wait(monitor.getCoolingTime());
}
+ monitor.incCoolingTime();
} catch (InterruptedException e) {
e.printStackTrace();
}
diff --git a/src/org/kne/cloud/network/klalb/Monitor.java b/src/org/kne/cloud/network/klalb/Monitor.java
index cf2eda3..7a027c5 100644
--- a/src/org/kne/cloud/network/klalb/Monitor.java
+++ b/src/org/kne/cloud/network/klalb/Monitor.java
@@ -1,9 +1,35 @@
package org.kne.cloud.network.klalb;
+import java.util.Timer;
+import java.util.TimerTask;
import java.util.concurrent.atomic.AtomicLong;
+import java.util.function.Consumer;
+
import static org.kne.cloud.network.klalb.KLALBUtils.*;
public class Monitor {
+ private static Timer t=new Timer("带宽测量线程",true);
+
+
+ private TimerTask ptt=new TimerTask() {
+
+ @Override
+ public void run() {
+ runMonitor();
+ }
+ };
+ protected void runMonitor() {
+ updateSpeed();
+
+ }
+ @Override
+ protected void finalize() throws Throwable {
+ ptt.cancel();
+ }
+ public Monitor() {
+
+ t.scheduleAtFixedRate(ptt, 1000, 1000);
+ }
private AtomicLong inTraffic=new AtomicLong();
private AtomicLong outTraffic=new AtomicLong();
@@ -19,21 +45,34 @@ public class Monitor {
private volatile long latency ;
private volatile long jitter ;
+
+ private volatile int state=0;
+ public static final int OFFLINE=0;
+ public static final int CONNECTING=1;
+ public static final int ONLINE=2;
+ private volatile long coolingTime;
+ {
+ resetCoolingTime();
+ }
+
+ public int getState() {
+ return state;
+ }
+ public void setState(int state) {
+ this.state = state;
+ }
public long getInSpeedMax() {
return inSpeedMax;
}
public void updateLatency(long newlatency) {
- long vj=Math.abs( latency-newlatency);;
- if(jitter>1;
+ //}
}
public long getOutSpeedMax() {
return outSpeedMax;
@@ -67,40 +106,98 @@ public class Monitor {
return outSpeed;
}
+ private long updatetime=System.nanoTime();
public void updateSpeed() {
+ long d=System.nanoTime()-updatetime;
long i=inTraffic.get()-inTrafficOld.get();
long o=outTraffic.get()-outTrafficOld.get();
inTrafficOld.set(inTraffic.get());
outTrafficOld.set(outTraffic.get());
-
- inSpeed= i;
- outSpeed= o;
+
+ inSpeed= i*1000000000/d;
+ outSpeed= o*1000000000/d;
+ updatetime=System.nanoTime();
+ //System.out.println(d+" "+o+" "+i);
if(inSpeed>inSpeedMax) {
inSpeedMax=inSpeed;
}else {
- inSpeedMax=(inSpeedMax*999+inSpeed)/1000;
+ inSpeedMax=(inSpeedMax*9999+inSpeed)/10000;
}
if(outSpeed>outSpeedMax) {
outSpeedMax=outSpeed;
}else {
- outSpeedMax=(outSpeedMax*999+outSpeed)/1000;
+ outSpeedMax=(outSpeedMax*9999+outSpeed)/10000;
+ }
+
+ if(changeListener!=null) {
+ changeListener.accept(this);
}
}
+ private ConsumerchangeListener;
+ public Consumer getChangeListener() {
+ return changeListener;
+ }
+ public void setChangeListener(Consumer changeListener) {
+ this.changeListener = changeListener;
+ }
public long getLatency() {
return latency;
}
public String toString() {
StringBuilder sb=new StringBuilder();
+ 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");
return sb.toString();
}
public String toString2() {
StringBuilder sb=new StringBuilder();
- sb.append(bytesUnit(outTraffic.get())).append("\t").append(bytesUnit(inTraffic.get())).append("\t").append(bytesUnit(outSpeed)).append("/s\t").append(bytesUnit(inSpeed)).append("/s\t").append(latency/1000000).append("ms").append("\t").append(jitter/1000000).append("ms");
- return sb.toString();
+ 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");
+
+ if(state==OFFLINE) {
+ sb.append((coolingTime-(System.currentTimeMillis()-mls))/1000L);
+ sb.append('s');
+ }else {
+ sb.append('/');
+ }
+
+ sb.append('\n');
+ sb.append(bytesUnit(outSpeedMax));
+
+ sb.append('\n');
+ sb.append(bytesUnit(inSpeedMax));
+
+ return sb.toString();
+ }
+ private String getDsc() {
+ switch(state) {
+ case OFFLINE:
+ return "○离线";
+ case CONNECTING:
+ return "◐连接中";
+ case ONLINE:
+ return "●在线";
+ }
+ return null;
+ }
+ private volatile long mls;
+ public long getCoolingTime() {
+ mls=System.currentTimeMillis();
+ return coolingTime;
+ }
+ public void incCoolingTime() {
+ coolingTime<<=1;
+ if(coolingTime>300000) {
+ coolingTime=300000;
+ }
+ }
+ public void resetCoolingTime() {
+ coolingTime=10000;
}
}
diff --git a/src/org/kne/cloud/network/klalb/MonitoredSocket.java b/src/org/kne/cloud/network/klalb/MonitoredSocket.java
index a19e8d8..8cf78f4 100644
--- a/src/org/kne/cloud/network/klalb/MonitoredSocket.java
+++ b/src/org/kne/cloud/network/klalb/MonitoredSocket.java
@@ -17,27 +17,20 @@ public class MonitoredSocket extends FilterSocket {
- private static Timer t=new Timer("带宽测量线程",true);
private Monitor monitor;
- private TimerTask ptt=new TimerTask() {
-
- @Override
- public void run() {
- runMonitor();
- }
- };
+
+ @Override
+ public synchronized void close() throws IOException {
+ super.close();
+ }
public MonitoredSocket(Socket socket) {
this(socket,new Monitor());
}
public MonitoredSocket(Socket socket,Monitor monitor) {
super(socket);
- t.scheduleAtFixedRate(ptt, 1000, 1000);
this.monitor=monitor;
}
- protected void runMonitor() {
- monitor.updateSpeed();
-
- }
+
diff --git a/src/org/kne/cloud/network/klalb/SendTask.java b/src/org/kne/cloud/network/klalb/SendTask.java
index 5d078e6..e738a96 100644
--- a/src/org/kne/cloud/network/klalb/SendTask.java
+++ b/src/org/kne/cloud/network/klalb/SendTask.java
@@ -56,8 +56,10 @@ public class SendTask {
+
public void check(int number) throws IOException {
- long limit= (1<{
KLALBVirtualSocket kvs=null;
try {
+ s.setTcpNoDelay(true);
kvs=new KLALBVirtualSocket(kc, kr.getRemoteVaddr(), 23333);
- new SocketBridge(s, kvs).run();
+ 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 {
@@ -56,7 +65,7 @@ public static void main(String[] args) throws UnknownHostException, IOException
String s=scn.nextLine();
switch(s) {
case "state":
- System.out.println("状态\t上传流量\t下载流量\t上传速度\t下载速度\t延迟\t抖动");
+ 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();
diff --git a/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java
index 8a2544b..8027082 100644
--- a/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java
+++ b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java
@@ -119,7 +119,14 @@ public class SimpleKLALBServer {
Socket s=null;
try {
s=mpsa.connectSocket();
- new SocketBridge(soc, s).run();
+ 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();
diff --git a/src/org/kne/cloud/network/mport/SocketBridge.java b/src/org/kne/cloud/network/mport/SocketBridge.java
index 3aac927..e4cedc2 100644
--- a/src/org/kne/cloud/network/mport/SocketBridge.java
+++ b/src/org/kne/cloud/network/mport/SocketBridge.java
@@ -1,25 +1,35 @@
package org.kne.cloud.network.mport;
import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
import java.net.*;
import org.kne.io.Task;
public class SocketBridge extends Task{
private Socket a;
private Socket b;
- public SocketBridge(Socket a, Socket b) {
+ public StreamBridge getSab() {
+ return sab;
+ }
+ public StreamBridge getSba() {
+ return sba;
+ }
+ 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());
}
@Override
protected void runTask() {
try {
- StreamBridge sba=new StreamBridge(a.getInputStream(), b.getOutputStream());
- StreamBridge sbb=new StreamBridge(b.getInputStream(), a.getOutputStream());
- sba.runAtNewThread("SocketBridge A->B thread");
- sbb.runAtNewThread("SocketBridge B->A thread");
+ sab.runAtNewThread("SocketBridge A->B thread");
+ sba.runAtNewThread("SocketBridge B->A thread");
+ sab.waitfortask();
sba.waitfortask();
- sbb.waitfortask();
} catch (Exception e) {
e.printStackTrace();
}finally {
@@ -37,5 +47,23 @@ public class SocketBridge extends Task{
}
}
}
+ public Socket getA() {
+ return a;
+ }
+ public Socket getB() {
+ return b;
+ }
+ protected OutputStream getBOUT() throws IOException {
+ return b.getOutputStream();
+ }
+ protected InputStream getBIN() throws IOException {
+ return b.getInputStream();
+ }
+ protected OutputStream getAOUT() throws IOException {
+ return a.getOutputStream();
+ }
+ protected InputStream getAIN() throws IOException {
+ return a.getInputStream();
+ }
}
diff --git a/src/org/kne/cloud/network/mport/StreamBridge.java b/src/org/kne/cloud/network/mport/StreamBridge.java
index c049ef0..73f1c7c 100644
--- a/src/org/kne/cloud/network/mport/StreamBridge.java
+++ b/src/org/kne/cloud/network/mport/StreamBridge.java
@@ -9,6 +9,7 @@ public class StreamBridge extends Task{
private InputStream in;
private OutputStream out;
private int blocksize=65535;
+ private long delay=0;
public StreamBridge(InputStream in, OutputStream out) {
super();
Objects.requireNonNull(in);
@@ -25,6 +26,11 @@ public class StreamBridge extends Task{
while ((v=in.read(b))!=-1) {
out.write(b,0,v);
out.flush();
+ try {
+ Thread.sleep(delay);
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
}
} catch (IOException e) {
// TODO 自动生成的 catch 块
@@ -51,6 +57,14 @@ public class StreamBridge extends Task{
public void setBlocksize(int blocksize) {
this.blocksize = blocksize;
}
+ public long getDelay() {
+ return delay;
+ }
+
+ public void setDelay(long delay) {
+ this.delay = delay;
+ }
+
public void runAtNewThread() {
runAtNewThread("StreamBridge thread");
}
diff --git a/src/org/kne/cloud/network/mport/ThreadTool.java b/src/org/kne/cloud/network/mport/ThreadTool.java
index a02920d..2982542 100644
--- a/src/org/kne/cloud/network/mport/ThreadTool.java
+++ b/src/org/kne/cloud/network/mport/ThreadTool.java
@@ -21,7 +21,7 @@ public class ThreadTool {
}catch(Throwable e) {
//e.printStackTrace();
if(first) {
- System.out.println("请使用java19以上版本以提高性能!");
+ System.out.println("请使用java19以上版本并开启--enable-preview选项以提高性能!");
first=false;
}
return new Thread(r, name);
diff --git a/src/org/kne/cloud/network/mport/VirtualServerSocket.java b/src/org/kne/cloud/network/mport/VirtualServerSocket.java
index 0d2bf88..0fab1ce 100644
--- a/src/org/kne/cloud/network/mport/VirtualServerSocket.java
+++ b/src/org/kne/cloud/network/mport/VirtualServerSocket.java
@@ -14,8 +14,8 @@ import javax.sql.rowset.RowSetMetaDataImpl;
public abstract class VirtualServerSocket extends ServerSocket {
public VirtualServerSocket(SocketImpl si) throws IOException {
- super();
- try {
+ super(si);
+ /*try {
Class c=ServerSocket.class;
Field con =c.getDeclaredField("impl");
con.setAccessible(true);
@@ -42,7 +42,7 @@ public abstract class VirtualServerSocket extends ServerSocket {
} catch (InvocationTargetException e) {
// TODO 自动生成的 catch 块
e.printStackTrace();
- }
+ }*/
}
@Override
diff --git a/src/org/kne/cloud/network/mport/VirtualSocket.java b/src/org/kne/cloud/network/mport/VirtualSocket.java
index ba5c088..25ef25c 100644
--- a/src/org/kne/cloud/network/mport/VirtualSocket.java
+++ b/src/org/kne/cloud/network/mport/VirtualSocket.java
@@ -13,8 +13,21 @@ import java.nio.channels.SocketChannel;
public abstract class VirtualSocket extends Socket {
- public VirtualSocket(SocketImpl si) throws SocketException {
+ private VirtualSocketImpl si;
+
+ public VirtualSocket(VirtualSocketImpl si) throws SocketException {
super(si);
+ this.si=si;
+ }
+
+ @Override
+ public InputStream getInputStream() throws IOException {
+ return si.getInputStream();
+ }
+
+ @Override
+ public OutputStream getOutputStream() throws IOException {
+ return si.getOutputStream();
}
diff --git a/src/org/kne/cloud/network/mport/VirtualSocketImpl.java b/src/org/kne/cloud/network/mport/VirtualSocketImpl.java
new file mode 100644
index 0000000..44a0097
--- /dev/null
+++ b/src/org/kne/cloud/network/mport/VirtualSocketImpl.java
@@ -0,0 +1,19 @@
+package org.kne.cloud.network.mport;
+
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.net.InetAddress;
+import java.net.SocketAddress;
+import java.net.SocketException;
+import java.net.SocketImpl;
+
+public abstract class VirtualSocketImpl extends SocketImpl {
+
+ @Override
+ public abstract InputStream getInputStream() throws IOException ;
+
+ @Override
+ public abstract OutputStream getOutputStream() throws IOException;
+
+}
diff --git a/src/org/kne/cloud/network/nathole/NatholeTestC.java b/src/org/kne/cloud/network/nathole/NatholeTestC.java
new file mode 100644
index 0000000..c0b2801
--- /dev/null
+++ b/src/org/kne/cloud/network/nathole/NatholeTestC.java
@@ -0,0 +1,32 @@
+package org.kne.cloud.network.nathole;
+
+import java.io.IOException;
+import java.net.BindException;
+import java.net.InetSocketAddress;
+import java.net.Socket;
+import java.net.SocketAddress;
+import java.net.SocketOption;
+import java.net.SocketOptions;
+import java.nio.channels.SocketChannel;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+public class NatholeTestC {
+ public static ExecutorService exc=Executors.newCachedThreadPool();
+ public static void main(String[] args) throws IOException {
+ for (int i1 = 1024; i1 < 65536; i1++) {
+ int i2=i1;
+ exc.execute(()->{
+ try {
+ Socket s=new Socket();
+ System.out.println(i2);
+ s.bind(new InetSocketAddress("0.0.0.0", i2));
+ s.connect(new InetSocketAddress("10.235.16.1", 10400), 100);
+ System.out.println("成功!");
+ System.exit(0);
+ } catch (IOException e) {
+ }
+ });
+ }
+ }
+}
diff --git a/src/org/kne/cloud/network/nathole/NatholeTestS.java b/src/org/kne/cloud/network/nathole/NatholeTestS.java
new file mode 100644
index 0000000..e08c3bf
--- /dev/null
+++ b/src/org/kne/cloud/network/nathole/NatholeTestS.java
@@ -0,0 +1,36 @@
+package org.kne.cloud.network.nathole;
+
+import java.io.IOException;
+import java.net.BindException;
+import java.net.InetSocketAddress;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.nio.channels.SocketChannel;
+
+public class NatholeTestS {
+ public static void main(String[] args) throws IOException {
+ for (int i1 = 1024; i1 < 65536; 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.close();
+ int i2=i1;
+ new Thread(()->{
+ try {
+ ServerSocket ssk=new ServerSocket(i2);
+ while(true) {
+ Socket s=ssk.accept();
+ System.out.println(s);
+ }
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ } ).start();
+ System.out.println(i1);
+ }catch(BindException e) {
+ }
+ }
+ }
+}