diff --git a/src/org/kne/cloud/network/klalb/Consts.java b/src/org/kne/cloud/network/klalb/Consts.java index 25394fa..01daf2c 100644 --- a/src/org/kne/cloud/network/klalb/Consts.java +++ b/src/org/kne/cloud/network/klalb/Consts.java @@ -1,6 +1,6 @@ package org.kne.cloud.network.klalb; public class Consts { - public static final int BLOCKSIZE=65536; - public static final long PINGTIMENS=500000000; + public static final int BLOCKSIZE=32768; + public static final long PINGTIMENS=1000000000; } diff --git a/src/org/kne/cloud/network/klalb/IOThreadManager.java b/src/org/kne/cloud/network/klalb/IOThreadManager.java index 5f6c945..98931fe 100644 --- a/src/org/kne/cloud/network/klalb/IOThreadManager.java +++ b/src/org/kne/cloud/network/klalb/IOThreadManager.java @@ -8,7 +8,7 @@ import java.util.List; import java.util.Vector; public class IOThreadManager { - private KLALBCore klc=new KLALBCore(100); + private KLALBCore klc=new KLALBCore(500); private Listtcps=new Vector<>(); @@ -24,7 +24,7 @@ public class IOThreadManager { while(true) { byte[]b=new byte[Consts.BLOCKSIZE]; int size=local.getDin().read(b); - if(size==-1) + if(size==-1) break; klc.packDataBlock(b,size); } diff --git a/src/org/kne/cloud/network/klalb/KLALBClient.java b/src/org/kne/cloud/network/klalb/KLALBClient.java index 4c580c2..e6a90d0 100644 --- a/src/org/kne/cloud/network/klalb/KLALBClient.java +++ b/src/org/kne/cloud/network/klalb/KLALBClient.java @@ -10,6 +10,7 @@ import java.net.URL; import java.net.UnknownHostException; import java.util.ArrayList; import java.util.List; +import java.util.Scanner; import java.util.UUID; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -17,6 +18,7 @@ import java.util.concurrent.atomic.AtomicInteger; public class KLALBClient { private List tls=new ArrayList<>(); private TCPListener tcpl; + private ServiceElement sel; public KLALBClient(int port) throws IOException { tcpl=new TCPListener(port); tcpl.setCon((s)->{ @@ -43,6 +45,11 @@ public class KLALBClient { TCPConnection tc=new TCPConnection(tll); tc.getDout().writeShort(59649); tc.getDout().write(1); + + tc.getDout().writeUTF(sel.ipport.getIp().getHostAddress()); + tc.getDout().writeInt(sel.ipport.getPort()); + tc.getDout().writeUTF(sel.proc); + tc.getDout().writeUTF(tll.getName()); tc.getDout().writeUTF(tll.getIp()); tc.getDout().writeInt(tll.getPort()); @@ -50,15 +57,15 @@ public class KLALBClient { tc.getDout().writeLong(uid.getLeastSignificantBits()); tc.getDout().flush(); aig.incrementAndGet(); - System.out.println("隧道"+tll+"已连接,可用线路数量:"+ (kcp.getTcps().size()+1)); + //System.out.println("隧道"+tll+"已连接,可用线路数量:"+ (kcp.getTcps().size()+1)); kcp.handleSocket(tc); - + tc.close(); int n=kcp.getTcps().size(); - System.out.println("隧道"+tll+"已断开,可用线路数量:"+ n); + //System.out.println("隧道"+tll+"已断开,可用线路数量:"+ n); if(n<=0) { kcp.closeRemote(); kcp.closeLocal(); - System.out.println("连接已断开"); + //System.out.println("连接已断开"); return; } } catch (IOException e) { @@ -66,7 +73,7 @@ public class KLALBClient { kcp.closeRemote(); kcp.closeLocal(); if(b.get()) { - System.out.println("连接失败"); + System.out.println("隧道连接失败"); b.set(false); } return; @@ -92,8 +99,11 @@ public class KLALBClient { public List getTls() { return tls; } - public static List services=new ArrayList(); + public List services=new ArrayList(); + public List getServices() { + return services; + } public void searchTunnels(String ipport) throws IOException { IPPort u=new IPPort(ipport); TCPConnection tcpc=null; @@ -108,47 +118,63 @@ public class KLALBClient { } int s2=tcpc.getDin().readInt(); for (int i = 0; i < s2; i++) { - tls.add(new Tunnel(tcpc.getDin().readUTF())); + tls.add(Tunnel.newTunnel(tcpc.getDin().readUTF())); } - System.out.println(services); - System.out.println(tls); + //System.out.println(services); + //System.out.println(tls); }finally { if(tcpc!=null) { tcpc.close(); } } } + + public ServiceElement getSel() { + return sel; + } + public void setSel(ServiceElement sel) { + this.sel = sel; + } public static void main(String[] args) throws IOException { + Scanner scn=new Scanner(System.in); KLALBClient kc=new KLALBClient(4568); kc.searchTunnels("153.36.240.12:65529"); - //Tunnel t= new Tunnel("Test1", "127.0.0.1", 4569); - /*Tunnel t= new Tunnel("Test1", "cn-hz-bgp-1.openfrp.top", 65529); - tls.add(t); - Tunnel t2=new Tunnel("Test2","cn-bj-bgp-3.openfrp.top",65529); - tls.add(t2); - Tunnel t3=new Tunnel("Test3","180.76.147.250",65529); - tls.add(t3); - Tunnel t4=new Tunnel("Test4","cn-sx-xa-bgp-1.openfrp.top",65529); - tls.add(t4); - */ -/* List tls=kc.getTls(); - Tunnel t5=new Tunnel("Test5","la.afrps.cn",49966); - tls.add(t5); - Tunnel t6=new Tunnel("Test6","sg.afrps.cn",49966); - tls.add(t6); - Tunnel t7=new Tunnel("Test7","ch.afrps.cn",49966); - tls.add(t7); - Tunnel t8=new Tunnel("Test8","sj.afrps.cn",49966); - tls.add(t8); - Tunnel t9=new Tunnel("Test9","frp.104300.xyz",49965); - tls.add(t9); - - Tunnel t11=new Tunnel("Test11","153.36.240.12",65529); - tls.add(t11); - Tunnel t12=new Tunnel("Test12","frp.freefrps.com",49965); - tls.add(t12);*/ + kc.setSel(kc.getServices().get(0)); kc.open(); + while(true) { + String s=scn.next(); + switch(s) { + case "states": + Tunnel.states(); + break; + } + } } - + //Tunnel t= new Tunnel("Test1", "127.0.0.1", 4569); + /*Tunnel t= new Tunnel("Test1", "cn-hz-bgp-1.openfrp.top", 65529); + tls.add(t); + Tunnel t2=new Tunnel("Test2","cn-bj-bgp-3.openfrp.top",65529); + tls.add(t2); + Tunnel t3=new Tunnel("Test3","180.76.147.250",65529); + tls.add(t3); + Tunnel t4=new Tunnel("Test4","cn-sx-xa-bgp-1.openfrp.top",65529); + tls.add(t4); + */ + /* List tls=kc.getTls(); + Tunnel t5=new Tunnel("Test5","la.afrps.cn",49966); + tls.add(t5); + Tunnel t6=new Tunnel("Test6","sg.afrps.cn",49966); + tls.add(t6); + Tunnel t7=new Tunnel("Test7","ch.afrps.cn",49966); + tls.add(t7); + Tunnel t8=new Tunnel("Test8","sj.afrps.cn",49966); + tls.add(t8); + Tunnel t9=new Tunnel("Test9","frp.104300.xyz",49965); + tls.add(t9); + + Tunnel t11=new Tunnel("Test11","153.36.240.12",65529); + tls.add(t11); + Tunnel t12=new Tunnel("Test12","frp.freefrps.com",49965); + tls.add(t12);*/ } diff --git a/src/org/kne/cloud/network/klalb/KLALBCore.java b/src/org/kne/cloud/network/klalb/KLALBCore.java index 2e22869..c4952d3 100644 --- a/src/org/kne/cloud/network/klalb/KLALBCore.java +++ b/src/org/kne/cloud/network/klalb/KLALBCore.java @@ -31,8 +31,8 @@ public class KLALBCore { private Set inputcache = Collections.synchronizedSet(new HashSet<>()); private List outputcache = Collections.synchronizedList(new ArrayList<>()); - private Map>acks=Collections.synchronizedMap(new WeakHashMap<>()); - + //private Map>acks=Collections.synchronizedMap(new WeakHashMap<>()); + private BlockingQueue ackq=new LinkedBlockingQueue<>(); private volatile boolean inlocal=true; private volatile boolean closeremote = false; @@ -94,8 +94,8 @@ public class KLALBCore { public void sendDataBlock(TCPConnection out) throws IOException, InterruptedException { if (closeremote) throw new InterruptedException(); - BlockingQueue bqk=acks.get(out); - while ((bqk==null||bqk.isEmpty())&& outputcache.isEmpty()) { + //BlockingQueue bqk=acks.get(out); + while (ackq.isEmpty()&& outputcache.isEmpty()) { if (closeremote) Thread.currentThread().interrupt(); Thread.sleep(1); @@ -103,10 +103,8 @@ public class KLALBCore { send0(out, new KLALBBlock(null,0 , 0)); } } - - if(bqk!=null&&bqk.size()>0) { - KLALBBlock klb=bqk.poll(); - if(klb!=null) + KLALBBlock klb=ackq.poll(); + if(klb!=null) { send0(out, klb); }else { KLALBBlock ks = null; @@ -115,21 +113,21 @@ public class KLALBCore { KLALBBlock kd = outputcache.get(i); if (kd.thread == null) { - //kd.time = System.nanoTime(); + kd.time = System.nanoTime(); kd.thread = Thread.currentThread(); ks = kd; } else { if (kd.thread.isAlive()) { - /*long timex = (System.nanoTime() - kd.time) / 1000000; + long timex = (System.nanoTime() - kd.time) / 1000000; if (timex > 10000) { - System.out.println("超时重传:"+kd); + //System.out.println("超时重传:"+kd); kd.time = System.nanoTime(); kd.thread = Thread.currentThread(); ks = kd; - }*/ + } } else { - System.out.println("掉线重传:"+kd); - //kd.time = System.nanoTime(); + //System.out.println("掉线重传:"+kd); + kd.time = System.nanoTime(); kd.thread = Thread.currentThread(); ks = kd; } @@ -160,14 +158,15 @@ public class KLALBCore { if (x.number > 0) { KLALBBlock klb=new KLALBBlock(null, 0, -x.number); - BlockingQueue bq=acks.get(in); + ackq.add(klb); + /*BlockingQueue bq=acks.get(in); if(bq==null) { BlockingQueue bqt=new LinkedBlockingQueue<>(); bqt.add(klb); acks.put(in,bqt); }else { acks.get(in).add(klb); - } + }*/ if (x.number >= inputcount) { diff --git a/src/org/kne/cloud/network/klalb/KLALBSM.java b/src/org/kne/cloud/network/klalb/KLALBSM.java index 8cfe301..a9b9326 100644 --- a/src/org/kne/cloud/network/klalb/KLALBSM.java +++ b/src/org/kne/cloud/network/klalb/KLALBSM.java @@ -3,12 +3,14 @@ package org.kne.cloud.network.klalb; import java.io.File; import java.io.IOException; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.Properties; import java.util.Scanner; public class KLALBSM { - public static List services=new ArrayList(); + public static Map services=new HashMap(); public static List tunnels=new ArrayList(); public static void main(String[] args) throws IOException { @@ -39,10 +41,10 @@ public class KLALBSM { case 0: ServiceElement se=new ServiceElement(s); System.out.println(se); - services.add(se); + services.put(se.proc,se); break; case 1: - Tunnel t=new Tunnel(s); + Tunnel t=Tunnel.newTunnel(s); System.out.println(t); tunnels.add(t); break; @@ -76,10 +78,10 @@ public class KLALBSM { case 0: ServiceElement se=new ServiceElement(s); System.out.println(se); - services.add(se); + services.put(se.proc,se); break; case 1: - Tunnel t=new Tunnel(s); + Tunnel t=Tunnel.newTunnel(s); System.out.println(t); tunnels.add(t); break; diff --git a/src/org/kne/cloud/network/klalb/KLALBServer.java b/src/org/kne/cloud/network/klalb/KLALBServer.java index bb7c5e1..cc15b89 100644 --- a/src/org/kne/cloud/network/klalb/KLALBServer.java +++ b/src/org/kne/cloud/network/klalb/KLALBServer.java @@ -4,14 +4,18 @@ import java.io.DataInputStream; import java.io.IOException; import java.net.ConnectException; import java.net.Socket; +import java.util.Collection; +import java.util.Collections; import java.util.HashMap; +import java.util.Iterator; import java.util.List; +import java.util.Map; import java.util.UUID; import java.util.WeakHashMap; public class KLALBServer { WeakHashMap whm=new WeakHashMap<>(); - public KLALBServer(int port, List services,List tunnels) throws IOException { + public KLALBServer(int port, Map services,List tunnels) throws IOException { TCPListener tcpl=new TCPListener(port); tcpl.setCon((s)->{ try { @@ -25,10 +29,11 @@ public class KLALBServer { int x=din.read(); if(x==0) { synchronized (services) { - int counts=services.size(); - tcc.getDout().writeInt(counts); - for (int i = 0; i < counts; i++) { - tcc.getDout().writeUTF(services.get(i).toString()); + tcc.getDout().writeInt(services.size()); + Collection cll=services.values(); + for (Iterator iterator = cll.iterator(); iterator.hasNext();) { + ServiceElement tunnel = iterator.next(); + tcc.getDout().writeUTF(tunnel.toString()); } } synchronized (tunnels) { @@ -41,10 +46,21 @@ public class KLALBServer { tcc.getDout().flush(); return; } + + String lip=din.readUTF(); + int lport=din.readInt(); + String lname=din.readUTF(); + ServiceElement eas=services.get(lname); + System.out.println(eas); + + if(eas==null||(!eas.ipport.equals(new IPPort(lip, lport)))){ + return; + } + String name = din.readUTF(); String ip = din.readUTF(); int portx=din.readInt(); - Tunnel tll=new Tunnel(name, ip, portx); + Tunnel tll=Tunnel.newTunnel(name, ip, portx); tcc.setTunnel(tll); UUID uid=new UUID(din.readLong(),din.readLong()); @@ -55,7 +71,7 @@ public class KLALBServer { nx=whm.get(uid); }else { nx=new IOThreadManager(); - Socket soc=new Socket("192.168.1.233",8444); + Socket soc=new Socket("192.168.1.233",eas.ipport.getPort()); nx.setLocal(new TCPConnection(null, soc)); nx.startLocal(); whm.put(uid, nx); diff --git a/src/org/kne/cloud/network/klalb/TCPConnection.java b/src/org/kne/cloud/network/klalb/TCPConnection.java index 3e6d235..c5a202c 100644 --- a/src/org/kne/cloud/network/klalb/TCPConnection.java +++ b/src/org/kne/cloud/network/klalb/TCPConnection.java @@ -20,11 +20,11 @@ public class TCPConnection { } private DataOutputStream dout; - private long delay; + private long delay=-1; - public Socket getConnect() { + /*public Socket getConnect() { return connect; - } + }*/ public DataInputStream getDin() { return din; @@ -56,6 +56,8 @@ public class TCPConnection { // TODO 自动生成的 catch 块 e.printStackTrace(); } + if(tunnel!=null&&connect!=null) + tunnel.getCCount().decrementAndGet(); connect = null; } @@ -65,9 +67,12 @@ public class TCPConnection { public void setDelay(long delay) { this.delay = delay; + if(tunnel!=null) { + tunnel.setDelay(delay); + } if(connect!=null) { try { - connect.setSoTimeout((int) (delay/500000)); + connect.setSoTimeout((int) (delay/100000)); } catch (SocketException e) { e.printStackTrace(); } @@ -84,6 +89,9 @@ public class TCPConnection { public TCPConnection(Tunnel t, Socket soc) throws IOException { connect = soc; tunnel = t; + 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)); @@ -108,4 +116,6 @@ public class TCPConnection { } } + + } diff --git a/src/org/kne/cloud/network/klalb/Tunnel.java b/src/org/kne/cloud/network/klalb/Tunnel.java index 55643a7..94d1d03 100644 --- a/src/org/kne/cloud/network/klalb/Tunnel.java +++ b/src/org/kne/cloud/network/klalb/Tunnel.java @@ -7,28 +7,57 @@ import java.io.DataOutputStream; import java.io.IOException; import java.net.Socket; import java.net.UnknownHostException; +import java.util.Iterator; +import java.util.Map; +import java.util.Map.Entry; import java.util.Objects; +import java.util.Set; +import java.util.WeakHashMap; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicInteger; public class Tunnel { private String name; private IPPort ipport; - - public Tunnel(String name, String ip, int port) throws UnknownHostException { + private long delay=-1; + public long getDelay() { + return delay; + } + public void setDelay(long delay) { + this.delay = delay; + } + private static Map tls=new ConcurrentHashMap<>(); + private Tunnel(String name, IPPort ipport) throws UnknownHostException { super(); this.name = name; - this.ipport=new IPPort(ip,port); + this.ipport=ipport; } - public Tunnel(String s) throws UnknownHostException { + public static Tunnel newTunnel(String name, IPPort ipport) throws UnknownHostException { + Tunnel tmp=tls.get(ipport); + if(tmp!=null) { + return tmp; + } + tmp=new Tunnel(name, ipport); + tls.put(ipport, tmp); + return tmp; + } + public static Tunnel newTunnel(String name, String ipport) throws UnknownHostException { + + return newTunnel(name, new IPPort(ipport)); + } + public static Tunnel newTunnel(String name, String ip, int port) throws UnknownHostException { + + return newTunnel(name, new IPPort(ip, port)); + } + + public static Tunnel newTunnel(String s) throws UnknownHostException { String[]t=s.split("\\$"); - name=t[0]; - ipport=new IPPort(t[1]); + return newTunnel(t[0], t[1]); } public String getName() { return name; } - public void setName(String name) { - this.name = name; - } + public String getIp() { return ipport.getIp().getHostAddress(); } @@ -59,7 +88,18 @@ public class Tunnel { Tunnel other = (Tunnel) obj; return Objects.equals(ipport, other.ipport) && Objects.equals(name, other.name); } - + public static void states() { + Set> s=tls.entrySet(); + for (Iterator iterator = s.iterator(); iterator.hasNext();) { + Entry entry = (Entry) iterator.next(); + Tunnel tll=entry.getValue(); + System.out.println(tll.getName()+" "+((tll.getDelay()==-1)?"未知":tll.getDelay()/1000000+"ms")+" 连接数"+tll.ati); + } + } + private AtomicInteger ati=new AtomicInteger(0); + public AtomicInteger getCCount() { + return ati; + }