diff --git a/src/klalb_en_US.properties b/src/klalb_en_US.properties index e2bf3d3..029dd46 100644 --- a/src/klalb_en_US.properties +++ b/src/klalb_en_US.properties @@ -117,3 +117,8 @@ enablewebapi=Enable Web API weblistenaddr=Web API listen address invaildweblistenaddr=Invalid Web API listen address enabletun=Enable TUN adapter +congestionmonitor=Congestion control monitor +rttbaseline=RTT baseline +rtt=RTT +rttmax=RTT max +rttmin=RTT min \ No newline at end of file diff --git a/src/klalb_zh_CN.properties b/src/klalb_zh_CN.properties index 9c46648..2b815e1 100644 --- a/src/klalb_zh_CN.properties +++ b/src/klalb_zh_CN.properties @@ -116,4 +116,9 @@ webapisettings=Web API 设置 enablewebapi=启用 Web API weblistenaddr=Web API 监听地址:端口 invaildweblistenaddr=无效的 Web API 监听地址:端口 -enabletun=启用TUN虚拟网卡 \ No newline at end of file +enabletun=启用TUN虚拟网卡 +congestionmonitor=拥塞控制监视器 +rttbaseline=往返延迟基线 +rtt=往返延迟 +rttmax=往返延迟上限 +rttmin=往返延迟下限 \ No newline at end of file diff --git a/src/org/kne/cloud/network/congestion/BBRCongestionAlgorithm.java b/src/org/kne/cloud/network/congestion/BBRCongestionAlgorithm.java index 1317607..5a95437 100644 --- a/src/org/kne/cloud/network/congestion/BBRCongestionAlgorithm.java +++ b/src/org/kne/cloud/network/congestion/BBRCongestionAlgorithm.java @@ -136,4 +136,19 @@ public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetConges return "BBR"; } + @Override + public double getUpperDelayBound() { + return MAX_GAIN; + } + + @Override + public double getLowerDelayBound() { + return MIN_GAIN; + } + + @Override + public double getBurstLimit() { + return burstLimit; + } + } diff --git a/src/org/kne/cloud/network/congestion/DCTCPCongestionAlgorithm.java b/src/org/kne/cloud/network/congestion/DCTCPCongestionAlgorithm.java index 936b431..b7ee01e 100644 --- a/src/org/kne/cloud/network/congestion/DCTCPCongestionAlgorithm.java +++ b/src/org/kne/cloud/network/congestion/DCTCPCongestionAlgorithm.java @@ -88,7 +88,7 @@ public class DCTCPCongestionAlgorithm implements CongestionAlgorithm { } }else { if(windowUsed*3L>=congressWindowSize) { - congressWindowSize+=4000; + congressWindowSize+=4096; updateWindowSize(); } } diff --git a/src/org/kne/cloud/network/congestion/DetnetCongestionAlgorithm.java b/src/org/kne/cloud/network/congestion/DetnetCongestionAlgorithm.java index 83d401a..1465e6f 100644 --- a/src/org/kne/cloud/network/congestion/DetnetCongestionAlgorithm.java +++ b/src/org/kne/cloud/network/congestion/DetnetCongestionAlgorithm.java @@ -4,4 +4,7 @@ public interface DetnetCongestionAlgorithm extends CongestionAlgorithm { public void setUpperDelayBound(double upper); public void setLowerDelayBound(double lower); public void setBurstLimit(double limit); + public double getUpperDelayBound( ); + public double getLowerDelayBound(); + public double getBurstLimit(); } diff --git a/src/org/kne/cloud/network/congestion/Vegas2CongestionAlgorithm.java b/src/org/kne/cloud/network/congestion/Vegas2CongestionAlgorithm.java index fc8bc39..87641e8 100644 --- a/src/org/kne/cloud/network/congestion/Vegas2CongestionAlgorithm.java +++ b/src/org/kne/cloud/network/congestion/Vegas2CongestionAlgorithm.java @@ -2,6 +2,8 @@ package org.kne.cloud.network.congestion; import java.util.function.Consumer; +import org.kne.cloud.network.monitor.DelaySampler; + public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCongestionAlgorithm { // 可配置的时延膨胀比上下限 - Vegas2.0核心参数 @@ -17,6 +19,7 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong private Consumer windowControlConsumer; private long MIN_RTTVAR = 30000000L; + private volatile long RTTCurr = 1000000000L; // RTT private volatile long RTTMin = 1000000000L; // BaseRTT private volatile long RTTVar = 1000000000L; private volatile long RTTAvg = 1000000000L; @@ -25,11 +28,6 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong private volatile long RTTCount = 0; private volatile long RTO = 1000000000L; - - private volatile long OWDTotal = 0; - private volatile long OWDCount = 0; - private volatile long OWDAvg2 = 1000000000L; // 当前平滑OWD - private volatile long OWDMin = 1000000000L; // BaseRTT private volatile long window = MIN_WINDOW; private volatile long window2 = MIN_WINDOW; @@ -58,106 +56,14 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong - @Override - public void putAck(long packetSize, long RTTns, long OWDup, boolean ecn) { - if(RTTns<0) { - throw new IllegalArgumentException("RTT:"+ RTTns+" is negative!"); - } - // 处理ECN信号 - 将其视为强烈的拥塞信号 - if (ecn) { - // 当收到ECN时,更激进地减少窗口 - window2 = Math.max(0, window2 * 3 / 4); - if (windowControlConsumer != null) { - windowControlConsumer.accept(window2); - } - } - - // 更新最小RTT(BaseRTT) - if (RTTns <= RTTMin) { - RTTMin = RTTns; - } else { - // 缓慢适应:当网络路径真正变化时,BaseRTT应能缓慢更新 - // 这里使用极慢的衰减因子,只有在持续观测到更低RTT时才快速更新 - RTTMin = (RTTMin * 999 + RTTns) / 1000; - } - - // 更新最小OWD(BaseOWD) - if (OWDup <= OWDMin) { - OWDMin = OWDup; - } else { - // 缓慢适应:当网络路径真正变化时,BaseRTT应能缓慢更新 - // 这里使用极慢的衰减因子,只有在持续观测到更低RTT时才快速更新 - OWDMin = (OWDMin * 999 + OWDup) / 1000; - } - - // 更新平滑RTT估计 - - RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - RTTns)) / 4; - RTTAvg = (RTTAvg * 7 + RTTns) / 8; - - - RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4); - - RTTTotal += RTTns; - RTTCount++; - - OWDTotal += OWDup; - OWDCount++; - - long curr = System.nanoTime(); - // 每个RTT周期调整一次窗口(使用当前估计的RTT) - if (curr - RTTstartTime > Math.max(BIAS, 2*RTTMin)) { - long RTTCountx=RTTCount; - if (RTTCountx > 0) { - long currentRTT = RTTTotal / RTTCountx; - //RTTAvg2 = (RTTAvg2 * 3 + currentRTT) / 4; // 平滑当前RTT - RTTAvg2=currentRTT; - RTTTotal = 0; - RTTCount = 0; - } - - long OWDCountx=OWDCount; - if (OWDCountx > 0) { - long currentOWD = OWDTotal / OWDCountx; - OWDAvg2=currentOWD; - OWDTotal = 0; - OWDCount = 0; - } - - - // ================== Vegas2.0核心逻辑 ================== - // 1. 计算时延膨胀比 - double delayRatio = (double) OWDAvg2 / (double) Math.max(BIAS, OWDMin); - //System.out.println(delayRatio+" "+RTTAvg2/1000000.0+" "+RTTMin/1000000.0); - // 2. 基于比值的窗口调整(取代原来的基于差值的调整) - if (delayRatio > delayUpperBound) { - // 时延过高:减小窗口,减少幅度与超标程度成正比 - window2 -= 4096; - } else if (delayRatio < delayLowerBound) { - // 时延过低:增大窗口,增加幅度与低于目标程度成正比 - window2 += 4096; - } else { - - } - // =================================================== - - - - - updateWindowSize(); - RTTstartTime = curr; - - - - - } - } + private DelaySampler dsp=new DelaySampler(); @Override public void putAck(long packetSize, long RTTns, boolean ecn) { if(RTTns<0) { throw new IllegalArgumentException("RTT:"+ RTTns+" is negative!"); } + this.RTTCurr=RTTns; // 处理ECN信号 - 将其视为强烈的拥塞信号 if (ecn) { // 当收到ECN时,更激进地减少窗口 @@ -167,14 +73,7 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong } } - // 更新最小RTT(BaseRTT) - if (RTTns <= RTTMin) { - RTTMin = RTTns; - } else { - // 缓慢适应:当网络路径真正变化时,BaseRTT应能缓慢更新 - // 这里使用极慢的衰减因子,只有在持续观测到更低RTT时才快速更新 - RTTMin = (RTTMin * 999 + RTTns) / 1000; - } + dsp.recordDelay(RTTns); // 更新平滑RTT估计 @@ -190,6 +89,15 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong long curr = System.nanoTime(); // 每个RTT周期调整一次窗口(使用当前估计的RTT) if (curr - RTTstartTime > Math.max(BIAS, 2*RTTMin)) { + dsp.update(); + long rmin=dsp.getMinDelayNanos(); + if(rmin<=RTTMin) { + RTTMin=rmin; + }else { + RTTMin=Math.min(rmin, RTTMin+1000); + } + + long RTTCountx=RTTCount; if (RTTCountx > 0) { long currentRTT = RTTTotal / RTTCountx; @@ -206,12 +114,12 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong // 2. 基于比值的窗口调整(取代原来的基于差值的调整) if (delayRatio > delayUpperBound) { // 时延过高:减小窗口,减少幅度与超标程度成正比 - window2 -= 4096; + window2 =window2-4096; } else if (delayRatio < delayLowerBound) { // 时延过低:增大窗口,增加幅度与低于目标程度成正比 - window2 += 4096; + window2 =window2+4096; } else { - + //window2=window2+4096 } // =================================================== @@ -271,10 +179,12 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong public long getBaseRTT() { return RTTMin; } - - public long getCurrentRTT() { + public long getAvgRTT() { return RTTAvg2; } + public long getCurrentRTT() { + return RTTCurr; + } @Override public void setUpperDelayBound(double upper) { @@ -295,4 +205,19 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong return "Vegas2"; } + @Override + public double getUpperDelayBound() { + return delayUpperBound; + } + + @Override + public double getLowerDelayBound() { + return delayLowerBound; + } + + @Override + public double getBurstLimit() { + return burstLimit; + } + } \ No newline at end of file diff --git a/src/org/kne/cloud/network/ipv6/AbstractIPv6NetworkLink.java b/src/org/kne/cloud/network/ipv6/AbstractIPv6NetworkLink.java index 165258a..45e6021 100644 --- a/src/org/kne/cloud/network/ipv6/AbstractIPv6NetworkLink.java +++ b/src/org/kne/cloud/network/ipv6/AbstractIPv6NetworkLink.java @@ -63,6 +63,8 @@ public abstract class AbstractIPv6NetworkLink implements IPv6NetworkLink { } private BiConsumer> receiveConsumer; + + private Supplier costSupplier; @Override public void setReceiveConsumer(BiConsumer> receiveConsumer) { this.receiveConsumer = receiveConsumer; @@ -73,6 +75,14 @@ public abstract class AbstractIPv6NetworkLink implements IPv6NetworkLink { return loopback; } + public Supplier getCostSupplier() { + return costSupplier; + } + + public void setCostSupplier(Supplier costSupplier) { + this.costSupplier = costSupplier; + } + @Override public List getRouteItems() { List rlist = new ArrayList<>(); @@ -86,12 +96,12 @@ public abstract class AbstractIPv6NetworkLink implements IPv6NetworkLink { for (Iterator iteratorx = getNeighborsInfo().iterator(); iteratorx.hasNext();) { Neighbor addresses = (Neighbor) iteratorx.next(); RouteItem ri = new RouteItem(new IPv6AddressGroup(addresses.getAddress().getAddress(), 128), - addresses.getAddress().getAddress(), this,RouteItem.SRv6_ENDXSID, 0, 128, null, "D"); + addresses.getAddress().getAddress(), this,RouteItem.SRv6_ENDXSID, 0, 128, costSupplier, "D"); rlist.add(ri); RouteItem ris = new RouteItem(addresses.getLocator(), addresses.getLocator().getAddress(), this, RouteItem.SRv6_ENDSID, 13, 128, - null, "D"); + costSupplier, "D"); rlist.add(ris); } diff --git a/src/org/kne/cloud/network/ipv6/IPv6Address.java b/src/org/kne/cloud/network/ipv6/IPv6Address.java index cbe3cdb..07890ed 100644 --- a/src/org/kne/cloud/network/ipv6/IPv6Address.java +++ b/src/org/kne/cloud/network/ipv6/IPv6Address.java @@ -61,6 +61,10 @@ public final class IPv6Address implements Comparable { ((long)(bytes[15] & 0xFF)); } + public IPv6Address(String string) throws UnknownHostException { + this((Inet6Address) Inet6Address.getByName(string)); + } + /** * 从两个long创建IPv6地址 */ diff --git a/src/org/kne/cloud/network/ipv6/IPv6AddressGroup.java b/src/org/kne/cloud/network/ipv6/IPv6AddressGroup.java index 10a059f..e938e96 100644 --- a/src/org/kne/cloud/network/ipv6/IPv6AddressGroup.java +++ b/src/org/kne/cloud/network/ipv6/IPv6AddressGroup.java @@ -67,7 +67,16 @@ public class IPv6AddressGroup implements Comparable{ prefixLength=in.read(); } public boolean checkMatch(IPv6Address address2) { - return address2.equals(address.maskWith(IPv6Address.createMask(prefixLength))); + IPv6Address cmsk= IPv6Address.createMask(prefixLength); + return address2.maskWith(cmsk).equals(address.maskWith(cmsk)); + } + /** + * 返回裁剪后的路由表地址(主机位清零)。 + * 例如 2001:db8:1::100/64 → 2001:db8:1::/64 + */ + public IPv6AddressGroup toNetworkRoute() { + IPv6Address mask = IPv6Address.createMask(prefixLength); + IPv6Address maskedAddr = address.maskWith(mask); + return new IPv6AddressGroup(maskedAddr, prefixLength); } - } diff --git a/src/org/kne/cloud/network/ipv6/IPv6Packet.java b/src/org/kne/cloud/network/ipv6/IPv6Packet.java index f84638f..6b502ee 100644 --- a/src/org/kne/cloud/network/ipv6/IPv6Packet.java +++ b/src/org/kne/cloud/network/ipv6/IPv6Packet.java @@ -1378,7 +1378,7 @@ public class IPv6Packet extends NetworkPacket { // 重路由计数器 private AtomicInteger rerouteCounter = new AtomicInteger(0); - public AtomicInteger getRerouteCounter() { + public AtomicInteger getRouteCounter() { return rerouteCounter; } diff --git a/src/org/kne/cloud/network/ipv6/RouteItem.java b/src/org/kne/cloud/network/ipv6/RouteItem.java index e3c2379..b7958ca 100644 --- a/src/org/kne/cloud/network/ipv6/RouteItem.java +++ b/src/org/kne/cloud/network/ipv6/RouteItem.java @@ -7,6 +7,7 @@ public class RouteItem implements Comparable,Cloneable{ public static final String DIRECT="Direct"; public static final String SRv6_ENDSID="SRv6 ENDSID"; public static final String SRv6_ENDXSID="SRv6 ENDXSID"; + public static final String KLALB_SRv6 = "KLALB_SRv6"; private IPv6AddressGroup destination; private IPv6Address nexthop; @@ -39,7 +40,7 @@ public class RouteItem implements Comparable,Cloneable{ public RouteItem(IPv6AddressGroup destination, IPv6Address nexthop, IPv6NetworkLink destlink, String proto, int pre, long cost,String flag) { super(); - this.destination = destination; + this.destination = destination.toNetworkRoute(); this.nexthop = nexthop; this.destlink = destlink; this.proto = proto; @@ -50,7 +51,7 @@ public class RouteItem implements Comparable,Cloneable{ public RouteItem(IPv6AddressGroup destination, IPv6Address nexthop, IPv6NetworkLink destlink, String proto, int pre, long cost,Supplier costSupplier,String flag) { super(); - this.destination = destination; + this.destination = destination.toNetworkRoute(); this.nexthop = nexthop; this.destlink = destlink; this.proto = proto; diff --git a/src/org/kne/cloud/network/ipv6/SRHInsertNetworkLink.java b/src/org/kne/cloud/network/ipv6/SRHInsertNetworkLink.java new file mode 100644 index 0000000..97a9cc4 --- /dev/null +++ b/src/org/kne/cloud/network/ipv6/SRHInsertNetworkLink.java @@ -0,0 +1,50 @@ +package org.kne.cloud.network.ipv6; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.List; + +import org.kne.cloud.network.srv6.SRv6Router; + +public class SRHInsertNetworkLink extends AbstractIPv6NetworkLink { + + private IPv6AddressGroup ipgroup; + + public SRHInsertNetworkLink(SRv6Router sRv6Router,IPv6AddressGroup ipgroup) { + super(false); + setSRv6Router(sRv6Router); + this.ipgroup=ipgroup; + } + + @Override + public List getAddressGroups() { + return new ArrayList<>(); + } + + @Override + public List getNeighborsInfo() { + return new ArrayList<>(); + } + + @Override + public void sendPacket(IPv6Packet pack, IPv6Address inet6Address) throws IOException { + //getSrv6Router().insertHeaderAndRoutePacket(null, pack); + System.out.println("insert:"+pack); + } + + @Override + public List getRouteItems() { + return List.of(new RouteItem(ipgroup, null, this, RouteItem.KLALB_SRv6, 30, 0,"G")); + } + + @Override + public String getName() { + return "SRHInsertNetworkLink"; + } + + @Override + public boolean isUp() { + return true; + } + +} diff --git a/src/org/kne/cloud/network/klalb/KLALBController.java b/src/org/kne/cloud/network/klalb/KLALBController.java index b2ae6c7..745238e 100644 --- a/src/org/kne/cloud/network/klalb/KLALBController.java +++ b/src/org/kne/cloud/network/klalb/KLALBController.java @@ -86,7 +86,7 @@ public class KLALBController { private static final boolean debug=false; private static final boolean showpacket = false; - private static final int PREFIX = 128;// 112 + private static final int PREFIX = 128; private static final int DISCOVERY_PORT = 4569; private static List ipmd = new ArrayList<>(); @@ -634,7 +634,8 @@ public class KLALBController { private void loadSRv6ProtocolStack(IPv6AddressGroup selfx, boolean enableVirtualAdapter) { networkInterfaceManager.getInetAddressesExcept().add(selfx.getAddress().toInet6Address()); - srv6Router = new SRv6Router(selfx, clock); + IPv6AddressGroup group=new IPv6AddressGroup(selfx.getAddress(),32); + srv6Router = new SRv6Router(selfx,group, clock); if(configItem!=null) { srv6Router.setPerformanceStrategy(PerformanceStrategy.fromDescription(configItem.getPerformanceStrategy())); srv6Router.setDeviceName(configItem.getDeviceName()); @@ -653,8 +654,7 @@ public class KLALBController { if (enableTUN) { Thread t=new Thread(()->{ try { - IPv6TUNLoopbackNetworkLink tunlink = new IPv6TUNLoopbackNetworkLink(name.trim(), - new IPv6AddressGroup(srv6Router.getLocator().getAddress(), 32), SRv6Router.MTU, dnsAddresses); + IPv6TUNLoopbackNetworkLink tunlink = new IPv6TUNLoopbackNetworkLink(name.trim(),group, SRv6Router.MTU, dnsAddresses); tunlink.setMonitor(datatMonitor); // srv6Router.getLinkTabel().add(tunlink); srv6Router.getinLoopback().setFallbackLink(tunlink); diff --git a/src/org/kne/cloud/network/klalb/KLALBRemoteLink.java b/src/org/kne/cloud/network/klalb/KLALBRemoteLink.java index 5bd3396..e420d85 100644 --- a/src/org/kne/cloud/network/klalb/KLALBRemoteLink.java +++ b/src/org/kne/cloud/network/klalb/KLALBRemoteLink.java @@ -149,6 +149,7 @@ public class KLALBRemoteLink extends AbstractControlledIPv6NetworkLink implement this.addressGroup = address; monitor=new SpeedAndTrafficAndDelayMonitorDataImpl(sysclk); name=kpl.toString(); + setCostSupplier( CostSupplierFactory.owdSupplier(monitor)); } private KLALBRemoteLink(KLALBController controller, MultiProtocolSocketAddress mpa, MultiProtocolSocketAddress bindaddr, IPv6AddressGroup address) { @@ -164,6 +165,7 @@ public class KLALBRemoteLink extends AbstractControlledIPv6NetworkLink implement } else { name=bindaddr.toString() + "→" + mpa.toString(); } + setCostSupplier( CostSupplierFactory.owdSupplier( monitor)); } diff --git a/src/org/kne/cloud/network/klalb/ui/LineMonitorGUI.java b/src/org/kne/cloud/network/klalb/ui/LineMonitorGUI.java index 1cbfe79..f0684ae 100644 --- a/src/org/kne/cloud/network/klalb/ui/LineMonitorGUI.java +++ b/src/org/kne/cloud/network/klalb/ui/LineMonitorGUI.java @@ -18,7 +18,10 @@ import org.jfree.chart.JFreeChart; import org.jfree.data.time.Millisecond; import org.jfree.data.time.TimeSeries; import org.jfree.data.time.TimeSeriesCollection; +import org.kne.cloud.network.congestion.CongestionAlgorithm; +import org.kne.cloud.network.congestion.DetnetCongestionAlgorithm; import org.kne.cloud.network.congestion.SendPacketSlidingWindow; +import org.kne.cloud.network.congestion.Vegas2CongestionAlgorithm; import org.kne.cloud.network.ipv6.IPv6Address; import org.kne.cloud.network.ipv6.IPv6AddressGroup; import org.kne.cloud.network.ipv6.IPv6Packet; @@ -42,12 +45,14 @@ public class LineMonitorGUI extends XFrame { private Timer timer1; private Timer timer2; private Timer timer3; + private Timer timer4; private TimerTask tsk3; - private TimerTask tsk4, tsk5,tsk6; + private TimerTask tsk4, tsk5,tsk6,tsk7; private TimerTask d1; private TimerTask d2; private TimerTask d3; + private TimerTask d4; private TimeSeries spdup = new TimeSeries(UIEnv.getRsb().getString("uploadspeed")); private TimeSeries spddown = new TimeSeries(UIEnv.getRsb().getString("downloadspeed")); @@ -62,13 +67,21 @@ public class LineMonitorGUI extends XFrame { private TimeSeries queueused = new TimeSeries(UIEnv.getRsb().getString("queueused")); private TimeSeries buffermax = new TimeSeries(UIEnv.getRsb().getString("maxbuffer")); + + private TimeSeries rttbaseline = new TimeSeries(UIEnv.getRsb().getString("rttbaseline")); + private TimeSeries rtt = new TimeSeries(UIEnv.getRsb().getString("rtt")); + private TimeSeries rttmin=new TimeSeries(UIEnv.getRsb().getString("rttmin")); + private TimeSeries rttmax=new TimeSeries(UIEnv.getRsb().getString("rttmax")); + private KLALBRemoteLink tr; private JFreeChart jfc; private JFreeChart jfce; private JFreeChart jbfc; + private JFreeChart jbfcv; private JLabel spdp; private JLabel delp; private JLabel sbpdp; + private JLabel sbpdpv; private long timeRange = 5000; private long[] values = new long[] { 1000, 2000, 5000, 10000, 20000, 50000, 100000, 200000, 500000, 1000000 }; private JPanel panel_1; @@ -155,6 +168,22 @@ public class LineMonitorGUI extends XFrame { jp.add(sbpdp); + TimeSeriesCollection tbsc1 = new TimeSeriesCollection(); + tbsc1.addSeries(rttbaseline); + tbsc1.addSeries(rtt); + tbsc1.addSeries(rttmax); + tbsc1.addSeries(rttmin); + jbfcv = ChartFactory.createTimeSeriesChart(UIEnv.getRsb().getString("congestionmonitor"), + UIEnv.getRsb().getString("time") + "(s)", UIEnv.getRsb().getString("delay") + "(ms)", tbsc1); + jbfcv.getXYPlot().setBackgroundPaint(Color.BLACK); + jbfcv.getXYPlot().getRenderer().setSeriesPaint(0, Color.CYAN); + jbfcv.getXYPlot().getRenderer().setSeriesPaint(1, Color.YELLOW); + jbfcv.getXYPlot().getRenderer().setSeriesPaint(2, Color.GREEN); + jbfcv.getXYPlot().getRenderer().setSeriesPaint(3, Color.RED); + changeFont(jbfcv); + sbpdpv = new JLabel(); + sbpdpv.setBorder(new LineBorder(Color.DARK_GRAY)); + jp.add(sbpdpv); @@ -262,6 +291,14 @@ public class LineMonitorGUI extends XFrame { buffermax.setMaximumItemCount((int) timeRange/5); queueused.setMaximumItemCount((int) timeRange/5); } + synchronized (jbfcv) { + + jbfcv.getXYPlot().getDomainAxis().setFixedAutoRange(timeRange); + bufferused.setMaximumItemCount((int) timeRange/5); + buffermax.setMaximumItemCount((int) timeRange/5); + queueused.setMaximumItemCount((int) timeRange/5); + } + windowSelect.setText(timeRange / 1000 + "s"); } @@ -328,6 +365,24 @@ public class LineMonitorGUI extends XFrame { } } } + private void recordCongestionData() { + if (isVisible()) { + Millisecond ms = new Millisecond( + new Date(tr.getMonitor().getClock().getCurrentTimeMillis())); + synchronized (jbfcv) { + CongestionAlgorithm window=tr.getWindow().getAlgorithm(); + if(window !=null&&window instanceof DetnetCongestionAlgorithm) { + Vegas2CongestionAlgorithm deta=(Vegas2CongestionAlgorithm) window; + double baseRTT=deta.getBaseRTT(); + rttbaseline.addOrUpdate(ms, baseRTT / 1000000.0); + rtt.addOrUpdate(ms, deta.getCurrentRTT() / 1000000.0); + rttmax.addOrUpdate(ms, baseRTT*deta.getUpperDelayBound() / 1000000.0); + rttmin.addOrUpdate(ms, baseRTT*deta.getLowerDelayBound() / 1000000.0); + + } + } + } + } @Override public void setVisible(boolean b) { if (b) { @@ -336,6 +391,7 @@ public class LineMonitorGUI extends XFrame { timer1 = new Timer(); timer2 = new Timer(); timer3 = new Timer(); + timer4 = new Timer(); tsk3 = new TimerTask() { @Override @@ -365,6 +421,13 @@ public class LineMonitorGUI extends XFrame { recordBufferData(); } }; + tsk7 = new TimerTask() { + + @Override + public void run() { + recordCongestionData(); + } + }; d1 = new TimerTask() { @Override @@ -406,6 +469,20 @@ public class LineMonitorGUI extends XFrame { } } }; + d4 = new TimerTask() { + + @Override + public void run() { + if (isVisible()) { + synchronized (jbfcv) { + ImageIcon i3 = new ImageIcon( + jbfcv.createBufferedImage(sbpdpv.getWidth(), sbpdpv.getHeight())); + sbpdpv.setIcon(i3); + } + + } + } + }; timer.scheduleAtFixedRate(tsk3, 25, 25); timer1.scheduleAtFixedRate(tsk4, 10, 10); timer1.scheduleAtFixedRate(d1, 25, 25); @@ -413,6 +490,8 @@ public class LineMonitorGUI extends XFrame { timer2.scheduleAtFixedRate(d2, 25, 25); timer3.scheduleAtFixedRate(tsk6, 10, 10); timer3.scheduleAtFixedRate(d3, 25, 25); + timer4.scheduleAtFixedRate(tsk7, 25, 25); + timer4.scheduleAtFixedRate(d4, 25, 25); } } else { @@ -422,12 +501,18 @@ public class LineMonitorGUI extends XFrame { tsk4.cancel(); if (tsk5 != null) tsk5.cancel(); + if (tsk6 != null) + tsk6.cancel(); + if (tsk7 != null) + tsk7.cancel(); if (d1 != null) d1.cancel(); if (d2 != null) d2.cancel(); if (d3 != null) d3.cancel(); + if (d4 != null) + d4.cancel(); if (timer != null) timer.cancel(); if (timer1 != null) @@ -438,6 +523,9 @@ public class LineMonitorGUI extends XFrame { if (timer3 != null) timer3.cancel(); + + if (timer4 != null) + timer4.cancel(); timer = null; } super.setVisible(b); diff --git a/src/org/kne/cloud/network/ntp/NTPContext.java b/src/org/kne/cloud/network/ntp/NTPContext.java index 80c2998..c3997bc 100644 --- a/src/org/kne/cloud/network/ntp/NTPContext.java +++ b/src/org/kne/cloud/network/ntp/NTPContext.java @@ -26,7 +26,7 @@ import org.kne.math.Long128; public class NTPContext implements Closeable, AutoCloseable { private HighAccuracyClock clock; private static final boolean debug = false; - private static final int REQUEST_COUNT = 5; + private static final int REQUEST_COUNT = 10; private Long128 systemFrequencyOffset = NTPTimestamps.nanosToNtp128BitTimeInterval(new Long128(5000)); private Long128 localPrecision = NTPTimestamps.nanosToNtp128BitTimeInterval(new Long128(1000)); diff --git a/src/org/kne/cloud/network/srv6/SRv6Router.java b/src/org/kne/cloud/network/srv6/SRv6Router.java index ea513ca..4975735 100644 --- a/src/org/kne/cloud/network/srv6/SRv6Router.java +++ b/src/org/kne/cloud/network/srv6/SRv6Router.java @@ -44,6 +44,7 @@ import org.kne.cloud.network.ipv6.Neighbor; import org.kne.cloud.network.ipv6.PadNHopByHopTLV; import org.kne.cloud.network.ipv6.PostcardEntry; import org.kne.cloud.network.ipv6.RouteItem; +import org.kne.cloud.network.ipv6.SRHInsertNetworkLink; import org.kne.cloud.network.klalb.KLALBRemoteLink; import org.kne.cloud.network.klalb.KLALBUtils; import org.kne.cloud.network.klalb.PerformanceStrategy; @@ -60,7 +61,7 @@ public class SRv6Router { // public static final int MTU = 16384; public static final int MTU = 16384; // 最大传输单元 - public static final int MAX_REROUTE_COUNT = 2;// 1 // 最大重路由次数 + public static final int MAX_REROUTE_COUNT = 2;// 最大重路由次数 private static final int IPv6_BITS = 128; // IPv6地址位数 @@ -104,7 +105,7 @@ public class SRv6Router { // 环回网络链路 private final LoopbackIPv6NetworkLink inLoopBack; - + private final SRHInsertNetworkLink srhInserter; // 获取路由表 public List getRouteTabel() { return routeTabel; @@ -121,7 +122,7 @@ public class SRv6Router { @Override public void accept(IPv6Packet t) { // HighPerformanceExecutor.defaultExecutor.execute(() -> { - routePacket(null, t, true); // 执行重路由 + routePacket(null, t); // 执行重路由 // }); } }; @@ -142,21 +143,11 @@ public class SRv6Router { if (oamtlv == null) return true; boolean bool = unduplicateSet.add(oamtlv.getUUID()); - // System.out.println(bool); return bool; } - // 默认接收处理器 - /* - * private BiConsumer defaultReceive = new - * BiConsumer() { - * - * @Override public void accept(IPv6NetworkLink link,IPv6Packet t) { // - * HighPerformanceExecutor.defaultExecutor.execute(() -> { if(unduplicate(t)) { - * routePacket(t); // 执行路由 t.putTimePassport("routed"); // 标记路由时间点 } // }); } }; - */ - // SRH接收处理器 - private BiConsumer> srhReceive = new BiConsumer>() { + // 接收处理器 + private BiConsumer> receiveConsumer = new BiConsumer>() { @Override public void accept(IPv6NetworkLink link, Supplier pack) { @@ -255,7 +246,7 @@ public class SRv6Router { nlink.setSRv6Router(this); nlink.addIPv6LinkStateListener(listener); // 添加链路状态监听器 - nlink.setReceiveConsumer(srhReceive); // 设置接收处理器 + nlink.setReceiveConsumer(receiveConsumer); // 设置接收处理器 if (nlink instanceof ControlledIPv6NetworkLink) { ((ControlledIPv6NetworkLink) nlink).setRerouteConsumer(defaultReroute); // 设置重路由处理器 ((ControlledIPv6NetworkLink) nlink).setCongressCondition(congressLock, congressCondition); // 设置拥塞条件 @@ -323,9 +314,9 @@ public class SRv6Router { } // 插入SRH并路由数据包 - public void insertSRHandRoutePacket(IPv6NetworkLink link, IPv6Packet ipp) { - insertHopByHopHeader(ipp); - insertSRHeader(ipp); // 插入段路由头 + public void insertHeaderAndRoutePacket(IPv6NetworkLink link, IPv6Packet ipp) { + insertHopByHopHeader(ipp);//插入逐跳头 + //insertSRHeader(ipp); // 插入段路由头 routePacket(link, ipp); // 路由数据包 } @@ -384,16 +375,13 @@ public class SRv6Router { } private void routePacket(IPv6Packet iPv6Packet) { - routePacket(null, iPv6Packet, false); + routePacket(null, iPv6Packet); } - // 路由数据包(默认非重路由) - private void routePacket(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet) { - routePacket(linkfrom, iPv6Packet, false); - } + // 路由数据包主方法 - private void routePacket(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet, boolean reroute) { + private void routePacket(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet) { StringBuilder dbg = null; if (debug) { dbg = new StringBuilder(); // 调试信息 @@ -406,7 +394,7 @@ public class SRv6Router { List searchResult = getTabelByAddress(dest); if (searchResult!=null&&(!searchResult.isEmpty())) { // 匹配到路由表,执行负载均衡路由 - routingLoadBalance(linkfrom, iPv6Packet, reroute, dbg, searchResult); + routingLoadBalance(linkfrom, iPv6Packet, dbg, searchResult); return; } @@ -441,7 +429,7 @@ public class SRv6Router { // 重新搜索路由表 List searchResult2 = getTabelByAddress(dest); if (searchResult2!=null&&(!searchResult2.isEmpty())) { - routingLoadBalance(linkfrom, iPv6Packet, reroute, dbg, searchResult2); + routingLoadBalance(linkfrom, iPv6Packet, dbg, searchResult2); return; } @@ -508,7 +496,7 @@ public class SRv6Router { private Condition congressCondition = congressLock.newCondition(); // 拥塞条件 // 负载均衡路由 - private boolean routingLoadBalance(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet, boolean reroute, + private boolean routingLoadBalance(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet, StringBuilder dbg, List routes) throws IllegalRawDataException, IOException { if(routes==null||routes.isEmpty()) { @@ -517,7 +505,7 @@ public class SRv6Router { } for (double i = 1; i < 100; i += 0.1) { - if (loadBalance(linkfrom, iPv6Packet, reroute, dbg, routes, i)) { + if (loadBalance(linkfrom, iPv6Packet, dbg, routes, i)) { return true; } @@ -536,7 +524,7 @@ public class SRv6Router { return false; // 路由失败 } - private boolean loadBalance(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet, boolean reroute, StringBuilder dbg, + private boolean loadBalance(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet, StringBuilder dbg, List routes, double cscale) throws IllegalRawDataException, IOException { // 遍历可用路由进行负载均衡 @@ -564,31 +552,39 @@ public class SRv6Router { dbg.append("matched.\n"); } // 找到可用链路,处理数据包 - processPacket(linkfrom, iPv6Packet, tri, reroute); + processPacket(linkfrom, iPv6Packet, tri); return true; } return false; } // 处理数据包转发 - private void processPacket(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet, RouteItem ri, boolean reroute) + private void processPacket(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet, RouteItem ri) throws IllegalRawDataException, IOException { int hop = iPv6Packet.getHopLimit(); if (ri.getDestlink().isLoopBack()) { // 环回链路处理 + IPv6HopByHopHeader iPv6HopByHopHeader=iPv6Packet.getHopByHopHeader(); + if(iPv6Packet.getRouteCounter().get()<=0) { + if (iPv6HopByHopHeader != null&&linkfrom instanceof KLALBRemoteLink) { + processKLALBOAM(iPv6Packet, linkfrom.getAddressGroups().get(0).getAddress(), + ((KLALBRemoteLink) linkfrom).getRemoteVaddr().getAddress(), iPv6HopByHopHeader); + } + } IPv6SegmentRoutingHeader srhh = iPv6Packet.getSegmentRoutingHeader(); if (srhh != null) { - processSRv6Packet(linkfrom, iPv6Packet, ri, srhh, iPv6Packet.getHopByHopHeader(), reroute); // SRv6特殊处理 + processSRv6Packet(linkfrom, iPv6Packet, ri, srhh); // SRv6特殊处理 } else { sendPacketToRouteItem(iPv6Packet, ri); } } else { // 普通链路处理 // 检查重路由计数 - if (iPv6Packet.getRerouteCounter().getAndIncrement() >= MAX_REROUTE_COUNT) { + int count=iPv6Packet.getRouteCounter().getAndIncrement(); + if (count >= MAX_REROUTE_COUNT) { return; // 超过最大重路由次数,丢弃 } - if (!reroute) { + if (count<=0) { hop--; // 减少TTL(非重路由时) } @@ -651,34 +647,16 @@ public class SRv6Router { // 处理SRv6数据包 private void processSRv6Packet(IPv6NetworkLink linkfrom, IPv6Packet iPv6Packet, RouteItem ri, - IPv6SegmentRoutingHeader srhh, IPv6HopByHopHeader iPv6HopByHopHeader, boolean reroute) + IPv6SegmentRoutingHeader srhh) throws IllegalRawDataException, IOException { + if (srhh.getSegmentsLeft() <= 0) { // 所有段已处理完毕 - if (iPv6HopByHopHeader != null) { - int sl = srhh.getSegmentsLeft(); - // IPv6Address prevSID - // =sl+1>=srhh.getAddresses().size()?iPv6Packet.getSourceAddress():srhh.getAddresses().get(sl+1); - if (linkfrom instanceof KLALBRemoteLink) { - processKLALBOAM(iPv6Packet, linkfrom.getAddressGroups().get(0).getAddress(), - ((KLALBRemoteLink) linkfrom).getRemoteVaddr().getAddress(), iPv6HopByHopHeader); - } - } - // System.out.println(iPv6HopByHopHeader); sendPacketToRouteItem(iPv6Packet, ri); } else { - if (reroute) { - // 重路由处理 - } else { + if (iPv6Packet.getRouteCounter().get()<=0) { + // 重路由不处理 int oldSL = srhh.getSegmentsLeft(); - if (iPv6HopByHopHeader != null) { - // IPv6Address prevSID - // =oldSL+1>=srhh.getAddresses().size()?iPv6Packet.getSourceAddress():srhh.getAddresses().get(oldSL+1); - if (linkfrom instanceof KLALBRemoteLink) { - processKLALBOAM(iPv6Packet, linkfrom.getAddressGroups().get(0).getAddress(), - ((KLALBRemoteLink) linkfrom).getRemoteVaddr().getAddress(), iPv6HopByHopHeader); - } - } // 正常SRv6处理:移动到下一个段 int newSL = oldSL - 1; srhh.setSegmentsLeft(newSL); // 更新剩余段数 @@ -687,7 +665,7 @@ public class SRv6Router { // 记录热点地址(用于流量工程) if (iPv6Packet.getPayload().getProtocolNumber() != KLALBRoutingProtocol.DEFAULT_PROTOCOL_NUMBER) klalbRouteProtol.putHotspotAddress(iPv6Packet.getSourceAddress()); - routePacket(linkfrom, iPv6Packet, reroute); // 继续路由 + routePacket(linkfrom, iPv6Packet); // 继续路由 } } @@ -814,7 +792,7 @@ public class SRv6Router { if(ecn) { packet.markCE(); } - insertSRHandRoutePacket(null, packet); // 插入SRH并路由 + insertHeaderAndRoutePacket(null, packet); // 插入SRH并路由 long time=System.nanoTime()-start; backplaneCount(time,ecn); } catch (Exception e) { @@ -827,7 +805,7 @@ public class SRv6Router { public void onReceive(IPv6NetworkLink link,IPv6Packet t){ if (unduplicate(t)) { - insertSRHandRoutePacket(link, t); // 插入SRH并路由数据包 + insertHeaderAndRoutePacket(link, t); // 插入SRH并路由数据包 } } @@ -916,15 +894,17 @@ public class SRv6Router { } // 构造函数 - public SRv6Router(IPv6AddressGroup hostAddress, HighAccuracyClock clock) { + public SRv6Router(IPv6AddressGroup endSID,IPv6AddressGroup networkGroup, HighAccuracyClock clock) { super(); - this.locator = hostAddress; + this.locator = endSID; this.clock = clock; // initWorkerThreads(); // 创建环回链路 this.inLoopBack = new LoopbackIPv6NetworkLink( - List.of(new IPv6AddressGroup(IPv6Address.LOOPBACK, 128), hostAddress), this); + List.of(new IPv6AddressGroup(IPv6Address.LOOPBACK, 128), endSID), this); linkTabel.add(inLoopBack); // 添加到链路表 + this.srhInserter=new SRHInsertNetworkLink(this,networkGroup); + linkTabel.add(srhInserter); backplaneTimer.scheduleAtFixedRate(new TimerTask() { @Override