package org.kne.cloud.network.klalb; import java.io.IOException; import java.io.StreamCorruptedException; import java.net.SocketTimeoutException; import java.nio.channels.UnresolvedAddressException; import java.security.SecureRandom; import java.util.ArrayList; import java.util.Iterator; import java.util.List; import java.util.Queue; import java.util.UUID; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.LockSupport; import java.util.function.Consumer; import org.jctools.queues.MpscArrayQueue; import org.kne.cloud.clock.HighAccuracyClock; import org.kne.cloud.clock.NTPTimestamps; import org.kne.cloud.clock.ReliabilityBackoffTimeClock; import org.kne.cloud.clock.WatchDogTimer; import org.kne.cloud.network.MultipurposeSocketAddress; import org.kne.cloud.network.ThreadTool; import org.kne.cloud.network.congestion.BBRCongestionAlgorithm; import org.kne.cloud.network.congestion.CongestionAlgorithm; import org.kne.cloud.network.congestion.CongestionAlgorithms; import org.kne.cloud.network.congestion.DetnetCongestionAlgorithm; import org.kne.cloud.network.congestion.SendPacketSlidingWindow; import org.kne.cloud.network.ipv6.ControlledIPv6NetworkLink; import org.kne.cloud.network.ipv6.IPv6Address; import org.kne.cloud.network.ipv6.IPv6HopByHopTLV; import org.kne.cloud.network.ipv6.IPv6Packet; import org.kne.cloud.network.ipv6.IPv6Packet.IPv6HopByHopHeader; import org.kne.cloud.network.ipv6.IPv6AddressGroup; import org.kne.cloud.network.ipv6.KLALBOAMHopByHopTLV; import org.kne.cloud.network.ipv6.KLALBPassportHopByHopTLV; import org.kne.cloud.network.ipv6.Neighbor; import org.kne.cloud.network.ipv6.PostcardEntry; import org.kne.cloud.network.ipv6.RouteItem; import org.kne.cloud.network.monitor.CostSupplierFactory; import org.kne.cloud.network.monitor.DelayMonitorData; import org.kne.cloud.network.monitor.LinkStatus; import org.kne.cloud.network.monitor.SpeedAndTrafficAndDelayMonitorDataImpl; import org.kne.cloud.network.perf.TESTNPacket; import org.kne.cloud.network.te.BandwidthDistributer; import org.kne.concurrent.SpinLock; import org.kne.concurrent.ThreadParker; import org.kne.math.Long128; public class KLALBRemoteLink extends AbstractControlledIPv6NetworkLink implements ControlledIPv6NetworkLink { private static final int HEADER_CALIBRATE = 40; private static final boolean debug = false; private static final boolean showpacket = false; private volatile boolean closed = false; private volatile KLALBController klalbController; public KLALBController getKlalbController() { return klalbController; } private volatile IPv6AddressGroup remoteVaddr; public IPv6AddressGroup getRemoteVaddr() { return remoteVaddr; } private HighAccuracyClock sysclk ; public void waitForRemoteVaddrAvaliable(long timeout) throws SocketTimeoutException { long s = System.nanoTime(); while (remoteVaddr == null) { try { Thread.sleep(1); } catch (InterruptedException e) { e.printStackTrace(); } if (System.nanoTime() - s > timeout * 1000000) throw new SocketTimeoutException(); } } public SpeedAndTrafficAndDelayMonitorDataImpl getMonitor() { return monitor; } public String toString() { return remoteVaddr + "\t" +name+"\n" +status+"\t"+monitor.toString(); } private String name; private LinkStatus status=new LinkStatus(); private SpeedAndTrafficAndDelayMonitorDataImpl monitor; private volatile KLALBPacketLink kplink; private MultipurposeSocketAddress bindAddress; private MultipurposeSocketAddress socketAddress; public void setKplink(KLALBPacketLink kplink) { this.kplink = kplink; } public KLALBPacketLink getKplink() { return kplink; } public MultipurposeSocketAddress getSocketAddress() { return socketAddress; } public MultipurposeSocketAddress getBindAddress() { return bindAddress; } private IPv6AddressGroup addressGroup; private IPv6AddressGroup peerAddress; @Override public List getAddressGroups() { List list = new ArrayList<>(); list.add(addressGroup); return list; } public IPv6AddressGroup getPeerAddress() { return peerAddress; } public KLALBRemoteLink(KLALBController controller, MultipurposeSocketAddress mpa, IPv6AddressGroup address) { this(controller,mpa, null,address); } public KLALBRemoteLink(KLALBController controller,KLALBPacketLink kpl, IPv6AddressGroup address) { this.klalbController=controller; sysclk = klalbController.getClock(); this.kplink = kpl; this.addressGroup = address; monitor=new SpeedAndTrafficAndDelayMonitorDataImpl(sysclk); name=kpl.toString(); } private KLALBRemoteLink(KLALBController controller,MultipurposeSocketAddress mpa, MultipurposeSocketAddress bindaddr, IPv6AddressGroup address) { this.klalbController=controller; sysclk = klalbController.getClock(); this.socketAddress = mpa; this.bindAddress = bindaddr; this.addressGroup = address; monitor=new SpeedAndTrafficAndDelayMonitorDataImpl(sysclk); if (bindaddr == null) { name="→" + mpa.toString(); } else { name=bindaddr.toString() + "→" + mpa.toString(); } } public KLALBRemoteLink(KLALBController controller,MultipurposeSocketAddress mpa, MultipurposeSocketAddress bind) { this(controller,mpa, bind, null); } public KLALBRemoteLink(KLALBController controller,MultipurposeSocketAddress mpa) { this(controller,mpa, (IPv6AddressGroup) null); } public KLALBRemoteLink(KLALBController controller,KLALBPacketLink link) { this(controller,link, (IPv6AddressGroup) null); } private CongestionAlgorithm algorithm; private SendPacketSlidingWindow sendWindow; private class RemoteSender implements Runnable { private MpscArrayQueue sendQueue = new MpscArrayQueue(8192); private WatchDogTimer watchdog = new WatchDogTimer(400000000L); private ThreadParker tlock = new ThreadParker(); private volatile long time = System.nanoTime(); private volatile long pingInterval = 100000000L; private Lock sendLock=new SpinLock(); private boolean checkPingTime() { long cu = System.nanoTime(); if (cu - time > pingInterval) { time = cu; return true; } else { return false; } } public Queue getSendQueue() { return sendQueue; } public WatchDogTimer getWatchdog() { return watchdog; } public long getPingInterval() { return pingInterval; } public void setPingInterval(long pingInterval) { this.pingInterval = pingInterval; } private volatile long flushInterval=1000000L; public long getFlushInterval() { return flushInterval; } public void setFlushInterval(long flushInterval) { this.flushInterval = flushInterval; } private long flushTime=System.nanoTime(); @Override public void run() { try { writeKLALBPacketToKPL(new ADDRREQPacket()); writeKLALBPacketToKPL(new ADDRREQPacket()); if (addressGroup != null) { writeKLALBPacketToKPL(new ADDRPacket(addressGroup)); writeKLALBPacketToKPL(new ADDRPacket(addressGroup)); } writeKLALBPacketToKPL(new VADDRREQPacket()); writeKLALBPacketToKPL(new VADDRREQPacket()); writeKLALBPacketToKPL(new VADDRPacket(klalbController.getSelf())); writeKLALBPacketToKPL(new VADDRPacket(klalbController.getSelf())); kplink.flush(); while ((!kplink.isClosed()) && (!closed)) { if (checkPingTime()) { ping(); if (remoteVaddr == null) sendPacket(new VADDRREQPacket()); if (peerAddress == null) sendPacket(new ADDRREQPacket()); }/*sendLock.lock(); int n; try { n = sendQueue.drain((tkpp)->{ try { kplink.writeKLALBPacket(tkpp); } catch (IOException e) { throw new RuntimeException(e); } },100); }finally { sendLock.unlock(); } if (n>0) { watchdog.feed(); } else { tlock.parkNanos(1000000L); }*/ long curr=System.nanoTime(); if(curr-flushTime>=flushInterval) { flushTime=curr; sendLock.lock(); try { kplink.flush(); watchdog.feed(); }finally { sendLock.unlock(); } } tlock.parkNanos(1000000L); } } catch (StreamCorruptedException sce) { sce.printStackTrace(); } catch (IOException | UnresolvedAddressException e) { if (debug) e.printStackTrace(); } finally { try { kplink.close(); } catch (IOException e) { e.printStackTrace(); } tpk.unpark(); } } public boolean send(KLALBPacket blk) { boolean b=sendQueue.add(blk); if(b) { tlock.unpark(); } return b; } public void sendDirect(KLALBPacket blk) { sendLock.lock(); try { writeKLALBPacketToKPL(blk); } catch (IOException | UnresolvedAddressException e) { if (debug) e.printStackTrace(); try { kplink.close(); } catch (IOException ex) { ex.printStackTrace(); } tpk.unpark(); }finally { sendLock.unlock(); } watchdog.feed(); tlock.unpark(); } public long getQueueUsed() { long used=0; for(KLALBPacket pack:sendQueue) { used+=pack.getTotalLength(); } return used; } } private class RemoteReceiver implements Runnable { private WatchDogTimer watchdog = new WatchDogTimer(400000000L); private long rtt=Long.MAX_VALUE/1024; private long rttMin=Long.MAX_VALUE/1024; AtomicBoolean fst = new AtomicBoolean(true); //ReentrantLock readLock = new ReentrantLock(); @Override public void run() { try { while ((!kplink.isClosed()) && (!closed)) { KLALBPacket kppb = readKLALBPacketFromKPL(); if (kppb == null) { break; } watchdog.feed(); KLALBPacket kpp = kppb; //SRv6Router.getDefaultExecutor().execute(() -> { debugShowPacket("RX", kpp); UUID val = null; long delay = Long.MAX_VALUE / 1024; switch (kpp.getType()) { case KLALBPacket.PING: long pct = System.nanoTime() - kpp.getRcvtime(); Long128 curr = sysclk.getCurrentTimeNanos(); Long128 rcv = curr.subtract(pct); Long128 local1281 = NTPTimestamps.nanosToNtp128BitTimestamp(curr); PONGPacket pong = new PONGPacket(((PINGPacket) kpp).getTime(), NTPTimestamps.ntp128To64(NTPTimestamps.nanosToNtp128BitTimestamp(rcv)).longValue(), NTPTimestamps.ntp128To64(NTPTimestamps.nanosToNtp128BitTimestamp(curr)).longValue(), monitor.getUploadBandwidth().calculateBandwidth(500000000L), monitor.getDownloadBandwidth().calculateBandwidth(500000000L)); sendPacket(pong); /* * long prcv1=pong.getTimepingrcv(); * System.out.println("prcv"+NTPTimestamps.nanosSince1970ToString( * NTPTimestamps.ntp128BitToNanosTimestamp( NTPTimestamps.inferNtp64To128( * prcv1,local1281)))); long posnd1=pong.getTimepongsnd(); * System.out.println("posnd"+NTPTimestamps.nanosSince1970ToString( * NTPTimestamps.ntp128BitToNanosTimestamp( NTPTimestamps.inferNtp64To128( * posnd1,local1281)))); */ break; case KLALBPacket.PONG: PONGPacket png = (PONGPacket) kpp; /* * long ul = png.getRcvtime() - png.getTimepingsnd() - (png.getTimepongsnd() - * png.getTimepingrcv()); */ Long128 local = sysclk.getCurrentTimeNanos(); Long128 local128 = NTPTimestamps.nanosToNtp128BitTimestamp(local); Long128 psnd = NTPTimestamps.ntp128BitToNanosTimestamp( NTPTimestamps.inferNtp64To128(png.getTimepingsnd(), local128)); // System.out.println("psnd"+NTPTimestamps.nanosSince1970ToString( // NTPTimestamps.ntp128BitToNanosTimestamp( NTPTimestamps.inferNtp64To128( // psnd,local128)))); Long128 prcv = NTPTimestamps.ntp128BitToNanosTimestamp( NTPTimestamps.inferNtp64To128(png.getTimepingrcv(), local128)); // System.out.println("prcv"+NTPTimestamps.nanosSince1970ToString( // NTPTimestamps.ntp128BitToNanosTimestamp( NTPTimestamps.inferNtp64To128( // prcv,local128)))); Long128 posnd = NTPTimestamps.ntp128BitToNanosTimestamp( NTPTimestamps.inferNtp64To128(png.getTimepongsnd(), local128)); // System.out.println("posnd"+NTPTimestamps.nanosSince1970ToString( // NTPTimestamps.ntp128BitToNanosTimestamp( NTPTimestamps.inferNtp64To128( // posnd,local128)))); long DsndDelay = prcv.subtract(psnd).longValue(); long DrcvDelay = local.subtract(posnd).longValue(); // System.out.println(DsndDelay+" "+DrcvDelay); rtt = DsndDelay + DrcvDelay; if(rtt= 0 && DrcvDelay >= 0) { delay = DrcvDelay; monitor.recordUploadDelay(DsndDelay); //sysclk.adjustClock(Long128.valueOf( DsndDelay-DrcvDelay), 0.001); } else { long Dfallback = rtt / 2; if (Dfallback >= 0) { delay = Dfallback; monitor.recordUploadDelay(Dfallback); } sysclk.adjustClock(Long128.valueOf( DsndDelay-DrcvDelay), 0.1); } algorithm.setCurrentBandwidth(png.getDownSpeed()); break; case KLALBPacket.BWINF: BWINFPacket bwi = (BWINFPacket) kpp; algorithm.setCurrentBandwidth(bwi.getDownSpeed()); break; case KLALBPacket.ADDRREQ: if (addressGroup != null) sendPacket(new ADDRPacket(addressGroup)); break; case KLALBPacket.ADDR: peerAddress = ((ADDRPacket) kpp).getAddr(); if (addressGroup == null) { byte[] ab = peerAddress.getAddress().toByteArray(); ab[15] = 2; IPv6AddressGroup ardg2 = new IPv6AddressGroup(IPv6Address.valueOf(ab), peerAddress.getPrefixLength()); if (ardg2 != null) addressGroup = ardg2; } KLALBRemoteLink.this.onAddressUpdate(); break; case KLALBPacket.VADDRREQ: sendPacket(new VADDRPacket(klalbController.getSelf())); break; case KLALBPacket.VADDR: IPv6AddressGroup vdr = ((VADDRPacket) kpp).getVaddr(); remoteVaddr = vdr; sendPacket(new VADDRACKPacket()); if (fst.compareAndSet(true, false)) { // coll.resetCoolingTime(); status.setState(LinkStatus.UP); KLALBRemoteLink.this.onOnlineStateUpdate(); } KLALBRemoteLink.this.onLocatorUpdate(); break; case KLALBPacket.VADDRACK: break; case KLALBPacket.TESTN: break; case KLALBPacket.IPV6OVERKLALB: IPv6OverKLALBPacket ivk = (IPv6OverKLALBPacket) kpp; IPv6Packet iv6 = ivk.getIPv6Packet(); if (getSrv6Router() != null) { getSrv6Router().onReceive(KLALBRemoteLink.this, iv6); } KLALBOAMHopByHopTLV ktlv = null; long btimes = 0; IPv6HopByHopHeader hop = iv6.getHopByHopHeader(); if (hop != null) { List ptlv = new ArrayList<>(); List tlvs = hop.getTlvs(); for (IPv6HopByHopTLV tlv : tlvs) { if (tlv instanceof KLALBOAMHopByHopTLV) { ktlv = (KLALBOAMHopByHopTLV) tlv; val = ktlv.getUUID(); btimes = ktlv.getTimestamp(); ptlv.add(btimes); } else if (tlv instanceof KLALBPassportHopByHopTLV) { ptlv.add(btimes + ((KLALBPassportHopByHopTLV) tlv).getTimestampOffsetLong()); } } if (ptlv.size() >= 2) { long delta = ptlv.get(ptlv.size() - 1) - ptlv.get(ptlv.size() - 2); long deltaxnanos = (long) (delta * 1000000000.0 / (1L << 32)); if (deltaxnanos > 0) delay = deltaxnanos; } } break; } long length = kpp.getTotalLength(); UUID pid = val == null ? KLALBUtils.createGlobalUUID() : val; //monitor.getDownloadBandwidth().recordPacket(pid,0,1, delay); if(delay!=Long.MAX_VALUE / 1024) { monitor.recordDownloadDelay(delay); //System.out.println(name+" "+delay); } monitor.recordDownloadPacket((int)length); if (klalbController != null) { //klalbController.getLinkMonitor().getDownloadBandwidth().recordPacket(pid, (int) length,delay); klalbController.getLinkMonitor().recordDownloadPacket((int)length); } //}); } } catch (StreamCorruptedException sce) { sce.printStackTrace(); } catch (IOException | UnresolvedAddressException e) { if (debug) e.printStackTrace(); } catch (Exception e) { e.printStackTrace(); } finally { try { if (kplink != null) kplink.close(); } catch (IOException e) { e.printStackTrace(); } tpk.unpark(); } } public WatchDogTimer getWatchdog() { return watchdog; } public long getRTT() { return rtt; } public long getRTTMin() { return rttMin; } } private ThreadParker tpk=new ThreadParker(); private StateManager stateManager; 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(); } } } } private RemoteSender sender; private RemoteReceiver receiver; private class StateManager implements Runnable { private Thread currThread; @Override public void run() { currThread = Thread.currentThread(); ReliabilityBackoffTimeClock coll = new ReliabilityBackoffTimeClock(); while (true) { try { if (socketAddress == null) { if (kplink.isClosed()) { close(); break; } if (addressGroup == null) { byte[] addr = new byte[16]; SecureRandom scr = new SecureRandom(); scr.nextBytes(addr); addr[0] = (byte) 0x24; addr[1] = (byte) 0x86; addr[2] = 0; addr[3] = 2; addr[14] = 0; addr[15] = 1; IPv6Address i6a = IPv6Address.valueOf(addr); addressGroup = new IPv6AddressGroup(i6a, 112); } } else { kplink = KLALBUtils.createKLALBPacketLink(bindAddress, socketAddress, true); } // startRecord(); kplink.setSoTimeout(3000); algorithm=getCongestionAlgorithm(); sender = new RemoteSender(); receiver = new RemoteReceiver(); sendWindow = new SendPacketSlidingWindow(algorithm, 65535, HEADER_CALIBRATE); sendWindow.setResendConsumer(rerouteConsumer); sendWindow.setClosingConsumer(rerouteConsumer); algorithm.setWindowControlConsumer((sendwindow) -> { sendWindow.setWindowSize(sendwindow); }); Thread t = ThreadTool.makeVDaemonThreadIfSupport("远程发送线程", sender); t.start(); int cores = Runtime.getRuntime().availableProcessors(); for (int i = 0; i < 1; i++) { ThreadTool.makeVDaemonThreadIfSupport("远程接收线程", receiver).start(); } while (!kplink.isClosed()) { if(klalbController!=null&& klalbController.getConfigItem()!=null) { KLALBControllerConfigItem config=klalbController.getConfigItem(); if(algorithm instanceof DetnetCongestionAlgorithm) { double v0=config.getBurstLimit(); if(v0!=0) { ((DetnetCongestionAlgorithm)algorithm).setBurstLimit(v0); } double v1=config.getDelayUpperBound(); if(v1!=0) { ((DetnetCongestionAlgorithm)algorithm).setUpperDelayBound(v1); } double v2=config.getDelayLowerBound(); if(v2!=0) { ((DetnetCongestionAlgorithm)algorithm).setLowerDelayBound(v2); } } sender.setFlushInterval( config.getLinkNagleDelayTime()); } boolean b = sender.getWatchdog().isBarking() || receiver.getWatchdog().isBarking(); boolean isup= (!b)&&(receiver.getRTT()<1000000000L);//&&(receiver.getRTT()(); @Override public List getNeighborsInfo() { List hs = new ArrayList<>(1); if (peerAddress != null) { hs.add(new Neighbor(peerAddress, remoteVaddr, monitor, bandwidthDistributer)); } return hs; } private Consumer rerouteConsumer; @Override public void sendPacket(IPv6Packet pack, IPv6Address inet6Address) throws IOException { // 只有已连接的链路才能传输数据 if (status.getState() == LinkStatus.UP) { IPv6HopByHopHeader hop = pack.getHopByHopHeader(); UUID val = null; if (hop != null) { List l = hop.getTlvs(); for (int i = 0; i < l.size(); i++) { IPv6HopByHopTLV tlv = l.get(i); if (tlv instanceof KLALBOAMHopByHopTLV) { KLALBOAMHopByHopTLV seqs = (KLALBOAMHopByHopTLV) tlv; // 获得数据包的序号 val = seqs.getUUID(); break; } } if (val != null) { //if (!(pack.getPayload() instanceof ICMPv6PostcardPacket)) { // 放入滑动窗口中并记录缓冲区大小占用 sendWindow.put(val, pack); //} } } sendIPv6PacketToLink(pack, val); } else { if (rerouteConsumer != null) { rerouteConsumer.accept(pack); } } } @Override public void setRerouteConsumer(Consumer rerouteConsumer) { this.rerouteConsumer = rerouteConsumer; } private void sendIPv6PacketToLink(IPv6Packet pack, UUID val) { IPv6OverKLALBPacket kipv6 = new IPv6OverKLALBPacket(pack); sendPacketToLink(kipv6, val); } private void sendPacketToLink(KLALBPacket packet, UUID val) { sendPacket(packet, val); } @Override public boolean isCongestion(IPv6Packet iPv6Packet, double scale) { SendPacketSlidingWindow sw = sendWindow; if (sw != null) { return sw.isCongestion( scale); } else { return true; } } @Override public String getName() { return name; } @Override public boolean isUp() { return (!isClosed()) && (status.getState() == LinkStatus.UP) ; } @Override public List getRouteItems() { List rlist = new ArrayList<>(); IPv6AddressGroup adg = addressGroup; /*if (adg != null) rlist.add(new RouteItem(new IPv6AddressGroup(adg.getAddress(), 128), adg.getAddress(), this, "Direct", 0, 0, null, "D", true));*/ 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, "Direct", 0, 128,CostSupplierFactory.expectedDelaySupplier((DelayMonitorData) addresses.getMonitor(),status,algorithm::getRTO), "D", false); rlist.add(ri); if (addresses.getLocator() != null) { RouteItem ris = new RouteItem(addresses.getLocator(), addresses.getLocator().getAddress(), this, "KLALB SRv6", 13, 128,CostSupplierFactory.expectedDelaySupplier((DelayMonitorData) addresses.getMonitor(),status,algorithm::getRTO), "D", false); rlist.add(ris); } } return rlist; } @Override public void setCongressCondition(Lock lock, Condition condition) { if (sendWindow != null) sendWindow.setCongressCondition(lock, condition); } public void onPostcardReceive(PostcardEntry postcard) { long deltaxnanos = (long) (postcard.getTimestampDeltaLong() * 1000000000.0 / (1L << 32)); if (deltaxnanos > 0) { Long128 psnd = NTPTimestamps.ntp128BitToNanosTimestamp( NTPTimestamps.inferNtp64To128(postcard.getTimestampBase(), sysclk.getCurrentTimeNTP128())); monitor.recordUploadDelay(deltaxnanos); } sendWindow.ack(postcard.getPacketUUID(),false,deltaxnanos); } public int getState() { return status.getState(); } public SendPacketSlidingWindow getWindow() { return sendWindow; } public long getQueueUsed() { if(sender==null) { return 0; } return sender.getQueueUsed(); } }