diff --git a/.idea/.gitignore b/.idea/.gitignore new file mode 100644 index 0000000..30cf57e --- /dev/null +++ b/.idea/.gitignore @@ -0,0 +1,10 @@ +# Default ignored files +/shelf/ +/workspace.xml +# Editor-based HTTP Client requests +/httpRequests/ +# Ignored default folder with query files +/queries/ +# Datasource local storage ignored files +/dataSources/ +/dataSources.local.xml diff --git a/.idea/misc.xml b/.idea/misc.xml new file mode 100644 index 0000000..28ee0a9 --- /dev/null +++ b/.idea/misc.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/.idea/modules.xml b/.idea/modules.xml new file mode 100644 index 0000000..79ba3af --- /dev/null +++ b/.idea/modules.xml @@ -0,0 +1,8 @@ + + + + + + + + \ No newline at end of file diff --git a/.idea/vcs.xml b/.idea/vcs.xml new file mode 100644 index 0000000..35eb1dd --- /dev/null +++ b/.idea/vcs.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/KLALB.iml b/KLALB.iml new file mode 100644 index 0000000..717bc29 --- /dev/null +++ b/KLALB.iml @@ -0,0 +1,253 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/lib/jctools-core-4.0.6-javadoc.jar b/lib/jctools-core-4.0.6-javadoc.jar new file mode 100644 index 0000000..d6d818e Binary files /dev/null and b/lib/jctools-core-4.0.6-javadoc.jar differ diff --git a/lib/jctools-core-4.0.6-sources.jar b/lib/jctools-core-4.0.6-sources.jar new file mode 100644 index 0000000..6c98e07 Binary files /dev/null and b/lib/jctools-core-4.0.6-sources.jar differ diff --git a/lib/jctools-core-4.0.6.jar b/lib/jctools-core-4.0.6.jar new file mode 100644 index 0000000..94a8307 Binary files /dev/null and b/lib/jctools-core-4.0.6.jar differ diff --git a/src/org/kne/cloud/network/congestion/BBRCongestionAlgorithmFactory.java b/src/org/kne/cloud/network/congestion/BBRCongestionAlgorithmFactory.java new file mode 100644 index 0000000..4b6f9b8 --- /dev/null +++ b/src/org/kne/cloud/network/congestion/BBRCongestionAlgorithmFactory.java @@ -0,0 +1,20 @@ +package org.kne.cloud.network.congestion; + +public class BBRCongestionAlgorithmFactory implements CongestionAlgorithmFactory { + + @Override + public CongestionAlgorithm create() { + return new BBRCongestionAlgorithm(); + } + + @Override + public String getName() { + return "BBR"; + } + + @Override + public String toString() { + return getName(); + } + +} diff --git a/src/org/kne/cloud/network/congestion/CongestionAlgorithmFactory.java b/src/org/kne/cloud/network/congestion/CongestionAlgorithmFactory.java new file mode 100644 index 0000000..c549000 --- /dev/null +++ b/src/org/kne/cloud/network/congestion/CongestionAlgorithmFactory.java @@ -0,0 +1,7 @@ +package org.kne.cloud.network.congestion; + +public interface CongestionAlgorithmFactory { + public CongestionAlgorithm create(); + + public String getName(); +} diff --git a/src/org/kne/cloud/network/congestion/CongestionAlgorithms.java b/src/org/kne/cloud/network/congestion/CongestionAlgorithms.java new file mode 100644 index 0000000..33f6d87 --- /dev/null +++ b/src/org/kne/cloud/network/congestion/CongestionAlgorithms.java @@ -0,0 +1,18 @@ +package org.kne.cloud.network.congestion; + +import java.util.HashMap; +import java.util.Map; + +public class CongestionAlgorithms { + private static Mapregs=new HashMap(); + static { + register(new BBRCongestionAlgorithmFactory()); + register(new Vegas2CongestionAlgorithmFactory()); + } + public static void register(CongestionAlgorithmFactory algorithm) { + regs.put(algorithm.getName(), algorithm); + } + public static CongestionAlgorithm get(String name) { + return regs.get(name).create(); + } +} diff --git a/src/org/kne/cloud/network/congestion/MpscMessageBatcher.java b/src/org/kne/cloud/network/congestion/MpscMessageBatcher.java new file mode 100644 index 0000000..cd21127 --- /dev/null +++ b/src/org/kne/cloud/network/congestion/MpscMessageBatcher.java @@ -0,0 +1,207 @@ +package org.kne.cloud.network.congestion; + +import org.jctools.queues.MpscArrayQueue; +import org.kne.cloud.network.ThreadTool; +import org.kne.opencl64.Releaser; + +import java.lang.ref.Cleaner; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.locks.LockSupport; +import java.util.function.Consumer; + +/** + * 基于 MPSC 队列的高性能消息批量发送器 + */ +public class MpscMessageBatcher implements MessageBatcher { + private static final Cleaner clr = Cleaner.create(); + + private final MpscMessageBatcher0 impl; + private final MpscMessageBatcherReleaser releaser; + + public MpscMessageBatcher(int batchSize, long maxDelayNanos) { + this.impl = new MpscMessageBatcher0<>(batchSize, maxDelayNanos, 65536); + this.releaser = new MpscMessageBatcherReleaser<>(impl); + clr.register(this, releaser); + } + + @Override + public int getBatchSize() { return impl.getBatchSize(); } + + @Override + public long getMaxDelay() { return impl.getMaxDelay(); } + + @Override + public void setConsumer(Consumer> consumer) { impl.setConsumer(consumer); } + + @Override + public void putMessage(T message) { impl.putMessage(message); } + + @Override + public void putMessages(List messages) { impl.putMessages(messages); } + + @Override + public void flush() { impl.flush(); } + + @Override + public int getQueueSize() { return impl.getQueueSize(); } + + @Override + public boolean isClosed() { return impl.isClosed(); } + + @Override + public void close() { impl.close(); } + + public static void main(String[] args) throws InterruptedException { + MpscMessageBatcher batcher = + new MpscMessageBatcher<>(10, 1_000_000L); + batcher.setConsumer(System.out::println); + + int n = 0; + for (;;) { + batcher.putMessage(n++); + batcher.putMessage(n++); + batcher.putMessage(n++); + batcher.putMessage(n++); + Thread.sleep(1); + } + } +} + +class MpscMessageBatcherReleaser extends Releaser> { + public MpscMessageBatcherReleaser(MpscMessageBatcher0 resource) { + super(resource); + } + + @Override + protected void release(MpscMessageBatcher0 resource) { + resource.close(); + } +} + +class MpscMessageBatcher0 implements Runnable, MessageBatcher { + + private final MpscArrayQueue queue; + private final int batchSize; + private final long maxDelay; + private volatile Consumer> consumer; + private final Thread batchThread; + private final AtomicBoolean closed = new AtomicBoolean(false); + private long firstTime; + + + public MpscMessageBatcher0(int batchSize, long maxDelay, int queueCapacity) { + this.batchSize = batchSize; + this.maxDelay = maxDelay; + this.queue = new MpscArrayQueue<>(queueCapacity); + + this.batchThread = ThreadTool.makeVDaemonThread("MpscMessageBatcher", this); + this.batchThread.start(); + } + + @Override + public void setConsumer(Consumer> consumer) { + this.consumer = consumer; + } + + @Override + public void putMessage(T message) { + if (message == null || closed.get()) return; + // 无锁入队 + while (!queue.offer(message)) { + // 队列满时,主动尝试 check 一次(自旋等待) + Thread.yield(); + } + check(false); + } + + @Override + public void putMessages(List messages) { + for (T msg : messages) { + putMessage(msg); + } + } + + @Override + public void run() { + while(!closed.get()) { + check(false); + // 更精确的等待策略 + long sleepTimeNanos = calculateSleepTime(); + if (sleepTimeNanos > 0) { + LockSupport.parkNanos(sleepTimeNanos); + } else { + // 避免忙等待 + LockSupport.parkNanos(1_000_000L); // 1ms + } + } + } + + private void check(boolean force) { + long currTime=System.nanoTime(); + if(queue.size()>=batchSize||(currTime-firstTime)>maxDelay||force) { + if(consumer!=null) { + ArrayListbatch=new ArrayList<>(batchSize); + queue.drain(batch::add, batchSize); + if(!batch.isEmpty()) + try { + consumer.accept(batch); + }catch(Exception e) { + e.printStackTrace(); + } + firstTime=currTime; + } + } + } + + /** + * 计算需要等待的时间 + * @return 等待时间(纳秒) + */ + private long calculateSleepTime() { + + + long elapsed = System.nanoTime() - firstTime; + long remaining = maxDelay - elapsed; + + if (remaining <= 0) { + return 0; // 立即处理 + } + + // 返回剩余时间或10ms中的较小值 + return Math.min(remaining, 10_000_000L); + } + + @Override + public void flush() { + check(true); + } + + @Override + public int getQueueSize() { + return queue.size(); + } + + @Override + public boolean isClosed() { + return closed.get(); + } + + @Override + public void close() { + closed.set(true); + LockSupport.unpark(batchThread); + flush(); // 确保剩余消息被处理 + } + + @Override + public int getBatchSize() { + return batchSize; + } + + @Override + public long getMaxDelay() { + return maxDelay; + } +} \ No newline at end of file diff --git a/src/org/kne/cloud/network/congestion/ReceivePacketSlidingWindow.java b/src/org/kne/cloud/network/congestion/ReceivePacketSlidingWindow.java new file mode 100644 index 0000000..3902222 --- /dev/null +++ b/src/org/kne/cloud/network/congestion/ReceivePacketSlidingWindow.java @@ -0,0 +1,146 @@ +package org.kne.cloud.network.congestion; + +import java.io.Closeable; +import java.io.IOException; +import java.net.SocketException; +import java.net.SocketTimeoutException; +import java.util.Collection; +import java.util.Iterator; +import java.util.Map; +import java.util.UUID; +import java.util.Map.Entry; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.LongAdder; +import java.util.concurrent.locks.Condition; +import java.util.concurrent.locks.Lock; +import java.util.function.Consumer; + +import org.kne.cloud.network.NetworkPacket; +import org.kne.cloud.network.ThreadTool; +import org.kne.cloud.network.ipv6.IPv6Packet; +import org.kne.cloud.network.klalb.DATATPacket; +import org.kne.cloud.network.klalb.SendItem; +import org.kne.cloud.network.kltp.KLTPPacket; +import org.kne.concurrent.ThreadParker; + +public class ReceivePacketSlidingWindow implements Closeable, AutoCloseable{ + private ThreadParker tp=new ThreadParker(); + + private Map> recvMap = new ConcurrentHashMap>(); + private AtomicLong recvWindowUsed=new AtomicLong(0); + //private LongAdder recvWindowUsed=new LongAdder(); + private long recvWindowSize; + private int headerCalibrate = 0; + + private volatile boolean closed = false; + + + public ReceivePacketSlidingWindow(long recvWindowSize) { + super(); + this.recvWindowSize = recvWindowSize; + } + + public ReceivePacketSlidingWindow(long recvWindowSize, int headerCalibrate) { + super(); + this.recvWindowSize = recvWindowSize; + this.headerCalibrate = headerCalibrate; + } + + public void setHeaderCalibrate(int headerCalibrate) { + this.headerCalibrate = headerCalibrate; + } + + public int getHeaderCalibrate() { + return headerCalibrate; + } + + public void put(K sequence, V packet) { + SendItem newitem = new SendItem(packet,headerCalibrate); + SendItem old = recvMap.putIfAbsent(sequence, newitem); + long dx=0; + if (old != null) { + dx=-old.getPacketLength() ; + } + recvWindowUsed.addAndGet(newitem.getPacketLength() +dx); + tp.unpark(); + } + + public V poll(K key) { + SendItem dtp2 = recvMap.remove(key); + if (dtp2 != null) { + recvWindowUsed.addAndGet(-dtp2.getPacketLength()); + return dtp2.getPacket(); + }else { + return null; + } + } + public V take(K key) throws SocketException { + while(true) { + V val=poll(key); + if(val!=null) { + return val; + } + if(isClosed()) { + throw new SocketException("Receive window closed!"); + }tp.parkNanos(1000000L); + } + } + + public V take(K key,long timeoutNanos) throws SocketException, SocketTimeoutException { + long start=System.nanoTime(); + while(true) { + if(System.nanoTime()-start>timeoutNanos) { + throw new SocketTimeoutException("Receive time out:"+(System.nanoTime()-start)+">"+timeoutNanos); + } + V val=poll(key); + if(val!=null) { + return val; + } + if(isClosed()) { + throw new SocketException("Receive window closed!"); + }tp.parkNanos(1000000L); + } + } + + public long getRecvWindowUsed() { + return recvWindowUsed.get(); + } + + public long getRecvWindowAvaliable() { + return Math.max(0, recvWindowSize-recvWindowUsed.get()); + } + + public Map> getRecvMap() { + return recvMap; + } + + public int getWindowPacketCount() { + return recvMap.size(); + } + + public boolean isEmpty() { + return recvMap.isEmpty(); + } + + public boolean isClosed() { + return closed; + } + + @Override + public void close() { + closed = true; + } + + @Override + public String toString() { + return "ReceivePacketSlidingWindow [recvWindowUsed=" + recvWindowUsed + ", recvWindowSize=" + recvWindowSize + + ", closed=" + closed + "]"; + } + + + + +} diff --git a/src/org/kne/cloud/network/congestion/Vegas2CongestionAlgorithmFactory.java b/src/org/kne/cloud/network/congestion/Vegas2CongestionAlgorithmFactory.java new file mode 100644 index 0000000..1eec4a9 --- /dev/null +++ b/src/org/kne/cloud/network/congestion/Vegas2CongestionAlgorithmFactory.java @@ -0,0 +1,20 @@ +package org.kne.cloud.network.congestion; + +public class Vegas2CongestionAlgorithmFactory implements CongestionAlgorithmFactory { + + @Override + public CongestionAlgorithm create() { + return new Vegas2CongestionAlgorithm(); + } + + @Override + public String getName() { + return "Vegas2"; + } + + @Override + public String toString() { + return getName(); + } + +} diff --git a/src/org/kne/cloud/network/ipv6/IPv6ProtocolRegister.java b/src/org/kne/cloud/network/ipv6/IPv6ProtocolRegister.java new file mode 100644 index 0000000..609075c --- /dev/null +++ b/src/org/kne/cloud/network/ipv6/IPv6ProtocolRegister.java @@ -0,0 +1,7 @@ +package org.kne.cloud.network.ipv6; + +import java.io.IOException; + +public interface IPv6ProtocolRegister { + public boolean onaccept(IPv6Packet packx)throws IOException; +} diff --git a/src/org/kne/cloud/network/klalb/KLALBProtocolRegister.java b/src/org/kne/cloud/network/klalb/KLALBProtocolRegister.java new file mode 100644 index 0000000..05c5f1a --- /dev/null +++ b/src/org/kne/cloud/network/klalb/KLALBProtocolRegister.java @@ -0,0 +1,44 @@ +package org.kne.cloud.network.klalb; + +import org.kne.cloud.network.ipv6.IPv6Address; +import org.kne.cloud.network.ipv6.IPv6Packet; +import org.kne.cloud.network.ipv6.IPv6Packet.IPv6Payload; +import org.kne.cloud.network.srv6.IPv6PacketConsumer; + +import java.io.*; +public class KLALBProtocolRegister extends PortBinder implements IPv6PacketConsumer { + private static final boolean showpacket=false; + private KLALBController controller; + + public KLALBProtocolRegister(KLALBController controller) { + super(controller.getSelf().getAddress()); + this.controller = controller; + } + + @Override + public void accept(IPv6Packet packx) throws IOException { + IPv6Payload pl = packx.getPayload(); + if (pl instanceof KLALBPacket) { + KLALBPacket rec = (KLALBPacket) pl; + if (showpacket) + System.out.println("KLALB_RX:" + rec); + if (rec instanceof PortPacket) { + rec.setCE(packx.isCE()); + IPv6Address srcA = packx.getSourceAddress(); + PortPacket pt = (PortPacket) rec; + BindableConsumer cons; + if ((cons=distributePacketToConsumer(srcA, pt))!=null) { + cons.accept(packx); + } else { + if (!(pt instanceof RSTPacket)) { + controller.getIpv6Router().enqueuePacketSendTask(() -> { + return controller.createPacketToAddress(srcA, 0, new RSTPacket(pt.getDstPort(), pt.getSrcPort()), + 2); + }); + } + } + } + } + } + + } \ No newline at end of file diff --git a/src/org/kne/cloud/network/kltp/KLTPInputStream.java b/src/org/kne/cloud/network/kltp/KLTPInputStream.java new file mode 100644 index 0000000..50bd8da --- /dev/null +++ b/src/org/kne/cloud/network/kltp/KLTPInputStream.java @@ -0,0 +1,302 @@ +package org.kne.cloud.network.kltp; + +import java.io.IOException; +import java.io.InputStream; +import java.net.BindException; +import java.net.SocketTimeoutException; +import java.nio.BufferOverflowException; +import java.nio.ByteBuffer; +import java.nio.channels.ReadableByteChannel; +import java.util.UUID; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; + +import org.kne.cloud.network.NetworkPacket; +import org.kne.cloud.network.congestion.MpscMessageBatcher; +import org.kne.cloud.network.congestion.ReceivePacketSlidingWindow; +import org.kne.cloud.network.ipv6.IPv6Address; +import org.kne.cloud.network.ipv6.IPv6Packet; +import org.kne.cloud.network.klalb.DATATPacket; +import org.kne.cloud.network.klalb.KLALBController; + +public class KLTPInputStream extends InputStream implements KLTPPacketConsumer, ReadableByteChannel{ + + + + private KLALBController controller; + + private IPv6Address remoteaddr; + + + private UUID streamUUID; + + + private ReceivePacketSlidingWindowrecvMap=new ReceivePacketSlidingWindow(Integer.MAX_VALUE,-20); + + + private AtomicLong inputcount = new AtomicLong(); + private KLTPPacket dataPack = null; + + private long soTimeout=0; + + public IPv6Address getRemoteAddress() { + return remoteaddr; + } + public KLTPInputStream(KLALBController controller,IPv6Address remoteaddr,UUID uuid) throws BindException { + this.controller=controller; + this.streamUUID =uuid; + this.remoteaddr=remoteaddr; + controller.getKLTPregister().registerReceiveStream(this); + + } + + + @Override + public int read() throws IOException { + if (dataPack == null ||(!dataPack.getKLTPData().hasRemaining())) { + dataPack=nextPacket(true); + } + if (dataPack.getDataSize() == 0) { + return -1; + } else { + int ret= dataPack.getKLTPData().get() & 0xff; + return ret; + } + } + private KLTPPacket nextPacket(boolean block) throws IOException { + try { + KLTPPacket dtp2 =null; + if(block) { + if(soTimeout==0) { + dtp2= recvMap.take(inputcount.get()); + + }else { + dtp2= recvMap.take(inputcount.get(),soTimeout); + } + }else { + dtp2=recvMap.poll(inputcount.get()); + } + if (dtp2 != null) { + inputcount.setPlain( inputcount.getPlain()+1); + int size=dtp2.getDataSize(); + //socketMonitor.getDownloadBandwidth().recordPacket(pid, size); + //controller.getDatatMonitor().getDownloadBandwidth().recordPacket(KLALBUtils.createGlobalUUID(), size); + //checkFlowControl(dtp2); + // System.out.println("序列号:"+dtp2.getSequence()); + return dtp2; + } + }catch(SocketTimeoutException e) { + close0(); + throw e; + } + return null; + } + + + + /*@Override + public int read(ByteBuffer dst) throws IOException { + int oldlmt=dst.limit(); + try { + + + if (dataPack == null ||(!dataPack.getKLTPData().hasRemaining())) { + nextPacket(); + } + if (dataPack.getDataSize() == 0) { + return -1; + } else { + int len = Math.min(dst.remaining(), available()); + dst.limit(dst.position()+len); + dst.put( dataPack.getKLTPData().get()) ; + + } + int i = 1; + try { + while (dst.hasRemaining()) { + + if (dataPack == null ||(!dataPack.getKLTPData().hasRemaining())) { + nextPacket(); + } + if (dataPack.getDataSize() == 0) { + break; + } + int min=Math.min(dataPack.getKLTPData().remaining(), dst.remaining()); + int oldlm=dataPack.getKLTPData().limit(); + dataPack.getKLTPData().limit(dataPack.getKLTPData().position()+min); + System.out.println("dst:"+dst+" datapack:"+dataPack); + dst.put(dataPack.getKLTPData()); + dataPack.getKLTPData().limit(oldlm); + i+=min; + + + } + } catch (IOException ee) { + } + return i; + }catch(BufferOverflowException e) { + System.err.println("dst:"+dst+" datapack:"+dataPack); + throw e; + }finally { + dst.limit(oldlmt); + } + }*/ + @Override + public int read(ByteBuffer dst) throws IOException { + if (!dst.hasRemaining()) { + return 0; + } + + int totalRead = 0; + + try { + // 如果当前没有数据包或当前数据包已读完,获取下一个 + if (dataPack == null || (!dataPack.getKLTPData().hasRemaining()&&(dataPack.getDataSize()!=0))) { + dataPack=nextPacket(true); + } + // EOF 检查 + if (dataPack.getDataSize() == 0) { + //System.out.println("EOF recv:"+dataPack); + return -1; + } + + // 循环读取直到 dst 满或没有更多数据 + while (dst.hasRemaining()) { + // 获取当前数据包的剩余数据 + ByteBuffer src = dataPack.getKLTPData(); + + if (!src.hasRemaining()) { + // 当前包读完,尝试获取下一个包 + dataPack=nextPacket(false); + if (dataPack==null||dataPack.getDataSize() == 0) { + break; // 下一个包还没来或EOF + } + src = dataPack.getKLTPData(); + } + + // 计算本次可拷贝的字节数 + int bytesToCopy = Math.min(src.remaining(), dst.remaining()); + + // 保存原 limit + int srcOldLimit = src.limit(); + int dstOldLimit = dst.limit(); + + try { + // 设置临时 limit + src.limit(src.position() + bytesToCopy); + dst.limit(dst.position() + bytesToCopy); + + // 执行拷贝 + dst.put(src); + totalRead += bytesToCopy; + } finally { + // 恢复 limit + src.limit(srcOldLimit); + dst.limit(dstOldLimit); + } + } + } catch (SocketTimeoutException e) { + close0(); + throw e; + } catch (BufferOverflowException e) { + // 不应该发生,因为我们做了 min() 检查 + throw new IOException("Buffer overflow in KLTPInputStream.read", e); + } + + return totalRead > 0 ? totalRead : -1; + } + + @Override + public int read(byte[] b, int off, int len) throws IOException { + return read(ByteBuffer.wrap(b,off,len)); + } + + @Override + public void close() throws IOException { + close0(); + } + private void close0() throws IOException{ + try { + recvMap.close(); + }finally { + controller.getKLTPregister().unregisterReceiveStream(this); + } + } + @Override + public int available() throws IOException { + //long i = recvMap.getRecvWindowUsed(); + long i=0; + if (dataPack != null) + i+=dataPack.getKLTPData().remaining(); + return (int) i; + } + + @Override + public boolean isOpen() { + return !recvMap.isClosed(); + } + + + public long read(ByteBuffer[] dsts, int offset, int length) throws IOException { + long lth=0; + for(int i=offset;i{ + KLTPPacket pack=new KLTPPacket(streamUUID,KLTPPacket.KLTP_TYPE_ACK,kltp.getSequence(),0); + pack.setCE(u.isCE()); + return controller.createPacketToAddress(remoteaddr,0,pack); + }); + break; + case KLTPPacket.KLTP_TYPE_DATAFIN: + //ackSequenceBatcher.putMessage(kseq2); + recvMap.put(kltp.getSequence(), kltp); + controller.getIpv6Router().runPacketSendTask(()->{ + KLTPPacket pack=new KLTPPacket(streamUUID,KLTPPacket.KLTP_TYPE_ACK,kltp.getSequence(),0); + pack.setCE(u.isCE()); + return controller.createPacketToAddress(remoteaddr,0,pack); + }); + //System.out.println(inputcount+" "+ recvMap.getRecvMap()); + break; + } + } + } + + + @Override + public UUID getStreamUUID() { + return streamUUID; + } + public boolean isClosed() { + return recvMap.isClosed(); + } + public void setSoTimeout(int value) { + soTimeout=value*1000000L; + } + public int getSoTimeout() { + return (int) (soTimeout/1000000L); + } + @Override + public String toString() { + return "KLTPInputStream [streamUUID=" + streamUUID + ", recvMap=" + recvMap + "]"; + } + + + + } \ No newline at end of file diff --git a/src/org/kne/cloud/network/kltp/KLTPOutputStream.java b/src/org/kne/cloud/network/kltp/KLTPOutputStream.java new file mode 100644 index 0000000..5bfd38e --- /dev/null +++ b/src/org/kne/cloud/network/kltp/KLTPOutputStream.java @@ -0,0 +1,375 @@ +package org.kne.cloud.network.kltp; + +import java.io.IOException; +import java.io.OutputStream; +import java.net.BindException; +import java.net.Inet6Address; +import java.net.SocketException; +import java.net.SocketTimeoutException; +import java.nio.ByteBuffer; +import java.nio.channels.WritableByteChannel; +import org.kne.concurrent.*; +import org.kne.cloud.network.NetworkPacket; +import org.kne.cloud.network.congestion.CongestionAlgorithm; +import org.kne.cloud.network.congestion.DCTCP2CongestionAlgorithm; +import org.kne.cloud.network.congestion.DCTCPCongestionAlgorithm; +import org.kne.cloud.network.congestion.SendPacketSlidingWindow; +import org.kne.cloud.network.ipv6.IPv6Address; +import org.kne.cloud.network.ipv6.IPv6Packet; +import org.kne.cloud.network.klalb.DATATPacket; +import org.kne.cloud.network.klalb.KLALBController; + +import java.util.UUID; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.locks.*; +public class KLTPOutputStream extends OutputStream implements KLTPPacketConsumer,WritableByteChannel{ + private static final int HEADER_CALIBRATE = 150; + private KLALBController controller; + + private IPv6Address remoteaddr; + public IPv6Address getRemoteAddress() { + return remoteaddr; + } + public KLTPOutputStream(KLALBController controller,IPv6Address remoteaddr,UUID uuid) throws BindException { + this.controller=controller; + this.streamUUID =uuid; + this.remoteaddr=remoteaddr; + if(controller.getConfigItem()!=null) + this.delaytime=controller.getConfigItem().getNagleDelayTime(); + + algorithm.setWindowControlConsumer((window)->{ + sendMap.setWindowSize(window); + }); + + controller.getKLTPregister().registerSendStream(this); + } + private UUID streamUUID; + + + private int MTU=8192; + + private CongestionAlgorithm algorithm=new DCTCPCongestionAlgorithm(); + private SendPacketSlidingWindowsendMap=new SendPacketSlidingWindow<>(algorithm, 4*MTU,HEADER_CALIBRATE,false); + { + sendMap.setResendConsumer((dtp)->{ + if(dtp.getSendCounter()>=50) { + System.err.println("send error!"); + try { + close0(false); + } catch (IOException e) { + e.printStackTrace(); + } + return; + } + + controller.getIpv6Router().runPacketSendTask(()->{ + long length= dtp.getTotalLength(); + //socketRawMonitor.getUploadBandwidth().recordPacket(pidg.generate(), (int) length); + //System.out.println("第"+(dtp.getSendCounter()-1)+"次重传:"+dtp); + return controller.createPacketToAddress( remoteaddr,0,dtp,1); + }); + dtp.incSendCounter(); + }); + } + + private AtomicLong outputcount=new AtomicLong( 0); + + private KLTPPacket dataPack =null; + + + private long delaytime=2; + + + + private Lock olock=new SpinLock(); + + private volatile boolean autoFlush=true; + @Override + public void write(int b) throws IOException { + if (sendMap.isClosed()) + throw new SocketException("Socket is closed"); + if(olock!=null) + olock.lock(); + createDataPack(); + try { + dataPack.getKLTPData().put((byte) b); + if (dataPack!=null&& dataPack.getKLTPData().hasRemaining()) { + if(autoFlush) + delayFlush(); + }else { + flush0(); + } + }finally { + if(olock!=null) + olock.unlock(); + } + } + + + + + + + + + @Override + public void write(byte[] b, int off, int len) throws IOException { + write(ByteBuffer.wrap(b, off, len)); + return ; + } + + private int writeWithoutFlush(ByteBuffer src)throws IOException { + if (sendMap.isClosed()) + throw new SocketException("Socket is closed"); + int counter=0; + if(olock!=null) + olock.lock(); + try { + + while(src.hasRemaining()) { + createDataPack(); + int min=Math.min(src.remaining(), dataPack.getKLTPData().remaining()); + int olm=src.limit(); + src.limit(src.position()+min); + dataPack.getKLTPData().put(src); + counter+=min; + src.limit(olm); + + if (!dataPack.getKLTPData().hasRemaining()) { + flush0(); + } + } + + + }finally { + if(olock!=null) + olock.unlock(); + } + return counter; + } + + @Override + public int write(ByteBuffer src) throws IOException { + if (sendMap.isClosed()) + throw new SocketException("Socket is closed"); + int counter=0; + if(olock!=null) + olock.lock(); + try { + + while(src.hasRemaining()) { + createDataPack(); + int min=Math.min(src.remaining(), dataPack.getKLTPData().remaining()); + int olm=src.limit(); + src.limit(src.position()+min); + dataPack.getKLTPData().put(src); + counter+=min; + src.limit(olm); + + if (!dataPack.getKLTPData().hasRemaining()) { + flush0(); + } + } + + if(dataPack!=null&& dataPack.getKLTPData().position()>0) { + if(autoFlush) + delayFlush(); + } + }finally { + if(olock!=null) + olock.unlock(); + } + return counter; + } + + + private void createDataPack() { + if(dataPack==null) { + long pl=outputcount.getPlain(); + dataPack=new KLTPPacket(streamUUID,KLTPPacket.KLTP_TYPE_DATA,pl++, MTU); + outputcount.setPlain(pl); + //System.out.println("EOF:"+dataPack); + } + } + public void waitForAllAcknowledged(int timeout) throws IOException { + long start=System.nanoTime(); + while(true){ + if (sendMap.isClosed()) + throw new SocketException("Socket is closed"); + if(sendMap.isEmpty()) + break; + if(timeout!=0&&(System.nanoTime()-start>timeout*1000000L)) + throw new SocketTimeoutException("wait for acknowledged timout"); + } + } + AtomicReference ioe=new AtomicReference<>(); + private volatile ScheduledFuture tt; + + @Override + public void flush() throws IOException { + if(olock!=null) + olock.lock(); + try { + delayFlush(); + }finally { + if(olock!=null) + olock.unlock(); + } + } + public void forceFlush() throws IOException{ + if(olock!=null) + olock.lock(); + try { + flush0(); + }finally { + if(olock!=null) + olock.unlock(); + } + } + public void delayFlush()throws IOException{ + if(delaytime<=0) { + flush0(); + }else { + if(tt==null) { + Runnable r= new Runnable() { + + @Override + public void run() { + if(sendMap.isClosed()) + tt.cancel(false); + try { + if(olock!=null) + olock.lock(); + try { + flush0(); + }finally { + if(olock!=null) + olock.unlock(); + } + } catch (IOException e) { + ioe.set(e); + } + } + }; + tt=controller.getScheduleTimer().scheduleAtFixedRate (r, delaytime, delaytime,TimeUnit.NANOSECONDS); + } + + IOException ioex=ioe.get(); + if(ioex!=null) { + ioex.fillInStackTrace(); + throw ioex; + } + } + } + + private void flush0() throws IOException { + KLTPPacket pack=dataPack; + if (pack!=null&&pack.getKLTPData().position() > 0) { + + sendMap.waitForAvaliable(); + + controller.getIpv6Router().runPacketSendTask(()->{ + pack.getKLTPData().flip(); + pack.incSendCounter(); + sendMap.put(pack.getSequence(),pack); + + int size= pack.getKLTPData().limit(); + //socketMonitor.getUploadBandwidth().recordPacket(pid, size); + //controller.getDatatMonitor().getUploadBandwidth().recordPacket(KLALBUtils.createGlobalUUID(), size); + + //socketRawMonitor.getUploadBandwidth().recordPacket(pid, size); + return controller.createPacketToAddress(remoteaddr,0,pack); + + }); + dataPack=null; + } + } + @Override + public void close() throws IOException { + close0(true); + } + private void close0(boolean grace) throws IOException { + try { + if(sendMap.isClosed()) + return; + if(grace) { + if(olock!=null) + olock.lock(); + try { + flush0(); + }finally { + if(olock!=null) + olock.unlock(); + } + + long pl=outputcount.getPlain(); + KLTPPacket pack=new KLTPPacket(streamUUID,KLTPPacket.KLTP_TYPE_DATAFIN, pl++ ,MTU); + outputcount.setPlain(pl); + pack.getKLTPData(). flip(); + pack.incSendCounter(); + controller.getIpv6Router().runPacketSendTask (()->{ + return controller.createPacketToAddress(remoteaddr,0,pack); + }); + sendMap.put(pack.getSequence(),pack); + // System.out.println("EOF send:"+pack); + } + + sendMap.close(); + }finally { + if(tt!=null) + tt.cancel(false); + controller.getKLTPregister().unregisterSendStream(this); + } + } + @Override + public boolean isOpen() { + return !sendMap.isClosed(); + } + public boolean isAutoFlush() { + return autoFlush; + } + public void setAutoFlush(boolean b) { + autoFlush=b; + } + public long write(ByteBuffer[] srcs, int offset, int length) throws IOException { + long lth=0; + for(int i=offset;i0) { + if(autoFlush) + flush(); + } + return lth; + } + + + @Override + public void accept(IPv6Packet u) { + if(u.getPayload() instanceof KLTPPacket) { + KLTPPacket kltp=(KLTPPacket) u.getPayload(); + switch(kltp.getType()) { + case KLTPPacket.KLTP_TYPE_ACK: + sendMap.ack(kltp.getSequence(), kltp.isCE()); + + break; + } + } + } + + @Override + public UUID getStreamUUID() { + return streamUUID; + } + public boolean isClosed() { + return sendMap.isClosed(); + } + @Override + public String toString() { + return "KLTPOutputStream [streamUUID=" + streamUUID + ", sendMap=" + sendMap + "]"; + } + + } \ No newline at end of file diff --git a/src/org/kne/cloud/network/kltp/KLTPPacket.java b/src/org/kne/cloud/network/kltp/KLTPPacket.java new file mode 100644 index 0000000..67e03aa --- /dev/null +++ b/src/org/kne/cloud/network/kltp/KLTPPacket.java @@ -0,0 +1,186 @@ +package org.kne.cloud.network.kltp; + +import java.io.IOException; +import java.nio.Buffer; +import java.nio.ByteBuffer; +import java.nio.channels.ReadableByteChannel; +import java.nio.channels.WritableByteChannel; +import java.util.List; +import java.util.UUID; +import java.util.function.Consumer; + +import org.kne.cloud.network.NetworkPacket; +import org.kne.cloud.network.ipv6.IPv6Packet.IPv6Payload; +import org.kne.io.KNEChannels; + +public class KLTPPacket extends IPv6Payload { + public static final int KLTP_PROTOCOL_NUMBER=253; + public static final int KLTP_HEADER_LENGTH=24+8; + + public static final int KLTP_TYPE_DATA=0; + public static final int KLTP_TYPE_DATAFIN=1; + public static final int KLTP_TYPE_ACK=2; + + protected ByteBuffer kltpHeader; + + protected ByteBuffer kltpData; + + + public KLTPPacket() { + super(KLTP_PROTOCOL_NUMBER); + kltpHeader=NetworkPacket.bufferAllocator.allocate(KLTP_HEADER_LENGTH); + } + + public KLTPPacket(UUID uuid,int type,long seq,int mtulimit) { + super(KLTP_PROTOCOL_NUMBER); + kltpHeader=NetworkPacket.bufferAllocator.allocate(KLTP_HEADER_LENGTH); + setUUID(uuid); + setType(type); + setSequence(seq); + kltpData=NetworkPacket.bufferAllocator.allocate(mtulimit); + } + + + + public UUID getUUID() { + long h=kltpHeader.getLong(0); + long l=kltpHeader.getLong(8); + return new UUID(h,l); + } + + public void setUUID(UUID uuid) { + kltpHeader.putLong(0, uuid.getMostSignificantBits()); + kltpHeader.putLong(8, uuid.getLeastSignificantBits()); + } + + public int getPayloadLength() { + return kltpHeader.getInt(16); + } + + public void setPayloadLength(int payloadLength) { + kltpHeader.putInt(16,payloadLength); + } + + public int getChecksum() { + return kltpHeader.getChar(20); + } + + public void setChecksum(int checksum) { + kltpHeader.putChar(20, (char) checksum); + } + + public int getType() { + return kltpHeader.get(22); + } + + public void setType(int type) { + kltpHeader.put(22, (byte) type); + } + + public boolean isCE() { + int v=kltpHeader.get(23)&1; + return v!=0; + } + + public void setCE(boolean b) { + kltpHeader.put(23, (byte) (b?1:0)); + } + + public long getSequence() { + return kltpHeader.getLong(24); + } + + + public void setSequence(long kseq) { + kltpHeader.putLong(24,kseq); + } + + + + @Override + public long getTotalLength() { + return KLTP_HEADER_LENGTH+kltpData.limit(); + } + + @Override + public void writeToChannel(WritableByteChannel dto) throws IOException { + setPayloadLength(kltpData.limit()); + dto.write(kltpHeader.slice(0, KLTP_HEADER_LENGTH)); + dto.write(kltpData.slice(0, kltpData.limit())); + } + + @Override + public void readFromChannel(ReadableByteChannel din, long length) throws IOException { + kltpHeader.clear(); + KNEChannels.readFully(din ,kltpHeader); + kltpHeader.flip(); + kltpData=NetworkPacket.bufferAllocator.allocate(getPayloadLength()); + KNEChannels.readFully(din, kltpData); + kltpData.flip(); + } + + public static IPv6Payload readKLTPPacketFromChannel(ReadableByteChannel din) throws IOException { + KLTPPacket pack=new KLTPPacket(); + pack.readFromChannel(din); + return pack; + } + + + + private int sendCounter=0; + + public int getSendCounter() { + return sendCounter; + } + + public void incSendCounter() { + sendCounter++; + } + + public int getDataSize() { + return kltpData.limit(); + } + + public String toString() { + StringBuilder sb=new StringBuilder(); + switch(getType()) { + case KLTP_TYPE_DATA: + sb.append("DATA "); + sb.append(getUUID()); + sb.append(' '); + sb.append(getSequence()); + sb.append(' '); + sb.append(getKLTPData()); + break; + case KLTP_TYPE_DATAFIN: + sb.append("DATAFIN "); + sb.append(getUUID()); + sb.append(' '); + sb.append(getSequence()); + sb.append(' '); + sb.append(getKLTPData()); + break; + case KLTP_TYPE_ACK: + sb.append("ACK "); + sb.append(getUUID()); + sb.append(' '); + sb.append(getSequence()); + break; + default: + sb.append("UNKNOWN "); + sb.append(getUUID()); + break; + } + return sb.toString(); + } + + public ByteBuffer getKLTPData() { + return kltpData; + } + + + + + + +} diff --git a/src/org/kne/cloud/network/kltp/KLTPPacketConsumer.java b/src/org/kne/cloud/network/kltp/KLTPPacketConsumer.java new file mode 100644 index 0000000..3bddb29 --- /dev/null +++ b/src/org/kne/cloud/network/kltp/KLTPPacketConsumer.java @@ -0,0 +1,10 @@ +package org.kne.cloud.network.kltp; + +import java.util.UUID; + +import org.kne.cloud.network.srv6.IPv6PacketConsumer; + + +public interface KLTPPacketConsumer extends IPv6PacketConsumer { + public UUID getStreamUUID(); +} diff --git a/src/org/kne/cloud/network/kltp/KLTPProtocolRegister.java b/src/org/kne/cloud/network/kltp/KLTPProtocolRegister.java new file mode 100644 index 0000000..353ac7c --- /dev/null +++ b/src/org/kne/cloud/network/kltp/KLTPProtocolRegister.java @@ -0,0 +1,108 @@ +package org.kne.cloud.network.kltp; + +import java.io.IOException; +import java.net.BindException; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; + +import org.kne.cloud.network.ipv6.IPv6Address; +import org.kne.cloud.network.ipv6.IPv6Packet; +import org.kne.cloud.network.ipv6.IPv6Packet.IPv6Payload; +import org.kne.cloud.network.ipv6.IPv6ProtocolRegister; +import org.kne.cloud.network.klalb.BindableConsumer; +import org.kne.cloud.network.klalb.KLALBController; +import org.kne.cloud.network.klalb.PortBinder; +import org.kne.cloud.network.srv6.IPv6PacketConsumer; + +public class KLTPProtocolRegister extends PortBinder implements IPv6ProtocolRegister { + private KLALBController controller; + + public KLTPProtocolRegister(KLALBController controller) { + super(controller.getSelf().getAddress()); + this.controller = controller; + } + private static final boolean showpacket=false; + private static final boolean debug=false; + private ConcurrentHashMap recvRegisterMap=new ConcurrentHashMap(); + private ConcurrentHashMap sendRegisterMap=new ConcurrentHashMap(); + public void registerReceiveStream(KLTPPacketConsumer kltp) throws BindException { + if(recvRegisterMap.putIfAbsent(kltp.getStreamUUID(), kltp)!=null) { + throw new BindException("KLTP receive UUID "+kltp.getStreamUUID()+" already used!"); + }else { + if(debug) + System.out.println("接收流打开:"+kltp.getStreamUUID()); + } + } + + public void unregisterReceiveStream(KLTPPacketConsumer kltp) { + recvRegisterMap.remove(kltp.getStreamUUID(), kltp); + if(debug) + System.out.println("接收流关闭:"+kltp.getStreamUUID()); + } + + public void registerSendStream(KLTPPacketConsumer kltp) throws BindException { + if(sendRegisterMap.putIfAbsent(kltp.getStreamUUID(), kltp)!=null) { + throw new BindException("KLTP send UUID "+kltp.getStreamUUID()+" already used!"); + }else { + if(debug) + System.out.println("发送流打开:"+kltp.getStreamUUID()); + } + } + + public void unregisterSendStream(KLTPPacketConsumer kltp) { + sendRegisterMap.remove(kltp.getStreamUUID(), kltp); + if(debug) + System.out.println("发送流关闭:"+kltp.getStreamUUID()); + } + + @Override + public boolean onaccept(IPv6Packet packx) throws IOException { + IPv6Payload pl = packx.getPayload(); + if (pl instanceof KLTPPacket) { + KLTPPacket kltp = (KLTPPacket) pl; + if (showpacket) + System.out.println("KLTP_RX:" + kltp); + if(kltp.getType() ==KLTPPacket.KLTP_TYPE_ACK) { + KLTPPacketConsumer cosu= sendRegisterMap.get(kltp.getUUID()); + if(cosu!=null) { + cosu.accept(packx); + return true; + } + }else { + KLTPPacketConsumer cosu= recvRegisterMap.get(kltp.getUUID()); + if(cosu!=null) { + cosu.accept(packx); + return true; + }else { + if(kltp.getSequence()==0) { + IPv6Address srca=packx.getSourceAddress(); + KLTPInputStream kins=new KLTPInputStream(controller, srca, kltp.getUUID()); + kins.accept(packx); + KLTPSessionPacket sess=new KLTPSessionPacket(kins); + sess.readFromChannel(kins); + BindableConsumer con; + if((con=distributePacketToConsumer(srca, sess))!=null) { + //System.out.println(this); + con.accept(sess); + System.out.println("接受连接:"+sess); + return true; + }else { + kins.close(); + System.out.println("丢弃连接:"+sess); + } + + }else { + //System.out.println("丢弃连接:"+kltp); + + } + } + } + + + + } + return false; + } + + } \ No newline at end of file diff --git a/src/org/kne/cloud/network/kltp/KLTPSessionPacket.java b/src/org/kne/cloud/network/kltp/KLTPSessionPacket.java new file mode 100644 index 0000000..964f0e3 --- /dev/null +++ b/src/org/kne/cloud/network/kltp/KLTPSessionPacket.java @@ -0,0 +1,105 @@ +package org.kne.cloud.network.kltp; + +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.channels.ReadableByteChannel; +import java.nio.channels.WritableByteChannel; + +import org.kne.cloud.network.ByteBufferAllocator; +import org.kne.cloud.network.NetworkPacket; +import org.kne.cloud.network.klalb.PortPacket; +import org.kne.io.KNEChannels; + +public class KLTPSessionPacket extends NetworkPacket implements PortPacket{ + private static final int KLTP_SESSION_HEADER_LENGTH=8; + private ByteBuffer header=NetworkPacket.bufferAllocator.allocate(KLTP_SESSION_HEADER_LENGTH); + + + private KLTPInputStream inputstream; + public KLTPSessionPacket(int sport, int dport,KLTPInputStream inputstream) { + setSrcPort(sport); + setDstPort(dport); + this.inputstream=inputstream; + } + public KLTPSessionPacket(int sport, int dport) { + setSrcPort(sport); + setDstPort(dport); + } + public KLTPSessionPacket() { + + } + + + public KLTPSessionPacket(KLTPInputStream inputstream) { + super(); + this.inputstream = inputstream; + } + public KLTPInputStream getInputstream() { + return inputstream; + } + + @Override + public long getTotalLength() { + return KLTP_SESSION_HEADER_LENGTH; + } + + @Override + public void writeToChannel(WritableByteChannel dto) throws IOException { + dto.write(header.slice(0, KLTP_SESSION_HEADER_LENGTH)); + //System.out.println("writesession:"+header); + } + + @Override + public void readFromChannel(ReadableByteChannel din, long length) throws IOException { + header.limit(KLTP_SESSION_HEADER_LENGTH); + KNEChannels.readFully(din, header); + header.flip(); + //System.out.println("readsesion:"+header); + } + + + + @Override + public int hashCode() { + return getSrcPort()^getDstPort(); + } + @Override + public boolean equals(Object obj) { + if (this == obj) + return true; + if (obj == null) + return false; + if (getClass() != obj.getClass()) + return false; + KLTPSessionPacket other = (KLTPSessionPacket) obj; + + return (getSrcPort()==other.getSrcPort())&&(getDstPort()==other.getDstPort()); + } + @Override + protected boolean needEndPosition() { + return false; + } + + public void setSrcPort(int sport) { + header.putInt(0,sport); + } + + public void setDstPort(int dport) { + header.putInt(4,dport); + } + + @Override + public int getSrcPort() { + return header.getInt(0); + } + + @Override + public int getDstPort() { + return header.getInt(4); + } + @Override + public String toString() { + return "KLTPSession "+getSrcPort()+"->"+getDstPort(); + } + +} diff --git a/src/org/kne/cloud/network/monitor/BandwidthSampler.java b/src/org/kne/cloud/network/monitor/BandwidthSampler.java new file mode 100644 index 0000000..542eb94 --- /dev/null +++ b/src/org/kne/cloud/network/monitor/BandwidthSampler.java @@ -0,0 +1,86 @@ +package org.kne.cloud.network.monitor; +import java.util.concurrent.atomic.LongAdder; + +/** + * 高性能网络流量统计器 + * + * 设计要点: + * 1. 数据面使用 LongAdder 无锁累加,完全不阻塞。 + * 2. 控制面使用快照缓存,避免每次都调用 sum() 遍历 Cell。 + * 3. 支持带宽(Bps)、包速率(PPS)、平均包大小(bytes/pkt)统计。 + */ +public class BandwidthSampler { + + // 数据面累加器(无锁) + private final LongAdder packetCount = new LongAdder(); + private final LongAdder byteCount = new LongAdder(); + + // 快照缓存(控制面使用,避免高频 sum()) + private volatile long cachedPacketCount = 0; + private volatile long cachedByteCount = 0; + private volatile long lastSnapshotTime = 0; + + // 统计结果缓存 + private volatile double currentBandwidthBps = 0.0; + private volatile double currentPacketRatePps = 0.0; + private volatile double currentAvgPacketSize = 0.0; // 新增:平均包大小(字节/包) + + /** + * 数据面调用:记录一个包 + * @param packetSizeBytes 包大小(字节) + */ + public void recordPacket(int packetSizeBytes) { + packetCount.increment(); + byteCount.add(packetSizeBytes); + } + + /** + * 控制面调用:更新统计快照(建议每 1 秒或每 1ms 调用一次) + * 计算带宽、PPS、平均包大小,并重置累加器 + */ + public void update() { + long now = System.nanoTime(); + + // 取当前累加值(会遍历 Cell,但频率低,可接受) + long currPackets = packetCount.sumThenReset(); + long currBytes = byteCount.sumThenReset(); + + // 计算时间间隔(秒) + double intervalSec = (lastSnapshotTime == 0) ? 1.0 : (now - lastSnapshotTime) / 1_000_000_000.0; + if (intervalSec <= 0) intervalSec = 1.0; + + // 更新缓存 + cachedPacketCount = currPackets; + cachedByteCount = currBytes; + + // 计算指标 + currentBandwidthBps = currBytes / intervalSec; + currentPacketRatePps = currPackets / intervalSec; + // 平均包大小 = 总字节数 / 总包数(若无包则为 0) + currentAvgPacketSize = (currPackets == 0) ? 0.0 : (double) currBytes / currPackets; + + lastSnapshotTime = now; + } + + // ========== 查询接口(直接返回缓存,无计算开销)========== + public double getBandwidthBps() { + return currentBandwidthBps; + } + + public double getPacketRatePps() { + return currentPacketRatePps; + } + + public double getAvgPacketSize() { + return currentAvgPacketSize; + } + + // 原始累加值 + public long getPacketCountSinceLastSnapshot() { + return cachedPacketCount; + } + + public long getByteCountSinceLastSnapshot() { + return cachedByteCount; + } +} \ No newline at end of file diff --git a/src/org/kne/cloud/network/monitor/CostSupplierFactory.java b/src/org/kne/cloud/network/monitor/CostSupplierFactory.java new file mode 100644 index 0000000..b36f661 --- /dev/null +++ b/src/org/kne/cloud/network/monitor/CostSupplierFactory.java @@ -0,0 +1,56 @@ +package org.kne.cloud.network.monitor; + +import java.util.function.Supplier; + +public class CostSupplierFactory { + /** + * 从 DelayMonitorData 获取 OWD 作为成本 + * + * @param monitor 延迟监控数据 + * @return 返回 OWD 的 Supplier + */ + public static Supplier owdSupplier(DelayMonitorData monitor) { + return () -> monitor.getOutDelay(); + } + + /** + * 静态成本 Supplier(用于测试或静态路由) + * + * @param cost 固定的成本值 + * @return 返回固定值的 Supplier + */ + public static Supplier staticSupplier(long cost) { + return () -> cost; + } + /** + * 使用你设计的“概率期望延迟”公式:OWD + RTO × (1 - Reliability) + * @param monitor 延迟监控数据(提供 OWD) + * @param linkStatus 链路状态(提供 Reliability) + * @param rtoNanos 超时重传时间(纳秒) + * @return 返回期望延迟的 Supplier + */ + public static Supplier expectedDelaySupplier(DelayMonitorData monitor, LinkStatus linkStatus, Supplier rtoNanos) { + return () -> { + long owd = monitor.getOutDelay(); + double reliability = linkStatus.getReliability(); +// 期望延迟 = OWD + RTO × (1 - Reliability) + long exp=(long) (owd + rtoNanos.get() * (1 - reliability)); + //System.out.println("OWD:"+owd+" EXP:"+exp); + return exp; + }; + } + + /** + * 组合两个 Supplier,取最大值(可用于 ECMP 场景下的保守调度) + */ + public static Supplier maxSupplier(Supplier a, Supplier b) { + return () -> Math.max(a.get(), b.get()); + } + + /** + * 组合两个 Supplier,取最小值(可用于 ECMP 场景下的乐观调度) + */ + public static Supplier minSupplier(Supplier a, Supplier b) { + return () -> Math.min(a.get(), b.get()); + } +} diff --git a/src/org/kne/cloud/network/monitor/DelaySampler.java b/src/org/kne/cloud/network/monitor/DelaySampler.java new file mode 100644 index 0000000..5acba62 --- /dev/null +++ b/src/org/kne/cloud/network/monitor/DelaySampler.java @@ -0,0 +1,160 @@ +package org.kne.cloud.network.monitor; + +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.LongAdder; + +/** + * 高性能网络延迟统计器 + * + * 设计要点: + * 1. 数据面使用 LongAdder 无锁累加总延迟,同时用 AtomicLong 原子记录最值。 + * 2. 控制面使用快照缓存,计算平均延迟、最小延迟、最大延迟和抖动。 + * 3. 采样方式支持:每包采样(高频)或每N包采样(低频),避免测量本身成为开销。 + */ +public class DelaySampler { + + // 数据面累加器(用于计算平均延迟) + private final LongAdder totalDelayNanos = new LongAdder(); + private final LongAdder packetCount = new LongAdder(); + + // 最值记录(使用 AtomicLong,保证原子更新,彻底避免读到中间状态) + private final AtomicLong minDelayNanos = new AtomicLong(Long.MAX_VALUE); + private final AtomicLong maxDelayNanos = new AtomicLong(0); + + // 快照缓存(控制面使用) + private volatile long snapshotTotalDelay = 0; + private volatile long snapshotPacketCount = 0; + private volatile long snapshotMinDelay = 0; + private volatile long snapshotMaxDelay = 0; + private volatile long lastSnapshotTime = 0; + private volatile long currentAvgDelayNanos = 0; + private volatile long currentMinDelayNanos = 0; + private volatile long currentMaxDelayNanos = 0; + private volatile long currentJitterNanos = 0; // 抖动:平均绝对偏差(基于相邻包延迟差) + + // 可选:用于计算抖动的历史延迟总和(或保留上次延迟值) + private volatile long lastDelayNanos = 0; + private final LongAdder jitterSumAbs = new LongAdder(); // 绝对偏差累积和(|delay_i - delay_{i-1}|) + + /** + * 数据面调用:记录一个包的延迟(纳秒) + * @param delayNanos 延迟(纳秒) + */ + public void recordDelay(long delayNanos) { + packetCount.increment(); + totalDelayNanos.add(delayNanos); + + // 更新最值(无锁自旋 CAS,线程安全) + updateMin(delayNanos); + updateMax(delayNanos); + + // 更新抖动:记录本次延迟与上次的差值绝对值 + long last = lastDelayNanos; + if (last != 0) { + jitterSumAbs.add(Math.abs(delayNanos - last)); + } + lastDelayNanos = delayNanos; + } + + /** + * 更新最小值(无锁 CAS 自旋) + */ + private void updateMin(long delayNanos) { + long min; + do { + min = minDelayNanos.get(); + if (delayNanos >= min) { + return; // 不是新最小值,直接返回 + } + } while (!minDelayNanos.compareAndSet(min, delayNanos)); + } + + /** + * 更新最大值(无锁 CAS 自旋) + */ + private void updateMax(long delayNanos) { + long max; + do { + max = maxDelayNanos.get(); + if (delayNanos <= max) { + return; // 不是新最大值,直接返回 + } + } while (!maxDelayNanos.compareAndSet(max, delayNanos)); + } + + /** + * 控制面调用:更新统计快照(建议与 BandwidthSampler.update() 同频调用) + * 计算平均延迟、最小延迟、最大延迟、抖动,并重置累加器 + */ + public void update() { + long now = System.nanoTime(); + + // 取当前累加值并重置 + long currPackets = packetCount.sumThenReset(); + long currTotalDelay = totalDelayNanos.sumThenReset(); + long currJitterSum = jitterSumAbs.sumThenReset(); + + // 取当前最值并重置(重置为初始值) + long currMin = minDelayNanos.getAndSet(Long.MAX_VALUE); + long currMax = maxDelayNanos.getAndSet(0); + + // 更新时间间隔(秒) + double intervalSec = (lastSnapshotTime == 0) ? 1.0 : (now - lastSnapshotTime) / 1_000_000_000.0; + if (intervalSec <= 0) intervalSec = 1.0; + + // 更新快照缓存 + snapshotPacketCount = currPackets; + snapshotTotalDelay = currTotalDelay; + snapshotMinDelay = currMin; + snapshotMaxDelay = currMax; + + // 计算统计指标 + if (currPackets > 0) { + currentAvgDelayNanos = (long) ((double) currTotalDelay / currPackets); + currentMinDelayNanos = currMin; + currentMaxDelayNanos = currMax; + + // 抖动:平均绝对偏差(MAD) = 累积绝对偏差 / (包数 - 1) + if (currJitterSum > 0 && currPackets > 1) { + currentJitterNanos = (long) ((double) currJitterSum / (currPackets - 1)); + } else { + currentJitterNanos = 0; + } + } else { + currentAvgDelayNanos = 0; + currentMinDelayNanos = 0; + currentMaxDelayNanos = 0; + currentJitterNanos = 0; + } + + // 重置 lastDelay,避免跨间隔的抖动误差 + lastDelayNanos = 0; + + lastSnapshotTime = now; + } + + // ========== 查询接口(直接返回缓存,无计算开销)========== + public long getAvgDelayNanos() { + return currentAvgDelayNanos; + } + + public long getMinDelayNanos() { + return currentMinDelayNanos; + } + + public long getMaxDelayNanos() { + return currentMaxDelayNanos; + } + + public long getJitterNanos() { + return currentJitterNanos; + } + + public long getSnapshotPacketCount() { + return snapshotPacketCount; + } + + public long getSnapshotTotalDelay() { + return snapshotTotalDelay; + } +} \ No newline at end of file diff --git a/src/org/kne/cloud/network/monitor/LinkStatus.java b/src/org/kne/cloud/network/monitor/LinkStatus.java new file mode 100644 index 0000000..a6e6292 --- /dev/null +++ b/src/org/kne/cloud/network/monitor/LinkStatus.java @@ -0,0 +1,93 @@ +package org.kne.cloud.network.monitor; + +import java.util.function.Consumer; + +/** + * 链路状态监控类,负责维护链路的在线状态和在线率。 + * 设计理念: + * 1. 状态变化时通过回调通知监听者。 + * 2. 在线率采用指数加权移动平均 (EWMA) 算法,平滑且对历史数据有衰减记忆。 + * 3. 自身不启动任何后台线程,状态的更新由外部(例如收到心跳包时)主动触发。 + * 4. 不包含链路名称,名称由外部管理(如 Map),实现关注点分离。 + */ +public class LinkStatus { + // 状态常量 + public static final int DOWN = 0; + public static final int UNSTABLE = 1; + public static final int UP = 2; + + private volatile int state; + private volatile double reliability; // 在线率,范围 [0.0, 1.0] + + private Consumer changeListener; + + // 用于EWMA计算的衰减因子 + private static final double EWMA_ALPHA = 0.9999; + + public LinkStatus() { + this.state = DOWN; + this.reliability = 0.0; + } + + // 状态 getter/setter + public int getState() { + return state; + } + + /** + * 更新链路状态,并在状态真正改变时通知监听器。 + * @param newState 新状态 (DOWN, UNSTABLE, UP) + */ + public void setState(int newState) { + if (this.state == newState) { + return; + } + this.state = newState; + if (changeListener != null) { + changeListener.accept(this); + } + } + + // 在线率 getter + public double getReliability() { + return reliability; + } + + /** + * 核心更新方法:基于当前的在线状态,更新在线率。 + * 此方法应由心跳检测等逻辑周期性调用(例如每秒调用一次)。 + * 使用 EWMA 算法: new_ewma = alpha * old_ewma + (1 - alpha) * current_value + */ + public void updateReliability() { + double currentOnline = (state == UP) ? 1.0 : 0.0; + this.reliability = EWMA_ALPHA * this.reliability + (1 - EWMA_ALPHA) * currentOnline; + } + + /** + * 注册状态变更监听器 + * @param listener 监听器函数 + */ + public void setChangeListener(Consumer listener) { + this.changeListener = listener; + } + + // 静态工具方法 + public static String stateToString(int state) { + switch (state) { + case DOWN: + return "○down"; + case UNSTABLE: + return "●unstable"; + case UP: + return "●up"; + default: + return "unknown"; + } + } + + @Override + public String toString() { + return String.format("[%s]%.2f%%", + stateToString(state), reliability * 100); + } +} \ No newline at end of file diff --git a/src/org/kne/cloud/network/tcp/UDPProtocolRegister.java b/src/org/kne/cloud/network/tcp/UDPProtocolRegister.java new file mode 100644 index 0000000..8da50aa --- /dev/null +++ b/src/org/kne/cloud/network/tcp/UDPProtocolRegister.java @@ -0,0 +1,44 @@ +package org.kne.cloud.network.tcp; + +import java.io.IOException; + +import org.kne.cloud.network.ipv6.IPv6Address; +import org.kne.cloud.network.ipv6.IPv6Packet; +import org.kne.cloud.network.ipv6.IPv6Packet.IPv6Payload; +import org.kne.cloud.network.ipv6.IPv6ProtocolRegister; +import org.kne.cloud.network.klalb.BindableConsumer; +import org.kne.cloud.network.klalb.KLALBController; +import org.kne.cloud.network.klalb.PortBinder; +import org.kne.cloud.network.klalb.PortPacket; +import org.kne.cloud.network.srv6.IPv6PacketConsumer; + +public class UDPProtocolRegister extends PortBinder implements IPv6ProtocolRegister { + private static final boolean showpacket=false; + + private KLALBController controller; + + public UDPProtocolRegister(KLALBController controller) { + super(controller.getSelf().getAddress()); + this.controller = controller; + } + + @Override + public boolean onaccept(IPv6Packet packx) throws IOException { + + IPv6Payload pl = packx.getPayload(); + if (pl instanceof UDPPacket) { + UDPPacket rec = (UDPPacket) pl; + if (showpacket) + System.out.println("UDP_RX:" + rec); + IPv6Address srcA = packx.getSourceAddress(); + PortPacket pt = (PortPacket) rec; + BindableConsumer cons; + if((cons=distributePacketToConsumer(srcA, pt))!=null) { + cons.accept(packx); + return true; + } + } + return false; + } + + } \ No newline at end of file diff --git a/src/org/kne/concurrent/HighPerformanceExecutor2.java b/src/org/kne/concurrent/HighPerformanceExecutor2.java new file mode 100644 index 0000000..f5c0fd4 --- /dev/null +++ b/src/org/kne/concurrent/HighPerformanceExecutor2.java @@ -0,0 +1,134 @@ +package org.kne.concurrent; + +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; +import java.util.Queue; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.Executor; +import java.util.concurrent.Executors; +import java.util.concurrent.LinkedBlockingDeque; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.locks.LockSupport; +import java.util.function.Consumer; + +import org.jctools.queues.MpscArrayQueue; + +public class HighPerformanceExecutor2 implements Executor { + + private ThreadElement[] threads; + + public HighPerformanceExecutor2(int threadcount) { + //this(threadcount,Executors.defaultThreadFactory()); + this(threadcount,Thread.ofVirtual().factory()); + } + public HighPerformanceExecutor2(int threadcount,ThreadFactory th) { + threads=new ThreadElement[threadcount]; + for (int i = 0; i < threadcount; i++) { + ThreadElement t= new ThreadElement(); + th.newThread(t).start(); + threads[i]=t; + } + } + private static class ThreadElement implements Runnable{ + private MpscArrayQueue queue=new MpscArrayQueue(2048); + private AtomicInteger size=new AtomicInteger(); + + private volatile ThreadParker parker=new ThreadParker(); + public Queue getQueue() { + return queue; + } + public int size() { + + return queue.size(); + } + public long prev=System.nanoTime(); + public boolean putTask(Runnable e) { + boolean b=queue.offer(e); + if(b) { + // size.incrementAndGet(); + long curr=System.nanoTime(); + if(curr-prev>1000L||queue.size()>=8) { + prev=curr; + parker.unpark(); + } + } + return b; + + } + @Override + public void run() { + while(true) { + Runnable r=queue.poll(); + if(r!=null) { + // size.decrementAndGet(); + try { + r.run(); + }catch(Throwable e) { + e.printStackTrace(); + } + ThreadYieldCheckpoint.yieldCheckpoint(1000000L); + }else { + parker.parkNanos(1000000L); + + } + } + } + } + @Override + public void execute(Runnable command) { + if(execute0((x)->{command.run();},1000)) { + return; + } + System.out.println("loss!"); + //backup.execute(command); + + + } + + /*private long vl=0; + private boolean execute0(Consumer command,int limit) { + long ord=vl++; + ThreadElement te= threads[(int) (ord%threads.length)]; + boolean b=te.size()>limit; + return te.putTask( ()->{command.accept(b);}); + }*/ + + private boolean execute0(Consumer command,int limit) { + for(int i=0;ilimit; + if(!b) { + if(te.putTask( ()->{command.accept(false);})){ + return true; + } + } + } + + for(int i=0;i{command.accept(true);})){ + return true; + } + } + return false; + } + + public void executeWithCongestionReport(Consumer command) { + if(execute0(command,1000)) { + return; + } + System.out.println("loss!"); + } + public void executeWithCongestionReport(Consumer command,int limit) { + if(execute0(command,limit)) { + return; + } + System.out.println("loss!"); + + + } +}