From 0797da8eb3f513b3b4b1798bf060d46dc83f170e Mon Sep 17 00:00:00 2001 From: Administrator Date: Sun, 18 Dec 2022 09:43:43 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E7=BA=BF=E7=A8=8B=E4=B8=8D?= =?UTF-8?q?=E8=83=BD=E6=AD=A3=E7=A1=AE=E5=85=B3=E9=97=AD=E7=9A=84=E9=97=AE?= =?UTF-8?q?=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../cloud/network/klalb/IOThreadManager.java | 23 ++++++++++++------- .../kne/cloud/network/klalb/KLALBBlock.java | 2 +- .../kne/cloud/network/klalb/KLALBClient.java | 10 +++++--- .../kne/cloud/network/klalb/KLALBCore.java | 11 +++++++-- .../network/klalb/LocalTCPConnection.java | 8 +++++++ .../network/klalb/RemoteTCPConnection.java | 6 +++++ src/org/kne/cloud/network/klalb/Tunnel.java | 2 ++ 7 files changed, 48 insertions(+), 14 deletions(-) diff --git a/src/org/kne/cloud/network/klalb/IOThreadManager.java b/src/org/kne/cloud/network/klalb/IOThreadManager.java index 2eb1d0f..f39a089 100644 --- a/src/org/kne/cloud/network/klalb/IOThreadManager.java +++ b/src/org/kne/cloud/network/klalb/IOThreadManager.java @@ -10,6 +10,7 @@ import java.util.Map; import java.util.UUID; import java.util.Vector; import java.util.WeakHashMap; +import java.util.concurrent.atomic.AtomicBoolean; import org.kne.cloud.network.mport.ServiceElement; import org.kne.cloud.network.mport.ThreadTool; @@ -24,6 +25,7 @@ public class IOThreadManager { public void handleLocal(LocalTCPConnection tc, boolean syn) { klc.getLocaltcps().put(tc.getCuid(), tc); + AtomicBoolean AB=new AtomicBoolean(true); try { if (syn) { KLALBBlock sbk = new KLALBBlock(); @@ -39,7 +41,7 @@ public class IOThreadManager { Thread lt = ThreadTool.makeVThreadIfSupport("数据包计时重发线程", () -> { try { - while (true) { + while (AB.get()) { klc.outputTimer(tc); } } catch (InterruptedException e) { @@ -49,16 +51,16 @@ public class IOThreadManager { Thread ls = ThreadTool.makeVThreadIfSupport("本地发送线程", () -> { try { - while (true) { + while (AB.get()) { KLALBBlock rd=klc.getDataBlock(tc); if(rd.data==null) { break; } tc.unpackBlock(rd); } - } catch (InterruptedException e) { - } catch (IOException e) { + } catch (IOException|InterruptedException e) { + AB.set(false); lt.interrupt(); tc.getOutputcache().clear(); KLALBBlock sbk = new KLALBBlock(); @@ -77,14 +79,15 @@ public class IOThreadManager { }); Thread lr = ThreadTool.makeVThreadIfSupport("本地接收线程", () -> { try { - while (true) { + while (AB.get()) { KLALBBlock ks = tc.packBlock(); klc.putDataBlock(tc, ks); if (ks.data == null) { break; } } - } catch (IOException e) { + } catch (IOException|InterruptedException e) { + AB.set(false); ls.interrupt(); lt.interrupt(); tc.getOutputcache().clear(); @@ -94,7 +97,6 @@ public class IOThreadManager { sbk.command = 3; klc.submitDataBlockNoDelay(sbk); e.printStackTrace(); - } catch (InterruptedException e) { } finally { try { tc.getDin().close(); @@ -104,7 +106,12 @@ public class IOThreadManager { } }); - + tc.setRSTHook((t)->{ + AB.set(false); + ls.interrupt(); + lt.interrupt(); + tc.getOutputcache().clear(); + }); ls.start(); lr.start(); lt.start(); diff --git a/src/org/kne/cloud/network/klalb/KLALBBlock.java b/src/org/kne/cloud/network/klalb/KLALBBlock.java index ef2ec9c..dcf6ffe 100644 --- a/src/org/kne/cloud/network/klalb/KLALBBlock.java +++ b/src/org/kne/cloud/network/klalb/KLALBBlock.java @@ -8,7 +8,7 @@ import java.util.concurrent.atomic.AtomicLong; import org.kne.cloud.network.mport.IPPort; public class KLALBBlock implements Comparable{ - private static AtomicLong sng=new AtomicLong(0); + private static AtomicLong sng=new AtomicLong(1); public long sn;//每个数据包的唯一编号 public UUID cuid;//用于识别数据包的stream ID号 diff --git a/src/org/kne/cloud/network/klalb/KLALBClient.java b/src/org/kne/cloud/network/klalb/KLALBClient.java index 792a72f..efb98d8 100644 --- a/src/org/kne/cloud/network/klalb/KLALBClient.java +++ b/src/org/kne/cloud/network/klalb/KLALBClient.java @@ -40,9 +40,9 @@ public class KLALBClient { try { tcc = new LocalTCPConnection(s); tcc.setServiceElement(sel); - System.out.println("TCP:"+tcc.getCuid()+"已连接"); + //System.out.println("TCP:"+tcc.getCuid()+"已连接"); iom.handleLocal(tcc,true); - System.out.println("TCP:"+tcc.getCuid()+"已关闭"); + //System.out.println("TCP:"+tcc.getCuid()+"已关闭"); } catch (IOException e) { e.printStackTrace(); }finally { @@ -55,9 +55,11 @@ public class KLALBClient { for (int i = 0; i < tls.size(); i++) { Tunnel tll=tls.get(i); ThreadTool.makeVThreadIfSupport("隧道监视线程", ()->{ + int re=0; RemoteTCPConnection tc=null; while(iom.isOpen()) { try { + re++; tc=new RemoteTCPConnection(tll); tc.getDout().writeShort(59649); tc.getDout().write(1); @@ -71,6 +73,7 @@ public class KLALBClient { if(v==-1) { throw new EOFException(); } + re=0; System.out.println(tll+":连接成功"); iom.handleRemote(tc); System.out.println(tll+":连接断开"); @@ -85,7 +88,8 @@ public class KLALBClient { tc.close(); } try { - Thread.sleep(5000); + //System.out.println(re); + Thread.sleep(500*(1<0){ @@ -244,9 +247,13 @@ public class KLALBCore { submitDataBlock(ks); ks.sendtime=System.nanoTime(); tc.getOutputcache().add(ks); + if(tc.getPeerCacheUsed()>cacheblocks) { while(tc.getPeerCacheUsed()>cacheblocks) { Thread.sleep(1); } + }else if(tc.getPeerCacheUsed()>cacheblocks/2){ + Thread.sleep(tc.getPeerCacheUsed()-cacheblocks/2); + } } private volatile long uackt=System.nanoTime(); @@ -256,7 +263,7 @@ 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+1000*i) { + if (timex > 200*(1< hook; public void setPeerCacheUsed(int cacheused) { pcu= cacheused; } public int getPeerCacheUsed() { return pcu; } + public Consumer getRSTHook() { + return hook; + } + public void setRSTHook(Consumer tc) { + this.hook=tc; + } } diff --git a/src/org/kne/cloud/network/klalb/RemoteTCPConnection.java b/src/org/kne/cloud/network/klalb/RemoteTCPConnection.java index f19a000..d0a21f5 100644 --- a/src/org/kne/cloud/network/klalb/RemoteTCPConnection.java +++ b/src/org/kne/cloud/network/klalb/RemoteTCPConnection.java @@ -46,8 +46,10 @@ public class RemoteTCPConnection extends TCPConnection { public void sendBlock(KLALBBlock data) throws IOException { dout.writeLong(data.sn); + if(data.sn!=0) { dout.writeLong(data.cuid.getMostSignificantBits()); dout.writeLong(data.cuid.getLeastSignificantBits()); + } dout.writeLong(data.number); if (data.number > 0) { if (data.data == null) { @@ -80,7 +82,11 @@ public class RemoteTCPConnection extends TCPConnection { public KLALBBlock receiveBlock() throws IOException { KLALBBlock klb = new KLALBBlock(); klb.sn = din.readLong(); + if(klb.sn!=0) { klb.cuid = new UUID(din.readLong(), din.readLong()); + }else { + klb.cuid=new UUID(0,0); + } klb.number = din.readLong(); if (klb.number > 0) { int size = din.readInt(); diff --git a/src/org/kne/cloud/network/klalb/Tunnel.java b/src/org/kne/cloud/network/klalb/Tunnel.java index 452ff41..0025c7b 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 javax.naming.spi.Resolver; + import org.kne.cloud.network.mport.IPPort; public class Tunnel {