This commit is contained in:
2026-07-07 00:22:32 +08:00
parent 620afef715
commit 74e4360d71
89 changed files with 1986 additions and 3374 deletions
@@ -78,11 +78,11 @@ public class ArrayListMessageBatcher<T> implements MessageBatcher<T>{
impl.close();
}
public static void main(String[] args) throws InterruptedException {
for(;;){
/*for(;;){
ArrayListMessageBatcher<Integer>mes=new ArrayListMessageBatcher<>(10, 1000000L);
System.gc();
}
/*ArrayListMessageBatcher<Integer>mes=new ArrayListMessageBatcher<>(10, 1000000L);
}*/
ArrayListMessageBatcher<Integer>mes=new ArrayListMessageBatcher<>(10, 1000000L);
mes.setConsumer((x)->{System.out.println(x);});
int n=0;
for(;;){
@@ -91,7 +91,7 @@ public class ArrayListMessageBatcher<T> implements MessageBatcher<T>{
mes.putMessage(n++);
mes.putMessage(n++);
Thread.sleep(1);
}*/
}
}
}
class ArrayListMessageBatcherReleaser<T> extends Releaser<ArrayListMessageBatcher0<T>>{
@@ -114,7 +114,7 @@ class ArrayListMessageBatcherReleaser<T> extends Releaser<ArrayListMessageBatche
private final long maxDelay;
private Consumer<List<T>> consumer;
private final Thread batchThread;
private final AtomicBoolean closed;
private final AtomicBoolean closed = new AtomicBoolean(false);
private long firstTime;
/**
@@ -125,7 +125,6 @@ class ArrayListMessageBatcherReleaser<T> extends Releaser<ArrayListMessageBatche
public ArrayListMessageBatcher0(int batchSize, long maxDelay) {
this.batchSize = batchSize;
this.maxDelay = maxDelay;
this.closed = new AtomicBoolean(false);
// 创建并启动批量处理线程
this.batchThread = ThreadTool.makeVDaemonThread("ArrayListMessageBatcher Thread", this );
@@ -12,9 +12,7 @@ public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetConges
private static long MIN_WINDOW = 16384;
private Consumer<Long> windowControlConsumer;
private BiConsumer<Long, Long> speedControlConsumer;
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
private long MIN_RTTVAR = 5000000L;
private volatile long RTTMin = 1000000000L;
private volatile long RTTVar = 1000000000L;
@@ -26,23 +24,19 @@ public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetConges
private volatile long window=MIN_WINDOW;
private volatile long speed=MIN_SPEED;
private volatile long maxwindow=99999999999L;
private long congressSpeed=MIN_SPEED;
private double MIN_GAIN=1.05;
private double MAX_GAIN=1.20;//2.8853900817779
private double windowGain=MIN_GAIN;
private double speedGain=MAX_GAIN;//1.20 MAX_GAIN
private double burstLimit=1.20;
@Override
public void reset() {
window = MIN_WINDOW;
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
}
speed = MIN_SPEED;
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
windowGain=MAX_GAIN;
}
@@ -51,10 +45,7 @@ public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetConges
this.windowControlConsumer = windowControlConsumer;
}
@Override
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
this.speedControlConsumer = speedControlConsumer;
}
private long RTTstartTime = System.nanoTime();
private long bandwidth;
@@ -65,14 +56,9 @@ public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetConges
} else {
RTTMin = (RTTMin * 99999 + latencyns) / 100000;
}
if (firstUpdate.compareAndSet(true, false)) {
RTTAvg = latencyns;
RTTVar = latencyns / 2;
RTTAvg2 = latencyns;
} else {
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
}
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4);// RTTVar*4
RTTTotal += latencyns;
@@ -100,14 +86,10 @@ public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetConges
} else {
congressSpeed = (congressSpeed * 49 + bandwidth) / 50;
}
speed=(long) (congressSpeed * speedGain)+MIN_SPEED ;
window= ((long) (congressSpeed*windowGain )*(RTTMin+BIAS)/1000000000L)+MIN_WINDOW;
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
}
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
}
}
@@ -124,6 +106,7 @@ public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetConges
@Override
public void setCurrentWindowUsed(long used) {
this.maxwindow = (long) (used*burstLimit) + MIN_WINDOW;
}
@@ -138,7 +121,6 @@ public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetConges
@Override
public void setUpperDelayBound(double upper) {
MAX_GAIN=upper;
speedGain=MAX_GAIN;
}
@Override
@@ -146,5 +128,14 @@ public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetConges
MIN_GAIN=lower;
windowGain=MIN_GAIN;
}
@Override
public void setBurstLimit(double limit) {
burstLimit=limit;
}
@Override
public String getName() {
return "BBR";
}
}
@@ -1,150 +0,0 @@
package org.kne.cloud.network.congestion;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
public class BBRVegasCongressAlgorithm implements CongestionAlgorithm {
private static final long BURST_TIME = 2000000L;
private static final long BIAS = 2000000L;
private static long MIN_SPEED = 128 * 1024L;
private static long MIN_WINDOW = 8192;
private Consumer<Long> windowControlConsumer;
private BiConsumer<Long, Long> speedControlConsumer;
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
private long MIN_RTTVAR = 50000000L;
private volatile long RTTMin = 1000000000L;
private volatile long RTTVar = 1000000000L;
private volatile long RTTAvg = 1000000000L;
private volatile long RTTAvg2 = 1000000000L;
private volatile long RTTTotal = 0;
private volatile long RTTCount = 0;
private volatile long RTO = 1000000000L;
private volatile long window=MIN_WINDOW;
private volatile long speed=MIN_SPEED;
private volatile long maxwindow=99999999999L;
private long congressSpeed=MIN_SPEED;
private double MAX_GAIN=2.8853900817779;//2.8853900817779
private double MIN_GAIN=1.1;
private double gain=MAX_GAIN;
@Override
public void reset() {
window = MIN_WINDOW;
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
}
speed = MIN_SPEED;
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
gain=MAX_GAIN;
}
@Override
public void setWindowControlConsumer(Consumer<Long> windowControlConsumer) {
this.windowControlConsumer = windowControlConsumer;
}
@Override
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
this.speedControlConsumer = speedControlConsumer;
}
private long RTTstartTime = System.nanoTime();
@Override
public synchronized void putAck(long packetSize, long latencyns, boolean ecn) {
if (latencyns <= RTTMin) {
RTTMin = latencyns;
} else {
RTTMin = (RTTMin * 99999 + latencyns) / 100000;
}
if (firstUpdate.compareAndSet(true, false)) {
RTTAvg = latencyns;
RTTVar = latencyns / 2;
RTTAvg2 = latencyns;
} else {
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
}
RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4);// RTTVar*4
RTTTotal += latencyns;
RTTCount++;
long curr=System.nanoTime();
if(curr-RTTstartTime>RTTMin) {
if(RTTCount>0) {
RTTAvg2=(RTTAvg2+RTTTotal/RTTCount)/2;
RTTTotal=0;
RTTCount=0;
}
long RTTAvg2x=Math.max(RTTAvg2-BIAS,RTTMin);
double excepted=window/(double)RTTMin;
double actual=window/(double)RTTAvg2x;
double diff=(excepted-actual)*RTTMin;
if(diff>65536*3+3) {
gain-=0.0005;
}
if(diff<65536*3+1) {
gain+=0.0005;
}
if(gain>=MAX_GAIN){
gain=MAX_GAIN;
//System.out.println("window:"+window);
}
if(gain<MIN_GAIN) {
gain=MIN_GAIN;
//System.out.println("window:"+window);
}
RTTstartTime=curr;
}
//System.out.println(speed+" "+window);
}
@Override
public void putLoss(long packetSize, long lossns, int losscounter) {
}
@Override
public long getRTO() {
return RTO;
}
@Override
public void setCurrentWindowUsed(long used) {
}
@Override
public void setCurrentBandwidth(long bandwidth) {
/*if (bandwidth >= congressSpeed) {
congressSpeed = (congressSpeed * 49 + bandwidth) / 50;
} else {
congressSpeed = (congressSpeed *2 + bandwidth) / 3;
}*/
if (bandwidth >= congressSpeed) {
//congressSpeed = bandwidth;
congressSpeed = (congressSpeed + bandwidth) / 2;
} else {
congressSpeed = (congressSpeed * 49 + bandwidth) / 50;
}
speed=Math.max(MIN_SPEED,(long) (congressSpeed * gain)) ;
window=Math.max( (Math.max(MIN_SPEED,(long) (congressSpeed*gain ))*Math.max(2000000L,RTTMin)/1000000000L),MIN_WINDOW);
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
}
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
//System.out.println(gain);
//System.out.println(bandwidth+" "+congressSpeed+" "+window);
}
}
@@ -6,11 +6,14 @@ import java.util.function.Consumer;
public interface CongestionAlgorithm {
public void reset();
public void setWindowControlConsumer(Consumer<Long>windowControlConsumer);
public void setSpeedControlConsumer(BiConsumer<Long,Long>speedControlConsumer);
public void putAck(long packetSize,long latencyns,boolean ecn);
public void putAck(long packetSize,long RTTns,boolean ecn);
default public void putAck(long packetSize,long RTTns,long OWDup,boolean ecn) {
putAck(packetSize,RTTns,ecn);
}
public void putLoss(long packetSize,long lossns,int losscounter);
public void setCurrentWindowUsed(long maxwindow);
public void setCurrentBandwidth(long bandwidth);
public long getRTO();
public String getName();
}
@@ -2,43 +2,37 @@ package org.kne.cloud.network.congestion;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.LongAdder;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
import org.jctools.counters.FixedSizeStripedLongCounter;
import org.kne.cloud.network.klalb.KLALBUtils;
public class DCTCP2CongestionAlgorithm implements CongestionAlgorithm {
private static final long BURST_TIME = 1000000L;
private static final long BIAS = 2000000L;
private static long MIN_SPEED = 64 * 1024L;
private static long MIN_WINDOW = 32768;
private static long MIN_WINDOW = 8192;
private Consumer<Long> windowControlConsumer;
private BiConsumer<Long, Long> speedControlConsumer;
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
private long MIN_RTTVAR = 400000000L;
private volatile long RTTMin = 1000000000L;
private long MIN_RTTVAR = 100000000L;
private volatile long RTTVar = 1000000000L;
private volatile long RTTAvg = 1000000000L;
private double ECNAvg2 = 0L;
private double alpha=0;
private volatile AtomicLong ECNSize = new AtomicLong(0);
private volatile AtomicLong TotalSize = new AtomicLong(0);
private volatile long RTO = 1000000000L;
private volatile FixedSizeStripedLongCounter ECNSize = KLALBUtils.createCounter();
private volatile FixedSizeStripedLongCounter TotalSize = KLALBUtils.createCounter();
private volatile long RTO = 500000000L;
private volatile long window=MIN_WINDOW;
private volatile long windowUsed=MIN_WINDOW;
private volatile long maxwindow = 99999999999L;
private volatile long congressWindowSize=MIN_WINDOW;
private volatile long speed=MIN_SPEED;
private long congressSpeed=MIN_SPEED;
private long RTTstartTime = System.nanoTime();
private long bandwidth=MIN_SPEED;
private double MAX_GAIN=1.5;//2.8853900817779
private double MIN_GAIN=1.2;
private double gain=MAX_GAIN;
@Override
public void reset() {
@@ -46,11 +40,6 @@ public class DCTCP2CongestionAlgorithm implements CongestionAlgorithm {
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
}
speed = MIN_SPEED;
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
gain=MAX_GAIN;
}
@Override
@@ -58,73 +47,44 @@ public class DCTCP2CongestionAlgorithm implements CongestionAlgorithm {
this.windowControlConsumer = windowControlConsumer;
}
@Override
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
this.speedControlConsumer = speedControlConsumer;
}
@Override
public synchronized void putAck(long packetSize, long latencyns, boolean ecn) {
if (latencyns <= RTTMin) {
RTTMin = latencyns;
} else {
RTTMin = (RTTMin * 99999 + latencyns) / 100000;
}
if (firstUpdate.compareAndSet(true, false)) {
RTTAvg = latencyns;
RTTVar = latencyns / 2;
} else {
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
}
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4);// RTTVar*4
if(ecn) {
ECNSize.addAndGet(packetSize);
ECNSize.inc(packetSize);
}
TotalSize.addAndGet(packetSize);
TotalSize.inc(packetSize);
long curr=System.nanoTime();
if(curr-RTTstartTime>(RTTAvg+BIAS)) {
RTTstartTime=curr;
long totalSize=TotalSize.getAndSet(0);
long ecnSize=ECNSize.getAndSet(0);
long totalSize=TotalSize.getAndReset();
long ecnSize=ECNSize.getAndReset();
if(totalSize!=0) {
ECNAvg2=ecnSize/(double)totalSize;
alpha=(alpha*15+ECNAvg2)/16;
alpha=(alpha*7+ECNAvg2)/8;
}
if(ECNAvg2>0.5) {
gain=MIN_GAIN;
if(congressWindowSize>MIN_WINDOW) {
congressWindowSize-=200;
updateWindowSize();
}
}else if(ECNAvg2>0.1){
if(windowUsed*3L>=congressWindowSize) {
}else if(ECNAvg2>0.01){
congressWindowSize+=200;
updateWindowSize();
}
}else {
if(windowUsed*3L>=congressWindowSize) {
congressWindowSize+=8000;
updateWindowSize();
}
congressWindowSize+=4000;
}
//System.out.println(ECNAvg2);
//System.out.println(speed+" "+window);
if (bandwidth >= congressSpeed) {
congressSpeed = bandwidth;
} else {
congressSpeed = (congressSpeed * 99 + bandwidth) / 100;
}
speed=Math.max(MIN_SPEED,(long) (congressSpeed * gain)) ;
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
updateWindowSize();
}
@@ -132,6 +92,11 @@ public class DCTCP2CongestionAlgorithm implements CongestionAlgorithm {
}
private void updateWindowSize() {
if(congressWindowSize<MIN_WINDOW) {
congressWindowSize=MIN_WINDOW;
}else if(congressWindowSize>maxwindow) {
congressWindowSize=maxwindow;
}
window=Math.max( congressWindowSize,MIN_WINDOW);
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
@@ -150,12 +115,17 @@ public class DCTCP2CongestionAlgorithm implements CongestionAlgorithm {
@Override
public void setCurrentWindowUsed(long used) {
this.windowUsed=used;
this.maxwindow=(long) (used*1.5)+MIN_WINDOW;
}
@Override
public void setCurrentBandwidth(long bandwidth) {
this.bandwidth=bandwidth;
}
@Override
public String getName() {
return "DCTCP2";
}
}
@@ -11,13 +11,12 @@ public class DCTCPCongestionAlgorithm implements CongestionAlgorithm {
private static final long BURST_TIME = 1000000L;
private static final long BIAS = 2000000L;
private static long MIN_SPEED = 64 * 1024L;
private static long MIN_WINDOW = 32768;
private static long MIN_WINDOW = 8192;
private Consumer<Long> windowControlConsumer;
private BiConsumer<Long, Long> speedControlConsumer;
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
private long MIN_RTTVAR = 400000000L;
private long MIN_RTTVAR = 100000000L;
private volatile long RTTMin = 1000000000L;
private volatile long RTTVar = 1000000000L;
private volatile long RTTAvg = 1000000000L;
@@ -31,14 +30,10 @@ public class DCTCPCongestionAlgorithm implements CongestionAlgorithm {
private volatile long window=MIN_WINDOW;
private volatile long windowUsed=MIN_WINDOW;
private volatile long congressWindowSize=MIN_WINDOW;
private volatile long speed=MIN_SPEED;
private long congressSpeed=MIN_SPEED;
private long RTTstartTime = System.nanoTime();
private long bandwidth=MIN_SPEED;
private double MAX_GAIN=2.8853900817779;//2.8853900817779
private double MIN_GAIN=1.2;
private double gain=MAX_GAIN;
@Override
public void reset() {
@@ -46,11 +41,6 @@ public class DCTCPCongestionAlgorithm implements CongestionAlgorithm {
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
}
speed = MIN_SPEED;
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
gain=MAX_GAIN;
}
@Override
@@ -58,10 +48,6 @@ public class DCTCPCongestionAlgorithm implements CongestionAlgorithm {
this.windowControlConsumer = windowControlConsumer;
}
@Override
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
this.speedControlConsumer = speedControlConsumer;
}
@Override
public synchronized void putAck(long packetSize, long latencyns, boolean ecn) {
@@ -86,7 +72,7 @@ public class DCTCPCongestionAlgorithm implements CongestionAlgorithm {
long curr=System.nanoTime();
if(curr-RTTstartTime>(RTTAvg+BIAS)) {
if(curr-RTTstartTime>Math.max(BIAS, 2*RTTMin)) {
RTTstartTime=curr;
long totalSize=TotalSize.getAndSet(0);
@@ -94,37 +80,23 @@ public class DCTCPCongestionAlgorithm implements CongestionAlgorithm {
if(totalSize!=0) {
ECNAvg2=ecnSize/(double)totalSize;
alpha=(alpha*15+ECNAvg2)/16;
//System.out.println(ECNAvg2);
}
if(ECNAvg2>0.6) {
gain=MIN_GAIN;
if(ECNAvg2>0.5) {
if(congressWindowSize>MIN_WINDOW) {
congressWindowSize=Math.max(MIN_WINDOW,(long) (congressWindowSize*(1-alpha/10)));
updateWindowSize();
}
}else if(ECNAvg2>0.2){
if(windowUsed*3L>=congressWindowSize) {
congressWindowSize+=200;
updateWindowSize();
}
}else {
if(windowUsed*3L>=congressWindowSize) {
congressWindowSize+=8000;
congressWindowSize+=4000;
updateWindowSize();
}
}
//System.out.println(ECNAvg2);
//System.out.println(speed+" "+window);
if (bandwidth >= congressSpeed) {
congressSpeed = bandwidth;
} else {
congressSpeed = (congressSpeed * 99 + bandwidth) / 100;
}
speed=Math.max(MIN_SPEED,(long) (congressSpeed * gain)) ;
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
}
@@ -155,7 +127,15 @@ public class DCTCPCongestionAlgorithm implements CongestionAlgorithm {
@Override
public void setCurrentBandwidth(long bandwidth) {
this.bandwidth=bandwidth;
// TODO 自动生成的方法存根
}
@Override
public String getName() {
return "DCTCP";
}
}
@@ -3,4 +3,5 @@ package org.kne.cloud.network.congestion;
public interface DetnetCongestionAlgorithm extends CongestionAlgorithm {
public void setUpperDelayBound(double upper);
public void setLowerDelayBound(double lower);
public void setBurstLimit(double limit);
}
@@ -3,7 +3,7 @@ package org.kne.cloud.network.congestion;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
public class EmptyCongestionAlgorithm implements CongestionAlgorithm{
public class NOCongestionAlgorithm implements CongestionAlgorithm{
private long timeout;
@@ -23,11 +23,7 @@ public class EmptyCongestionAlgorithm implements CongestionAlgorithm{
}
@Override
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
// TODO 自动生成的方法存根
}
@Override
public void putAck(long packetSize, long latencyns, boolean ecn) {
@@ -53,7 +49,7 @@ public class EmptyCongestionAlgorithm implements CongestionAlgorithm{
}
public EmptyCongestionAlgorithm(long timeout) {
public NOCongestionAlgorithm(long timeout) {
super();
this.timeout = timeout;
}
@@ -63,4 +59,9 @@ public class EmptyCongestionAlgorithm implements CongestionAlgorithm{
return timeout;
}
@Override
public String getName() {
return "NO";
}
}
@@ -1,45 +0,0 @@
package org.kne.cloud.network.congestion;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReferenceArray;
public class ReceiveByteSlidingWindow<E> {
private AtomicReferenceArray<E> array;
private AtomicInteger windowSize =new AtomicInteger();
private AtomicLong windowPosition=new AtomicLong();
public int getWindowSize() {
return windowSize.get();
}
public void setWindowSize(int windowSize) {
this.windowSize.set(windowSize);
}
public long getWindowPosition() {
return windowPosition.get();
}
public void setWindowPosition(long windowPosition) {
this.windowPosition.set(windowPosition);
}
public ReceiveByteSlidingWindow(int capacity){
array=new AtomicReferenceArray<E>(capacity);
}
public boolean push(E object) {
Objects.requireNonNull(object);
if(array.get((int) ((windowPosition.get()-windowSize.get())%array.length()))==null ) {
array.set((int) (windowPosition.getAndIncrement()%array.length()), object);
return true;
}
return false;
}
public E pull(long number) {
return array.getAndSet((int) (number%array.length()), null);
}
}
@@ -1,96 +0,0 @@
package org.kne.cloud.network.congestion;
import java.util.Iterator;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReferenceArray;
public class SendByteSlidingWindow<E> implements Iterable<E>{
private Object[] array;
private volatile int windowSize ;
private volatile long windowPosition;
public int getWindowSize() {
return windowSize;
}
public void setWindowSize(int windowSize) {
this.windowSize=windowSize;
}
public long getWindowPosition() {
return windowPosition;
}
public void setWindowPosition(long windowPosition) {
this.windowPosition=windowPosition;
}
public SendByteSlidingWindow(int capacity,int windowSize){
if(windowSize>capacity) {
throw new IllegalArgumentException("windowSize>capacity!");
}
array= new Object[capacity];
this.windowSize=windowSize;
}
public boolean checkPush() {
if(array[calcPosition(windowPosition-windowSize)]==null ) {
return true;
}
return false;
}
public void push(E object) {
array[calcPosition(windowPosition)]= object;
}
public E remove(long number) {
if(number<windowPosition&&number>=windowPosition-windowSize) {
int indp=calcPosition(number);
E old=(E) array[indp];
array[indp]=null;
return old;
}else {
return null;
}
}
public boolean isEmpty() {
for (long i = windowPosition-windowSize; i < windowPosition; i++) {
if(array[calcPosition(i)]!=null) {
return false;
}
}
return true;
}
public E get(long number) {
if(number<windowPosition&&number>=windowPosition-windowSize) {
return (E) array[calcPosition(number)];
}else {
return null;
}
}
@Override
public Iterator<E> iterator() {
return new Iterator<E>() {
private long pointer=windowPosition-windowSize;
@Override
public boolean hasNext() {
return pointer<windowPosition;
}
@Override
public E next() {
return (E) array[calcPosition(pointer++)];
}
};
}
private int calcPosition(long number) {
return (int) ((number+array.length)%array.length);
}
}
@@ -1,30 +1,27 @@
package org.kne.cloud.network.congestion;
import java.io.Closeable;
import java.io.IOException;
import java.net.SocketException;
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.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.SendItem;
import org.kne.concurrent.ThreadParker;
public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Closeable, AutoCloseable {
private Map<K, SendItem<V>> sendMap = new ConcurrentHashMap<>(1024, 0.2f);
private AtomicLong sendmapWindowUsed=new AtomicLong(0);
//private LongAdder sendmapWindowUsed = new LongAdder();
private ThreadParker tp=new ThreadParker();
private Map<K, SendItem<V>> sendMap = new ConcurrentHashMap<>(2048);
private AtomicLong sendWindowUsed=new AtomicLong(0);
//private LongAdder sendWindowUsed = new LongAdder();
private int headerCalibrate = 0;
@@ -36,14 +33,14 @@ public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Clos
private Consumer<V> resendConsumer;
private Consumer<V> closingConsumer;
private volatile boolean closed = false;
private ResendChecker checker = new ResendChecker();
private AtomicLong windowSize = new AtomicLong();
private Lock lock;
private Condition condition;
private boolean removeAtResend;
private volatile long windowSize =0;
private volatile boolean closed = false;
public SendPacketSlidingWindow(CongestionAlgorithm algorithm, long windowSize) {
this(algorithm, windowSize, 0, true);
@@ -58,17 +55,17 @@ public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Clos
this.headerCalibrate = headerCalibrate;
this.removeAtResend = removeAtResend;
this.algorithm = algorithm;
this.windowSize.set(windowSize);
this.windowSize=windowSize;
Thread t = ThreadTool.makeVThread("重路由计时器线程", checker);
t.start();
}
public long getWindowSize() {
return windowSize.get();
return windowSize;
}
public void setWindowSize(long windowSize) {
this.windowSize.set(windowSize);
this.windowSize=windowSize;
}
public Consumer<V> getResendConsumer() {
@@ -87,8 +84,8 @@ public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Clos
this.closingConsumer = closingConsumer;
}
public long getSendmapWindowUsed() {
return sendmapWindowUsed.get();
public long getSendWindowUsed() {
return sendWindowUsed.get();
}
public int getHeaderCalibrate() {
@@ -99,20 +96,38 @@ public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Clos
return algorithm;
}
public boolean isCongress(IPv6Packet iPv6Packet, double scale) {
boolean congress = sendmapWindowUsed.get() > windowSize.get() * scale;
public boolean isCongestion( double scale) {
boolean congress = sendWindowUsed.get() > windowSize * scale;
return congress;
}
public boolean isCongestion() {
boolean congress = sendWindowUsed.get() > windowSize ;
return congress;
}
public void waitForAvaliable() throws SocketException {
while(isCongestion()) {
if(closed) {
throw new SocketException("Send window closed!");
}
tp.parkNanos(1000000L);
}
}
public void putBlocked(K sequence, V packet) throws SocketException {
waitForAvaliable();
put(sequence,packet);
}
public void put(K sequence, V packet) {
SendItem<V> newitem = new SendItem<V>(packet);
SendItem<V> newitem = new SendItem<V>(packet,headerCalibrate);
SendItem<V> old = sendMap.put(sequence, newitem);
long dx=0;
if (old != null) {
dx=-(old.getPacketLength() + headerCalibrate);
dx=-old.getPacketLength() ;
}
long newWindow =sendmapWindowUsed.addAndGet(newitem.getPacketLength() + headerCalibrate+dx);
long newWindow =sendWindowUsed.addAndGet(newitem.getPacketLength() +dx);
algorithm.setCurrentWindowUsed(newWindow);
}
public V ack(K sequence) {
@@ -122,10 +137,11 @@ public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Clos
public V ack(K sequence,boolean ecn) {
SendItem<V> ipv6;
if ((ipv6 = sendMap.remove(sequence)) != null) {
long newWindow =sendmapWindowUsed.addAndGet(-(ipv6.getPacketLength() + headerCalibrate));
long newWindow =sendWindowUsed.addAndGet(-ipv6.getPacketLength() );
algorithm.setCurrentWindowUsed(newWindow);
long RTTC = System.nanoTime() - ipv6.getSendtime();
algorithm.putAck(ipv6.getPacketLength(), RTTC, ecn);
tp.unpark();
// System.out.println("remove:"+aseq.getSequence()+" "+sendMap.size());
Lock lockx = lock;
if (lockx != null) {
@@ -146,6 +162,34 @@ public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Clos
}
}
public V ack(K sequence,boolean ecn,long OWDdown) {
SendItem<V> ipv6;
if ((ipv6 = sendMap.remove(sequence)) != null) {
long newWindow =sendWindowUsed.addAndGet(-ipv6.getPacketLength() );
algorithm.setCurrentWindowUsed(newWindow);
long RTTC = System.nanoTime() - ipv6.getSendtime();
algorithm.putAck(ipv6.getPacketLength(), RTTC,OWDdown, ecn);
tp.unpark();
// System.out.println("remove:"+aseq.getSequence()+" "+sendMap.size());
Lock lockx = lock;
if (lockx != null) {
lockx.lock();
try {
Condition cds = condition;
if (cds != null) {
cds.signalAll();
}
} finally {
lockx.unlock();
}
}
return ipv6.getPacket();
} else {
// System.out.println("miss:"+aseq.getSequence()+" "+sendMap.size());
return null;
}
}
private class ResendChecker implements Runnable {
@Override
@@ -159,8 +203,9 @@ public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Clos
if (current - sitm.getSendtime() > algorithm.getRTO()) {
if (removeAtResend) {
iterator.remove();
long newWindow = sendmapWindowUsed.addAndGet((int) -(sitm.getPacketLength() + headerCalibrate));
long newWindow = sendWindowUsed.addAndGet((int) -sitm.getPacketLength() );
algorithm.setCurrentWindowUsed(newWindow);
tp.unpark();
} else {
sitm.setSendtime(current);
}
@@ -183,7 +228,7 @@ public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Clos
}
}
sendMap.clear();
sendmapWindowUsed.set(0);
sendWindowUsed.set(0);
}
}
@@ -215,7 +260,20 @@ public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Clos
@Override
public String toString() {
return "SendPacketSlidingWindow [sendmapWindowUsed=" + sendmapWindowUsed + ", windowSize=" + windowSize + "]";
return "SendPacketSlidingWindow [sendWindowUsed=" + sendWindowUsed + ", windowSize=" + windowSize + ", closed="
+ closed + "]";
}
public void waitForEmpty() throws SocketException {
while(!isEmpty()) {
if(isClosed()) {
throw new SocketException("Send window closed!");
}
tp.parkNanos(1000000L);
}
}
}
@@ -9,19 +9,15 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong
// 可配置的时延膨胀比上下限 - Vegas2.0核心参数
private double delayUpperBound = 1.20; // RTT上限为基础RTT的1.2倍
private double delayLowerBound = 1.15; // RTT下限为基础RTT的1.05倍
private double burstLimit=1.20;
private static final long BURST_TIME = 2000000L;
private static final long BIAS = 500000L;
private static long MIN_SPEED = 256 * 1024L;
private static final long BIAS = 1000000L;
private static long MIN_WINDOW = 16384;
// 新增:窗口调整参数(可基于当前窗口大小动态调整)
private static final double WINDOW_ADJUST_FACTOR = 0.01; // 1%的窗口调整幅度
private Consumer<Long> windowControlConsumer;
private BiConsumer<Long, Long> speedControlConsumer;
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
private long MIN_RTTVAR = 30000000L;
private volatile long RTTMin = 1000000000L; // BaseRTT
private volatile long RTTVar = 1000000000L;
@@ -31,12 +27,15 @@ 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;
private volatile long speed = MIN_SPEED;
private volatile long maxwindow = 99999999999L;
private long congressSpeed=MIN_SPEED;
public Vegas2CongestionAlgorithm() {
reset();
}
@@ -47,10 +46,7 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
}
speed = MIN_SPEED;
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
}
@Override
@@ -58,15 +54,17 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong
this.windowControlConsumer = windowControlConsumer;
}
@Override
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
this.speedControlConsumer = speedControlConsumer;
}
private long RTTstartTime = System.nanoTime();
@Override
public synchronized void putAck(long packetSize, long latencyns, boolean ecn) {
@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时,更激进地减少窗口
@@ -77,35 +75,128 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong
}
// 更新最小RTTBaseRTT
if (latencyns <= RTTMin) {
RTTMin = latencyns;
if (RTTns <= RTTMin) {
RTTMin = RTTns;
} else {
// 缓慢适应:当网络路径真正变化时,BaseRTT应能缓慢更新
// 这里使用极慢的衰减因子,只有在持续观测到更低RTT时才快速更新
RTTMin = (RTTMin * 9999 + latencyns) / 10000;
RTTMin = (RTTMin * 999 + RTTns) / 1000;
}
// 更新最小OWDBaseOWD
if (OWDup <= OWDMin) {
OWDMin = OWDup;
} else {
// 缓慢适应:当网络路径真正变化时,BaseRTT应能缓慢更新
// 这里使用极慢的衰减因子,只有在持续观测到更低RTT时才快速更新
OWDMin = (OWDMin * 999 + OWDup) / 1000;
}
// 更新平滑RTT估计
if (firstUpdate.compareAndSet(true, false)) {
RTTAvg = latencyns;
RTTVar = latencyns / 2;
RTTAvg2 = latencyns;
} else {
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
}
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - RTTns)) / 4;
RTTAvg = (RTTAvg * 7 + RTTns) / 8;
RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4);
RTTTotal += latencyns;
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;
}
}
@Override
public void putAck(long packetSize, long RTTns, 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);
}
}
// 更新最小RTTBaseRTT
if (RTTns <= RTTMin) {
RTTMin = RTTns;
} else {
// 缓慢适应:当网络路径真正变化时,BaseRTT应能缓慢更新
// 这里使用极慢的衰减因子,只有在持续观测到更低RTT时才快速更新
RTTMin = (RTTMin * 999 + RTTns) / 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++;
long curr = System.nanoTime();
// 每个RTT周期调整一次窗口(使用当前估计的RTT)
if (curr - RTTstartTime > Math.max(BIAS, RTTMin)) {
if (RTTCount > 0) {
long currentRTT = RTTTotal / RTTCount;
RTTAvg2 = (RTTAvg2 * 3 + currentRTT) / 4; // 平滑当前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;
}
@@ -113,64 +204,47 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong
// ================== Vegas2.0核心逻辑 ==================
// 1. 计算时延膨胀比
double delayRatio = (double) RTTAvg2 / (double) Math.max(BIAS, RTTMin);
//System.out.println(delayRatio+" "+RTTAvg2/1000000.0+" "+RTTMin/1000000.0);
// 2. 基于比值的窗口调整(取代原来的基于差值的调整)
long adjustStep = Math.max(512, (long)(window2 * WINDOW_ADJUST_FACTOR));
if (delayRatio > delayUpperBound) {
// 时延过高:减小窗口,减少幅度与超标程度成正比
double exceedRatio = delayRatio / delayUpperBound;
window2 -= (long)(adjustStep * exceedRatio);
window2 -= 4096;
} else if (delayRatio < delayLowerBound) {
// 时延过低:增大窗口,增加幅度与低于目标程度成正比
double belowRatio = delayLowerBound / delayRatio;
window2 += (long)(adjustStep * belowRatio);
window2 += 4096;
} else {
// 在理想区间内:微调以维持稳定
// 计算目标RTT = BaseRTT × 目标膨胀比
long targetRTT = (long)(Math.max(BIAS, RTTMin) * (delayLowerBound + delayUpperBound) / 2.0);
if (RTTAvg2 > targetRTT) {
// 略高于目标:小幅减少
window2 = Math.max(0, window2 - 256);
} else if (RTTAvg2 < targetRTT) {
// 略低于目标:小幅增加
window2 = Math.min(maxwindow, window2 + 256);
}
// 非常接近目标:保持窗口不变
}
// ===================================================
// 窗口边界检查
if (window2 >= maxwindow) {
window2 = maxwindow;
}
if (window2 < 0) {
window2 = 0;
}
updateWindowSize();
RTTstartTime = curr;
if (bandwidth >= congressSpeed) {
congressSpeed = bandwidth;
//congressSpeed = (congressSpeed + bandwidth) / 2;
} else {
congressSpeed = (congressSpeed * 49 + bandwidth) / 50;
}
// 速度计算:基于当前窗口和RTT,考虑一个安全边界
// 使用目标膨胀比而非固定1.2倍,确保与窗口控制逻辑一致
window=window2+MIN_WINDOW;
speed=(long) (congressSpeed * delayUpperBound)+MIN_SPEED ;
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
}
if (speedControlConsumer != null) {
speedControlConsumer.accept(speed, BURST_TIME);
}
}
}
private void updateWindowSize() {
// 窗口边界检查
if (window2 >= maxwindow) {
window2 = maxwindow;
}
if (window2 < 0) {
window2 = 0;
}
// 速度计算:基于当前窗口和RTT,考虑一个安全边界
// 使用目标膨胀比而非固定1.2倍,确保与窗口控制逻辑一致
window=window2+MIN_WINDOW;
if (windowControlConsumer != null) {
windowControlConsumer.accept(window);
}
}
@Override
public void putLoss(long packetSize, long lossns, int losscounter) {
@@ -183,14 +257,12 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong
@Override
public void setCurrentWindowUsed(long used) {
this.maxwindow = used * 16 + 65536;
this.maxwindow = (long) (used * burstLimit )+ MIN_WINDOW;
}
private long bandwidth;
@Override
public void setCurrentBandwidth(long bandwidth) {
// 可选:基于已知带宽信息调整参数
this.bandwidth=bandwidth;
}
// 新增方法:获取当前状态信息(用于监控和调试)
@@ -215,4 +287,14 @@ public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCong
public void setLowerDelayBound(double lower) {
delayLowerBound=lower;
}
@Override
public void setBurstLimit(double limit) {
burstLimit=limit;
}
@Override
public String getName() {
return "Vegas2";
}
}