diff --git a/.gitignore b/.gitignore index ae3c172..52b1fb8 100644 --- a/.gitignore +++ b/.gitignore @@ -1 +1,5 @@ /bin/ +/frpc/ +/klalbs4.json +/klalbs.json +/klalbs2.json diff --git a/src/org/kne/cloud/network/klalb/Consts.java b/src/org/kne/cloud/network/klalb/Consts.java index c3a27f4..4593e53 100644 --- a/src/org/kne/cloud/network/klalb/Consts.java +++ b/src/org/kne/cloud/network/klalb/Consts.java @@ -1,9 +1,8 @@ package org.kne.cloud.network.klalb; public class Consts { - public static final int BLOCKSIZE=16384; + public static final int BLOCKSIZE=8192; public static final long PINGTIMENS=1000000000L; - public static final double A = 0.125; public static final int SO_TIMEOUT = 20000; public static final long SN_KEEP = 60000000000L; public static final int MAX_RESEND = 100; diff --git a/src/org/kne/cloud/network/klalb/KLALBBlock.java b/src/org/kne/cloud/network/klalb/KLALBBlock.java index bc06daf..012e3c0 100644 --- a/src/org/kne/cloud/network/klalb/KLALBBlock.java +++ b/src/org/kne/cloud/network/klalb/KLALBBlock.java @@ -23,6 +23,8 @@ public class KLALBBlock implements Comparable{ public transient volatile long sendtime; public transient volatile int resend; + + public long sendtimeForRTT; public KLALBBlock() { diff --git a/src/org/kne/cloud/network/klalb/KLALBClient.java b/src/org/kne/cloud/network/klalb/KLALBClient.java index f4fbc34..447bb09 100644 --- a/src/org/kne/cloud/network/klalb/KLALBClient.java +++ b/src/org/kne/cloud/network/klalb/KLALBClient.java @@ -67,6 +67,7 @@ public class KLALBClient { public void listenProxy() { for (int i = 0; i < pet.size(); i++) { ProxyElement pe=pet.get(i); + if(pe.getIpport().getPort()>0&&pe.getIpport().getPort()<65536) try { open(pe); } catch (IOException e) { diff --git a/src/org/kne/cloud/network/klalb/KLALBCore.java b/src/org/kne/cloud/network/klalb/KLALBCore.java index 539106a..3cf5cda 100644 --- a/src/org/kne/cloud/network/klalb/KLALBCore.java +++ b/src/org/kne/cloud/network/klalb/KLALBCore.java @@ -32,10 +32,15 @@ import java.util.concurrent.Executors; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Predicate; +import java.util.zip.DataFormatException; import org.kne.cloud.network.mport.ThreadTool; public class KLALBCore { + private volatile long RTT=1000000000; + private volatile long RTO=1000000000; + private volatile long DevRTT=0; + private volatile int cacheblocks; private Map localtcps = new ConcurrentHashMap<>(); @@ -104,7 +109,7 @@ public class KLALBCore { pdb.pingtime=System.nanoTime(); tc.sendBlock(pdb); } - if(!tc.getNDSendDeque().isEmpty()||!tc.getSendDeque().isEmpty()) + if(!tc.getNDSendDeque().isEmpty()||!tc.getSendDeque().isEmpty()||!tc.getReSendDeque().isEmpty()) break; synchronized (tc.getSendDeque()) { tc.getSendDeque().wait(10); @@ -112,19 +117,32 @@ public class KLALBCore { } KLALBBlock k0=tc.getNDSendDeque().poll(); if(k0!=null) { - tc.sendBlock(k0,tc.getNDSendDeque().isEmpty()); + if(k0.cuid==0&&k0.number==0&&k0.command==1) { + k0.pingtime+=System.nanoTime()- k0.sendtime; + } + tc.sendBlock(k0); }else { + KLALBBlock k1=tc.getReSendDeque().poll(); + if(k1!=null) { + tc.sendBlock(k1); + }else { KLALBBlock k=tc.getSendDeque().poll(); if(k!=null) { //if(localtcps.containsKey(k.cuid)||k.cuid.equals(ZERO_UUID)) - tc.sendBlock(k,tc.getSendDeque().isEmpty()); + tc.sendBlock(k); } + } } } public void remoteReceive(RemoteTCPConnection tc) throws IOException { KLALBBlock brc=null; for(;;) { - brc=tc.receiveBlock(); + try { + brc=tc.receiveBlock(); + } catch (DataFormatException e) { + throw new StreamCorruptedException("ZIP error"); + } + brc.sendtime=System.nanoTime(); if(brc.number!=0||brc.sn==0) { break; } @@ -156,11 +174,26 @@ public class KLALBCore { pdb.number=0; pdb.command=1; pdb.pingtime=brc.pingtime; + pdb.sendtime=brc.sendtime; tc.getNDSendDeque().offer(pdb); }else if(brc.command==1) { long cur=System.nanoTime(); - long delay=(cur-brc.pingtime)/2; + long delay=cur-brc.pingtime; tc.setDelayAvg(delay); + synchronized (remotetcps) { + Collections.sort(remotetcps, new Comparator() { + @Override + public int compare(RemoteTCPConnection o1, RemoteTCPConnection o2) { + if(o1.getDelay()>o2.getDelay()) { + return 1; + }else if(o1.getDelay() { + List l=ltc.getOutputcache(); + synchronized (l) { + for (Iterator iterator = l.iterator(); iterator.hasNext();) { + KLALBBlock klalbBlock = (KLALBBlock) iterator.next(); + if(klalbBlock.number==nx) { + iterator.remove(); + long SampleRTT=System.nanoTime()-klalbBlock.sendtimeForRTT; + RTT=(RTT*19+SampleRTT)/20; + DevRTT = (3*DevRTT + Math.abs( SampleRTT - RTT))/4; + RTO=RTT+DevRTT; + //System.out.println(RTO/1000000); + } + } + } + /*ltc.getOutputcache().removeIf((b) -> { return b.number == nx; - }); + });*/ ltc.setPeerCacheUsed(brc.cacheused); } } } } - public void submitDataBlock(KLALBBlock kb) throws InterruptedException { + public void submitDataBlock(KLALBBlock ks) throws InterruptedException { + submitDataBlock(ks,1); + } + public void submitDataBlock(KLALBBlock kb,int i) throws InterruptedException { + int ni=Math.min(i, remotetcps.size()); while(true) { synchronized (remotetcps) { - Collections.sort(remotetcps, new Comparator() { - @Override - public int compare(RemoteTCPConnection o1, RemoteTCPConnection o2) { - if(o1.getDelay()>o2.getDelay()) { - return 1; - }else if(o1.getDelay() iterator = remotetcps.iterator(); iterator.hasNext();) { RemoteTCPConnection remoteTCPConnection = (RemoteTCPConnection) iterator.next(); Queue bdq=remoteTCPConnection.getSendDeque(); + if(bdq.size()<1) { + bdq.add(kb); + synchronized (bdq) { + bdq.notifyAll(); + } + ni--; + if(ni<=0) + return; + } + } + } + Thread.sleep(1); + } + + } + public void submitReDataBlock(KLALBBlock kb) throws InterruptedException { + while(true) { + synchronized (remotetcps) { + for (Iterator iterator = remotetcps.iterator(); iterator.hasNext();) { + RemoteTCPConnection remoteTCPConnection = (RemoteTCPConnection) iterator.next(); + Queue bdq=remoteTCPConnection.getReSendDeque(); if(bdq.size()<1) { bdq.add(kb); synchronized (bdq) { @@ -252,20 +311,54 @@ public class KLALBCore { } } } - + public void submitAckBlockNoDelay(KLALBBlock kb,int limit) { + synchronized (remotetcps) { + int coun=0; + for (Iterator iterator = remotetcps.iterator(); iterator.hasNext();) { + if(coun>=limit) { + return; + } + RemoteTCPConnection rmt = (RemoteTCPConnection) iterator.next(); + Queue sq = rmt.getNDSendDeque(); + sq.offer(kb); + synchronized (sq) { + sq.notifyAll(); + } + coun++; + } + } + } public void putDataBlock(LocalTCPConnection tc, KLALBBlock ks) throws InterruptedException { - while (!tc.getOutputcache().isEmpty()&&(tc.getOutputcount()-tc.getOutputcache().get(0).number>cacheblocks)) { + while (true) { + boolean b; + synchronized (ks) { + b=!tc.getOutputcache().isEmpty()&&(tc.getOutputcount()-tc.getOutputcache().get(0).number>cacheblocks); + } + if(!b) { + break; + } Thread.sleep(1); } ks.sendtime=System.nanoTime(); + ks.sendtimeForRTT=ks.sendtime; + //System.out.println(tc.getOutputcache().size()); + if(ks.data==null) { + submitDataBlock(ks); + }else { + if(ks.data.lengthcacheblocks) { - while(tc.getPeerCacheUsed()>cacheblocks) { + int peerCacheUsed = tc.getPeerCacheUsed(); + if(peerCacheUsed>cacheblocks) { + while( tc.getPeerCacheUsed()>cacheblocks) { Thread.sleep(1); } - }else if(tc.getPeerCacheUsed()>cacheblocks/2){ - Thread.sleep(tc.getPeerCacheUsed()-cacheblocks/2); + }else if(peerCacheUsed>cacheblocks/2){ + Thread.sleep(peerCacheUsed-cacheblocks/2); } } @@ -276,15 +369,15 @@ public class KLALBCore { for (int i = 0; i < l.size(); i++) { KLALBBlock block=l.get(i); long timex = (System.nanoTime() - block.sendtime) / 1000000; - if (timex > 200*(1<100)||timex > (RTO/1000000)*(1<=Consts.MAX_RESEND) { l.remove(i); i--; } - System.out.println("超时重传:"+block); + System.out.println("超时重传:"+block+" 位置:"+i); } } } diff --git a/src/org/kne/cloud/network/klalb/RemoteTCPConnection.java b/src/org/kne/cloud/network/klalb/RemoteTCPConnection.java index 7e934e0..074abab 100644 --- a/src/org/kne/cloud/network/klalb/RemoteTCPConnection.java +++ b/src/org/kne/cloud/network/klalb/RemoteTCPConnection.java @@ -10,6 +10,7 @@ import java.io.OutputStream; import java.io.StreamCorruptedException; import java.net.Socket; import java.net.UnknownHostException; +import java.util.Arrays; import java.util.Queue; import java.util.Random; import java.util.UUID; @@ -18,8 +19,11 @@ import java.util.concurrent.BlockingDeque; import java.util.concurrent.BlockingQueue; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.LinkedBlockingDeque; +import java.util.zip.DataFormatException; +import java.util.zip.Deflater; import java.util.zip.GZIPInputStream; import java.util.zip.GZIPOutputStream; +import java.util.zip.Inflater; import org.kne.cloud.network.mport.IPPort; import org.kne.io.AddInputStream; @@ -40,6 +44,8 @@ public class RemoteTCPConnection extends TCPConnection { } protected void initIOs() throws IOException { + connect.setTrafficClass(0x10); + connect.setTcpNoDelay(true); OutputStream outm=null; InputStream inm=null; if(tunnel!=null) { @@ -58,8 +64,8 @@ public class RemoteTCPConnection extends TCPConnection { } - dout = new DataOutputStream(new BufferedOutputStream(outm,Consts.BLOCKSIZE*2+1) ); - din = new DataInputStream(new BufferedInputStream(inm,Consts.BLOCKSIZE*2+1)); + dout = new DataOutputStream(outm); + din = new DataInputStream(inm); } @Override @@ -133,8 +139,17 @@ public class RemoteTCPConnection extends TCPConnection { if (data.data == null) { dout.writeShort(-1); } else { - dout.writeShort(data.data.length); - dout.write(data.data); + Deflater def = new Deflater(Deflater.BEST_COMPRESSION, true); + def.setInput(data.data); + def.finish(); + byte[] b = new byte[(int)(data.data.length * 1.1D) + 64]; + int nsize = def.deflate(b); + def.end(); + this.dout.writeShort(nsize); + this.dout.write(b, 0, nsize); + /*dout.writeShort(data.data.length); + dout.write(data.data);*/ + //System.out.println(Arrays.toString(data.data)); } } else if (data.number == 0) { dout.write(data.command); @@ -165,7 +180,7 @@ public class RemoteTCPConnection extends TCPConnection { //System.out.println("SEND:" + data); } - public KLALBBlock receiveBlock() throws IOException { + public KLALBBlock receiveBlock() throws IOException, DataFormatException { KLALBBlock klb = new KLALBBlock(); klb.cuid = din.readInt(); klb.number = din.readLong(); @@ -176,7 +191,18 @@ public class RemoteTCPConnection extends TCPConnection { } else { byte[] d = new byte[size]; din.readFully(d); - klb.data = d; + Inflater in = new Inflater(true); + in.setInput(d); + byte[] b = new byte[8192]; + int nsize = in.inflate(b); + in.end(); + if (nsize == b.length) { + klb.data = b; + } else { + klb.data = Arrays.copyOf(b, nsize); + } + //System.out.println(Arrays.toString(d)); + //klb.data = d; } } else if (klb.number == 0) { klb.command = din.read(); @@ -222,7 +248,11 @@ public class RemoteTCPConnection extends TCPConnection { public Queue getNDSendDeque() { return NDsendDeque; } - + private Queue resendDeque = new ConcurrentLinkedQueue<>(); + + public Queue getReSendDeque() { + return sendDeque; + } private Queue sendDeque = new ConcurrentLinkedQueue<>(); public Queue getSendDeque() { diff --git a/src/org/kne/cloud/network/klalb/Tunnel.java b/src/org/kne/cloud/network/klalb/Tunnel.java index 548234e..5b11281 100644 --- a/src/org/kne/cloud/network/klalb/Tunnel.java +++ b/src/org/kne/cloud/network/klalb/Tunnel.java @@ -138,7 +138,7 @@ public class Tunnel{ try { - p = Runtime.getRuntime().exec("./frpc/frpc_windows_386.exe",null,f); + p = Runtime.getRuntime().exec("./frpc/frpc_windows_amd64.exe",null,f); BufferedReader brd=new BufferedReader(new InputStreamReader( p.getInputStream())); diff --git a/src/org/kne/debug/TimeDebugger.java b/src/org/kne/debug/TimeDebugger.java new file mode 100644 index 0000000..c9f3a79 --- /dev/null +++ b/src/org/kne/debug/TimeDebugger.java @@ -0,0 +1,25 @@ +package org.kne.debug; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.Hashtable; +import java.util.List; +import java.util.Map; + +public class TimeDebugger { + Map l=new Hashtable(); + long tmp=-1; + public void putTime(String name){ + + long v=System.currentTimeMillis(); + if(tmp==-1){ + l.put(name, 0l); + }else{ + l.put(name, v-tmp); + } + tmp=v; + } + public void print(){ + System.err.println(l+"(ms)"); + } +} diff --git a/src/org/kne/io/Compress.java b/src/org/kne/io/Compress.java new file mode 100644 index 0000000..98703cc --- /dev/null +++ b/src/org/kne/io/Compress.java @@ -0,0 +1,188 @@ +package org.kne.io; + +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.io.DataInputStream; +import java.io.DataOutputStream; +import java.io.EOFException; +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.util.ArrayList; +import java.util.concurrent.Executors; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.zip.DataFormatException; +import java.util.zip.Deflater; +import java.util.zip.Inflater; +import org.kne.debug.TimeDebugger; + +public class Compress { + static int n=0; + public static ThreadPoolExecutor exf=(ThreadPoolExecutor) Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors()*2, new ThreadFactory() { + + @Override + public Thread newThread(Runnable r) { + Thread t=new Thread(r,"Compress #"+(n++)); + t.setDaemon(true); + return t; + } + }); + + + + public static Data compress(Data b) { + + Deflater def=new Deflater(Deflater.BEST_COMPRESSION, true); + def.setInput(b.b,b.off,b.len); + def.finish(); + byte[]bo=new byte[(int) (b.len*1.1)+64]; + int vp=def.deflate(bo,4,bo.length-4); + + def.end(); + int v=b.len; + bo[0]=(byte) ((v >>> 24) & 0xFF); + bo[1]=(byte) ((v >>> 16) & 0xFF); + bo[2]=(byte) ((v >>> 8) & 0xFF); + bo[3]=(byte) ((v >>> 0) & 0xFF); + return new Data(bo,0,vp+4); + } +public static Data uncompress(Data b) throws IOException, DataFormatException { + if(b.len<4){ + throw new EOFException("!!"); + } + int ch1 =b.b[b.off+0]&0xff; + int ch2 = b.b[b.off+1]&0xff; + int ch3 =b. b[b.off+2]&0xff; + int ch4 =b.b[b.off+3]&0xff; + + int zix= ((ch1 << 24) + (ch2 << 16) + (ch3 << 8) + (ch4 << 0)); + + Inflater inf=new Inflater(true); + inf.setInput(b.b, b.off+4, b.len-4); + byte[]d=new byte[zix]; + int l=inf.inflate(d); + inf.end(); + + return new Data(d,0,l); + } +private static class compresstask extends Task{ + volatile Data input,output; + + public compresstask(Data input) { + super(); + this.input = input; + } + + public Data getOutput() { + return output; + } + + @Override + protected void runTask() { + output=compress(input); + } + +} +public static void mulcomp(DataOutputStream dos,byte[]in, TimeDebugger dbg) throws IOException { + ArrayList l=new ArrayList(); + int count =4; + int sze=in.length/count; + int y=in.length%count; + + for (int i = 0; i < count; i++) { + Task t=new compresstask(new Data(in,i*sze,sze)); + exf.execute(t); + l.add(t); + } + Data tmp=null; + if(y>0) { + tmp=compress(new Data(in,count*sze,y)); + } + dbg.putTime("compress"); + dos.writeInt(l.size()+(tmp==null?0:1)); + for (int i = 0; i < l.size(); i++) { + compresstask n=(compresstask) l.get(i); + n.waitfortask(); + Data g=n.getOutput(); + dos.writeInt(g.len); + dos.write(g.b, g.off, g.len); + + } + + + if(tmp!=null) { + dos.writeInt(tmp.len); + dos.write(tmp.b, tmp.off, tmp.len); + } +} +private static class uncompresstask extends Task{ + volatile Data input,output; + + public uncompresstask(Data input) { + super(); + this.input = input; + } + + public Data getOutput() { + return output; + } + + @Override + protected void runTask() { + try { + output=uncompress(input); + } catch (IOException e) { + // TODO �Զ����ɵ� catch �� + e.printStackTrace(); + } catch (DataFormatException e) { + // TODO �Զ����ɵ� catch �� + e.printStackTrace(); + } + } + +} +public static byte[] muluncomp(DataInputStream dis,TimeDebugger dbg) throws IOException, DataFormatException { + ArrayList l=new ArrayList(); + ByteArrayOutputStream bod=new ByteArrayOutputStream(1024); + int bc=dis.readInt(); + for(int x=0;x= 0) { + offset += numRead; + } + Task tsk=new uncompresstask(new Data(input,0,input.length)); + exf.execute(tsk); + l.add(tsk); + + //Data rez=uncompress(new Data(input,0,input.length)); + + } dbg.putTime("read"); + for (int i = 0; i < l.size(); i++) { + uncompresstask uct=(uncompresstask) l.get(i); + uct.waitfortask(); + Data rez=uct.getOutput(); + bod.write(rez.b,rez.off,rez.len); + } + + + //rez.b.length=rez.off; + return bod.toByteArray(); +} + public static void main(String[] args) throws IOException, DataFormatException { + String s="jjjjjjjjjjjjjjjjjjhjhjhm"; + byte[] ib=s.getBytes(); + Data i=new Data(ib,0,ib.length); + Data o=compress(i); + Data d=uncompress(o); + System.out.println(o.len+" "+d.len); + //for (int j = d.off; j < d.off+d.len; j++) { + System.out.println(new String(d.b,d.off,d.len)); + //} + } + +} diff --git a/src/org/kne/io/Data.java b/src/org/kne/io/Data.java new file mode 100644 index 0000000..59f7ad3 --- /dev/null +++ b/src/org/kne/io/Data.java @@ -0,0 +1,13 @@ +package org.kne.io; + +public class Data { + byte[]b; + int off,len; + public Data(byte[] b, int off, int len) { + super(); + this.b = b; + this.off = off; + this.len = len; + } + +} diff --git a/src/org/kne/io/Task.java b/src/org/kne/io/Task.java new file mode 100644 index 0000000..827debb --- /dev/null +++ b/src/org/kne/io/Task.java @@ -0,0 +1,31 @@ +package org.kne.io; + +public abstract class Task implements Runnable{ + private volatile boolean begin=false,end=false; + private volatile Object t=new Object(); + @Override + public void run() { + begin=true; + runTask(); + + end=true; + if(t!=null) { + synchronized(t) { + t.notifyAll(); + } + } + } + public void waitfortask() { + synchronized(t){ + + try { + while(!end) { + t.wait(1); + } + } catch (InterruptedException e) { + + e.printStackTrace(); + }} + } + protected abstract void runTask(); +}