package org.kne.cloud.network.perf; import java.io.IOException; import java.nio.channels.UnresolvedAddressException; import java.util.Objects; import java.util.concurrent.ConcurrentLinkedQueue; import org.kne.cloud.clock.AdjustedNanoClock; import org.kne.cloud.clock.HighAccuracyClock; import org.kne.cloud.network.MultiProtocolSocketAddress; import org.kne.cloud.network.SpeedLimiter; import org.kne.cloud.network.ThreadTool; import org.kne.cloud.network.klalb.KLALBPacket; import org.kne.cloud.network.klalb.KLALBPacketLink; import org.kne.cloud.network.klalb.KLALBUtils; import org.kne.cloud.network.klalb.PINGPacket; import org.kne.cloud.network.klalb.PONGPacket; import org.kne.cloud.network.monitor.LinkStatus; import org.kne.cloud.network.monitor.SpeedAndTrafficAndDelayMonitorDataImpl; import org.kne.concurrent.ThreadParker; public class Kperf implements Runnable{ private KperfReports reports; private ConcurrentLinkedQueue IsendDequeList = new ConcurrentLinkedQueue(); //private PriorityBlockingQueue sendDequeList = new PriorityBlockingQueue(); private static final boolean debug = true; private static final boolean showpacket = false; private volatile AdjustedNanoClock adjnc = new AdjustedNanoClock(); private MultiProtocolSocketAddress targetAddress; private MultiProtocolSocketAddress bindAddress; private volatile boolean connected=false; private volatile KLALBPacketLink link; private LinkStatus status=new LinkStatus(); private SpeedAndTrafficAndDelayMonitorDataImpl monitor=new SpeedAndTrafficAndDelayMonitorDataImpl(HighAccuracyClock.SYSTEM_CLOCK); private volatile boolean closed = false; private volatile ThreadParker tlock=new ThreadParker(); public Kperf(MultiProtocolSocketAddress targetAddress, MultiProtocolSocketAddress bindAddress) { Objects.requireNonNull(targetAddress); this.targetAddress=targetAddress; this.bindAddress=bindAddress; } public Kperf(MultiProtocolSocketAddress targetAddress) { this(targetAddress,new MultiProtocolSocketAddress("[::0]:0")); } public Kperf(KLALBPacketLink link) { Objects.requireNonNull(link); this.link=link; } public KperfReports getReports() { return reports; } public void setReports(KperfReports reports) { this.reports = reports; } private SpeedLimiter test = new SpeedLimiter(MIN_SPEED,10000000L); private static long MIN_SPEED=64*1024L; KLALBPacket testPacket=null; private long starttime=System.nanoTime(); long totaltime=1000; long usedtime=0; private MonitorUpdater monitorUpdater; private class MonitorUpdater implements Runnable{ @Override public void run() { while (!closed) { status.updateReliability(); monitor.update(); try { Thread.sleep(1); } catch (InterruptedException e) { e.printStackTrace(); } } } } //kperf {KLALB_Stream}[2486:1:f189:7589:f1f7:6f5e:b2eb:1]:4564 //kperf {KLALB_Stream}[2486:1::8888]:4564 @Override public void run() { Thread tcc = ThreadTool.makeVDaemonThreadIfSupport("测速控制线程", () -> { long startTime=System.nanoTime(); try { System.out.println("time\tspd↑\tspd↓\tPPS↑\tPPS↓\tdelay↑\tdelay↓\tjitter↑\tjitter↓"); while(!closed) { System.out.printf("%.2fs\t%s\n", (System.nanoTime()- startTime)/1000000000.0,monitor.toString()); Thread.sleep(1000); } } catch (InterruptedException e) { } }); Thread tc = ThreadTool.makeVDaemonThreadIfSupport("测速控制线程", () -> { final int wtimes=5; totaltime=(20*1000*1+20*1000*4+wtimes*1000*5)*1000000L; starttime=System.nanoTime(); System.out.println("测速开始!"); try { long startTime; System.out.println("测试1/3:延迟测试"); test.setLimitspeed(0); startTime=System.nanoTime(); for (int i = 0; i < 20; i++) { Thread.sleep(1000); updateUsedTime(); } System.out.printf("上传延迟 最小:%.3fms 平均:%.3fms 抖动:%.3fms\n",monitor.getOutDelay()/1000000.0,monitor.getOutDelay()/1000000.0,monitor.getOutJitter()/1000000.0); System.out.printf("下载延迟 最小:%.3fms 平均:%.3fms 抖动:%.3fms\n",monitor.getInDelay()/1000000.0,monitor.getInDelay()/1000000.0,monitor.getInJitter()/1000000.0); if(reports!=null) { reports.uploadDelayFinish(monitor.getOutDelay()); reports.downloadDelayFinish(monitor.getInDelay()); } System.out.println("等待时间"+wtimes+"s"); Thread.sleep(wtimes*1000); startTime=System.nanoTime(); System.out.println("测试2/3:带宽测试"); test.setLimitspeed(1000L*1024L*1024L*1024L); testPacket=new TESTNPacket(65535); for (int i = 0; i < 60; i++) { Thread.sleep(1000); } System.out.printf("上传带宽:%s/s\t",KLALBUtils.convertIBUint(monitor.getOutSpeed())); System.out.printf("下载带宽:%s/s\n",KLALBUtils.convertIBUint(monitor.getInSpeed())); if(reports!=null) { reports.downloadSpeedFinish(monitor.getInSpeed()); reports.uploadSpeedFinish(monitor.getOutSpeed()); } System.out.println("等待时间"+wtimes+"s"); Thread.sleep(wtimes*1000); startTime=System.nanoTime(); System.out.println("测试3/3:包转发率测试"); test.setLimitspeed(1000L*1024L*1024L*1024L); testPacket=new TEST1Packet(); for (int i = 0; i < 20; i++) { Thread.sleep(1000); } if(reports!=null) { reports.uploadPPSFinish(monitor.getOutPPS()); reports.downloadPPSFinish(monitor.getInPPS()); } System.out.println("等待时间"+wtimes+"s"); Thread.sleep(wtimes*1000); } catch (InterruptedException e) { //System.out.println("测速被中断!"); }finally { close(); if(reports!=null) reports.testFinish(); System.out.println("测速结束!"); } }); try { status.setState(LinkStatus.UNSTABLE); if(link==null) link=KLALBUtils.createKLALBPacketLink(bindAddress, targetAddress,false,20000); link.setSoTimeout(20000); Thread t = ThreadTool.makeVDaemonThreadIfSupport("测速发送线程", () -> { try { KLALBPacket kpp = null; while ((!link.isClosed()) && (!closed)) { if (checkPingTimeSleep()) { writePacketToKPL(new PINGPacket(System.nanoTime(),0, 0)); }else { KLALBPacket kpip = IsendDequeList.poll(); if (kpip != null) { writePacketToKPL(kpip); }else { KLALBPacket testp=testPacket; if (testp!=null&&test.checkTransmit(testp.getTotalLength())) { writePacketToKPL(testp); }else { tlock.parkNanos(1000000); } } } } } catch (IOException | UnresolvedAddressException e) { if (debug) e.printStackTrace(); } finally { try { link.close(); } catch (IOException e) { e.printStackTrace(); } } }); t.start(); boolean fst = true; while ((!link.isClosed()) && (!closed)) { KLALBPacket kpp = readPacketFromKPL(); if (kpp == null) { break; } switch (kpp.getType()) { case KLALBPacket.PING: sendIPacket(new PONGPacket(((PINGPacket) kpp).getTime() ,kpp.getRcvtime(), System.nanoTime(),0, 0)); break; case KLALBPacket.PONG: if (fst) { // coll.resetCoolingTime(); status.setState(LinkStatus.UP); connected=true; if(targetAddress!=null) { System.out.println("连接地址"+targetAddress+"成功!"); }else { System.out.println("接受地址"+link.toString()+"连接成功!"); } tcc.start(); tc.start(); fst = false; } PONGPacket png = (PONGPacket) kpp; long ul = png.getRcvtime() - png.getTimepingsnd()-(png.getTimepongsnd()-png.getTimepingrcv()); long DsndDelayFactor0 = png.getTimepingrcv() - png.getTimepingsnd(); long DrcvDelayFactor0 = png.getRcvtime() - png.getTimepongsnd(); long DsndDelayFactor; long DrcvDelayFactor; adjnc.getLock().lock(); try { adjnc.calibrate(png.getTimepongsnd() + ul / 2 - png.getRcvtime(), ul / 2); DsndDelayFactor = DsndDelayFactor0 - adjnc.getDelta(); DrcvDelayFactor = DrcvDelayFactor0 + adjnc.getDelta(); if (DsndDelayFactor < 0) { if (DrcvDelayFactor < 0) { } else { adjnc.setDelta(adjnc.getDelta() + DsndDelayFactor); DsndDelayFactor = DsndDelayFactor0 - adjnc.getDelta(); DrcvDelayFactor = DrcvDelayFactor0 + adjnc.getDelta(); } } else { if (DrcvDelayFactor < 0) { adjnc.setDelta(adjnc.getDelta() - DrcvDelayFactor); DsndDelayFactor = DsndDelayFactor0 - adjnc.getDelta(); DrcvDelayFactor = DrcvDelayFactor0 + adjnc.getDelta(); } } } finally { adjnc.getLock().unlock(); } tlock.unpark(); updateUsedTime() ; if(reports!=null) reports.updateDial(monitor); break; case KLALBPacket.TESTN: break; } } } catch (IOException e) { if(status.getState()==LinkStatus.UNSTABLE) System.out.println("连接目标地址失败!"); if(debug) e.printStackTrace(); }finally { status.setState(LinkStatus.DOWN); if(link!=null) { try { link.close(); } catch (IOException e) { e.printStackTrace(); } } tc.interrupt(); tcc.interrupt(); } } private void updateUsedTime() { usedtime=System.nanoTime()-starttime; //System.out.println(usedtime+" "+totaltime); if(reports!=null) reports.testProcess(usedtime/(double)totaltime); } private volatile long time = System.nanoTime(); private volatile long pingInterval = 10000000L; private volatile long pingIntervalSleep = 200000000L; private boolean checkPingTime() { long cu = System.nanoTime(); if (cu - time > pingInterval) { time = cu; return true; } else { return false; } } private boolean checkPingTimeSleep() { long cu = System.nanoTime(); if (cu - time > pingIntervalSleep) { time = cu; return true; } else { return false; } } private volatile long bwtime = System.nanoTime(); private volatile long bwpingInterval = 10000000L; private volatile long bwpingIntervalSleep = 100000000L; private boolean checkBandwidthReportTime(){ long cu = System.nanoTime(); if (cu - bwtime > bwpingIntervalSleep) { bwtime = cu; return true; } else { return false; } } public void close() { closed = true; status.setState(LinkStatus.DOWN); try { if (link != null) link.close(); } catch (IOException e) { e.printStackTrace(); } } /*protected void sendPacket(KLALBPacket blk) { blk.markJoinqueuetime(); sendDequeList.add(blk); LockSupport.unpark(tlock); }*/ private void sendIPacket(KLALBPacket blk) { IsendDequeList.add(blk); tlock.unpark(); } private KLALBPacket readPacketFromKPL() throws IOException { KLALBPacket packet = link.readKLALBPacket(); if (showpacket) { System.out.println("RX:" + packet); } if (packet != null) { long length=packet.getTotalLength(); //monitor.getDownloadBandwidth().recordPacket(KLALBUtils.createGlobalUUID(), (int)length); monitor.recordDownloadPacket((int)length); } return packet; } private void writePacketToKPL(KLALBPacket packet) throws IOException { long length=packet.getTotalLength(); //monitor.getUploadBandwidth().recordPacket(KLALBUtils.createGlobalUUID(), (int)length); monitor.recordUploadPacket((int) length); link.writeKLALBPacket(packet); link.flush(); if (showpacket) { System.err.println("TX:" + packet); } } public void startPerfing() { monitorUpdater=new MonitorUpdater(); ThreadTool.makeVDaemonThread("线路性能采样线程", monitorUpdater).start(); ThreadTool.makeVDaemonThreadIfSupport("测速接收线程" ,this).start(); } }