This commit is contained in:
Administrator
2024-08-18 17:47:11 +08:00
parent e4f4c77c19
commit e243203535
121 changed files with 7499 additions and 2169 deletions
@@ -1,6 +1,8 @@
package org.kne.cloud.network.klalb;
import java.io.File;
import java.io.FileNotFoundException;
import java.io.FileOutputStream;
import java.io.IOException;
import java.io.PrintStream;
import java.io.PrintWriter;
@@ -10,30 +12,62 @@ import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.net.SocketTimeoutException;
import java.nio.channels.UnresolvedAddressException;
import java.time.Clock;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List;
import java.util.PriorityQueue;
import java.util.UUID;
import java.util.Vector;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.locks.LockSupport;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
import java.util.function.Supplier;
import javax.management.monitor.Monitor;
import org.kne.cloud.clock.AdjustedNanoClock;
import org.kne.cloud.clock.ExponentialBackoffTimeClock;
import org.kne.cloud.clock.ReliabilityBackoffTimeClock;
import org.kne.cloud.network.MultipurposeSocketAddress;
import org.kne.cloud.network.SpeedLimiter;
import org.kne.cloud.network.ThreadTool;
import org.kne.cloud.network.monitor.MonitorData;
import org.kne.cloud.network.monitor.SpeedAndTrafficAndDelayMonitorDataImpl;
import org.kne.cloud.network.monitor.SpeedAndTrafficMonitorDataImpl;
import org.kne.debug.TimeDebugger;
public class KLALBRemoteLine implements Comparable<KLALBRemoteLine>{
public class KLALBRemoteLine implements Comparable<KLALBRemoteLine> {
private static final boolean debug = true;
private static final boolean showpacket = false;
private ReliabilityBackoffTimeClock coll = new ReliabilityBackoffTimeClock();
private volatile boolean closed = false;
private volatile boolean closed=false;
private volatile Supplier<Inet6Address> localVaddrSupplier;
private volatile KLALBController klalbController;
public KLALBController getKlalbController() {
return klalbController;
}
public void setKlalbController(KLALBController klalbController) {
this.klalbController = klalbController;
if (klalbController != null && remoteVaddr != null) {
adjnc = klalbController.getAdjustedClockByVaddr(remoteVaddr);
}
}
public Supplier<Inet6Address> getLocalVaddrSupplier() {
return localVaddrSupplier;
}
@@ -41,45 +75,44 @@ public class KLALBRemoteLine implements Comparable<KLALBRemoteLine>{
public void setLocalVaddrSupplier(Supplier<Inet6Address> localVaddrSupplier) {
this.localVaddrSupplier = localVaddrSupplier;
}
private volatile Inet6Address remoteVaddr;
public Inet6Address getRemoteVaddr() throws SocketTimeoutException {
public Inet6Address getRemoteVaddr() {
return remoteVaddr;
}
private volatile AdjustedNanoClock adjnc = new AdjustedNanoClock();
public void waitForRemoteVaddrAvaliable(long timeout) throws SocketTimeoutException {
long s=System.nanoTime();
while(remoteVaddr==null) {
long s = System.nanoTime();
while (remoteVaddr == null) {
try {
Thread.sleep(1);
} catch (InterruptedException e) {
e.printStackTrace();
}
if(System.nanoTime()-s>timeout*1000000)
if (System.nanoTime() - s > timeout * 1000000)
throw new SocketTimeoutException();
}
}
public Monitor getMonitor() {
public SpeedAndTrafficAndDelayMonitorDataImpl getMonitor() {
return monitor;
}
public String toString() {
return remoteVaddr+"\t"+monitor.toString()+"\t"+sumNextPacketCount();
return remoteVaddr + "\t" + monitor.toString();
}
public String toString2() {
return remoteVaddr+"\t"+monitor.toString2();
}
private Monitor monitor;
private SpeedAndTrafficAndDelayMonitorDataImpl monitor;
private volatile KLALBPacketLink kplink;
private MultipurposeSocketAddress bindAddress;
private MultipurposeSocketAddress socketAddress;
private Thread lthd;
public void setKplink(KLALBPacketLink kplink) {
this.kplink = kplink;
}
@@ -88,415 +121,554 @@ public class KLALBRemoteLine implements Comparable<KLALBRemoteLine>{
return kplink;
}
public MultipurposeSocketAddress getSocketAddress() {
return socketAddress;
}
public MultipurposeSocketAddress getBindAddress() {
return bindAddress;
}
protected static KLALBPacketLink createLink(MultipurposeSocketAddress bindAddress,MultipurposeSocketAddress mpa) throws IOException {
if(mpa.isStream()) {
if(bindAddress!=null) {
return new StreamKLALBPacketLink(mpa.connectSocket(InetAddress.getByName( bindAddress.getHost()),bindAddress.getPort()));
}else {
return new StreamKLALBPacketLink(mpa.connectSocket());
protected static KLALBPacketLink createLink(MultipurposeSocketAddress bindAddress, MultipurposeSocketAddress mpa)
throws IOException {
if (mpa.isStream()) {
if (mpa.supportNIO()) {
if (bindAddress != null) {
return new StreamChannelKLALBPacketLink(mpa
.connectSocketChannel(InetAddress.getByName(bindAddress.getHost()), bindAddress.getPort()));
} else {
return new StreamChannelKLALBPacketLink(mpa.connectSocketChannel());
}
} else {
if (bindAddress != null) {
return new StreamKLALBPacketLink(
mpa.connectSocket(InetAddress.getByName(bindAddress.getHost()), bindAddress.getPort()));
} else {
return new StreamKLALBPacketLink(mpa.connectSocket());
}
}
} else {
if (bindAddress != null) {
return new SplitedDatagramKLALBPacketLink(
mpa.connectDatagramSocket(InetAddress.getByName(bindAddress.getHost()), bindAddress.getPort()));
} else {
return new SplitedDatagramKLALBPacketLink(mpa.connectDatagramSocket());
}
}
else {
if(bindAddress!=null) {
return new DatagramKLALBPacketLink(mpa.connectDatagramSocket(InetAddress.getByName( bindAddress.getHost()),bindAddress.getPort()));
}else {
return new DatagramKLALBPacketLink(mpa.connectDatagramSocket());
}
}
}
public KLALBRemoteLine(MultipurposeSocketAddress mpa) {
this(mpa,new Monitor());
}
public KLALBRemoteLine(MultipurposeSocketAddress mpa,Monitor monitor) {
this(mpa,null,monitor);
}
public KLALBRemoteLine(KLALBPacketLink kpl) {
this(kpl,new Monitor());
}
public KLALBRemoteLine(KLALBPacketLink kpl,Monitor monitor) {
this.kplink=kpl;
this.monitor=monitor;
monitor.setName(kpl.toString());
}
public KLALBRemoteLine(MultipurposeSocketAddress mpa,MultipurposeSocketAddress bindaddr,Monitor monitor ) {
this.socketAddress=mpa;
this.bindAddress=bindaddr;
this.monitor=monitor;
monitor.setName(mpa.toString());
public KLALBRemoteLine(MultipurposeSocketAddress mpa) {
this(mpa, new SpeedAndTrafficAndDelayMonitorDataImpl());
}
public KLALBRemoteLine(MultipurposeSocketAddress mpa, SpeedAndTrafficAndDelayMonitorDataImpl monitor) {
this(mpa, null, monitor);
}
public KLALBRemoteLine(KLALBPacketLink kpl) {
this(kpl, new SpeedAndTrafficAndDelayMonitorDataImpl());
}
public KLALBRemoteLine(KLALBPacketLink kpl, SpeedAndTrafficAndDelayMonitorDataImpl monitor) {
this.kplink = kpl;
this.monitor = monitor;
monitor.setName(kpl.toString());
}
public KLALBRemoteLine(MultipurposeSocketAddress mpa, MultipurposeSocketAddress bindaddr,
SpeedAndTrafficAndDelayMonitorDataImpl monitor) {
this.socketAddress = mpa;
this.bindAddress = bindaddr;
this.monitor = monitor;
if (bindaddr == null) {
monitor.setName("" + mpa.toString());
} else {
monitor.setName(bindaddr.toString() + "" + mpa.toString());
}
}
public KLALBRemoteLine(MultipurposeSocketAddress mpa, MultipurposeSocketAddress bindaddr) {
this(mpa,bindaddr,new Monitor());
this(mpa, bindaddr, new SpeedAndTrafficAndDelayMonitorDataImpl());
}
PrintStream pw;
private void startRecord() {
try {
pw=new PrintStream(UUID.randomUUID()+".csv");
pw.println("时间,状态,上传速度,下载速度,上传延迟,下载延迟,往返延迟,抖动");
ThreadTool.makePDaemonThreadIfSupport("BigData Record", ()->{
while(true) {
PrintStream pw;
private boolean flag=false;
private void startRecord() {
if(flag)
return;
flag=true;
try {
File fl = new File("logs");
if (!fl.exists()) {
fl.mkdirs();
}
File fx = new File(fl, getMonitor().getName().replace(':', '') + ".csv");
pw = new PrintStream(new FileOutputStream(fx, true));
if (fx.length() <= 0)
pw.println("时间,状态,上传速度,下载速度,上传延迟,下载延迟,往返延迟,抖动");
ThreadTool.makeVDaemonThreadIfSupport("BigData Record", () -> {
while (!closed) {
try {
Thread.sleep(1000);
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
pw.println(System.currentTimeMillis()+","+monitor.getState()+","+monitor.getOutSpeedAvg()+","+monitor.getInSpeedAvg()+","+sndDelayFactor+","+rcvDelayFactor+","+monitor.getLatencyAvg()+","+monitor.getJitter());
pw.println(System.currentTimeMillis() + "," + monitor.getState() + "," + monitor.getOutSpeed() + ","
+ monitor.getInSpeed() + "," + monitor.getOutDelay() + "," + monitor.getInDelay() + ","
+ monitor.getLatency() + "," + monitor.getTotalJitter());
}
pw.close();
}).start();
} catch (FileNotFoundException e) {
e.printStackTrace();
}
}
private void endRecord() {
if(pw!=null)
pw.close();
}
protected void startIO() {
ThreadTool.makePDaemonThreadIfSupport("远程接收线程", () -> {
//startRecord();
while(true) {
try {
monitor.setState(Monitor.CONNECTING);
if(socketAddress==null) {
if(kplink.isClosed()) {
monitor.setState(Monitor.OFFLINE);
closed=true;
break;
}
}else {
try {
kplink=createLink(bindAddress,socketAddress);
}catch(BindException be){
//be.printStackTrace();
close();
throw be;
}
}
kplink.setSoTimeout(10000);
Thread t=ThreadTool.makePDaemonThreadIfSupport("远程发送线程", () -> {
private SpeedLimiter congress = new SpeedLimiter();
private long congressSpeed = 0;
private double increaceFactor = 1.2;
protected void startIO() {
ThreadTool.makeVDaemonThreadIfSupport("远程接收线程", () -> {
while (true) {
try {
writePacketToKPL(new VADDRREQPacket());
writePacketToKPL(new VADDRREQPacket());
writePacketToKPL(new VADDRREQPacket());
while ((!kplink.isClosed())&&(!closed)) {
if(checkPingTime()) {
writePacketToKPL(new PINGPacket(System.nanoTime()));
}else {
KLALBPacket kpp=getNextPacket();
long lth=statLengthBefore(10);
if(lth>65536*2) {
monitor.updateOutSpeedMax();
}
if(kpp!=null) {
pingInterval= 100000000L;
writePacketToKPL(kpp);
}else {
pingInterval=1000000000L;
tlock=Thread.currentThread();
LockSupport.parkNanos(1000000);
monitor.setState(MonitorData.CONNECTING);
if (socketAddress == null) {
if (kplink.isClosed()) {
close();
break;
}
} else {
kplink = createLink(bindAddress, socketAddress);
}
startRecord();
kplink.setSoTimeout(10000);
Thread t = ThreadTool.makeVDaemonThreadIfSupport("远程发送线程", () -> {
try {
writePacketToKPL(new VADDRREQPacket());
writePacketToKPL(new VADDRREQPacket());
writePacketToKPL(new VADDRPacket(localVaddrSupplier.get()));
writePacketToKPL(new VADDRPacket(localVaddrSupplier.get()));
flushKPL();
KLALBPacket kpp = null;
while ((!kplink.isClosed()) && (!closed)) {
// TimeDebugger tdb=new TimeDebugger();
// tdb.putTime("start");
boolean flsh = false;
if (checkPingTime()) {
writePacketToKPL(new PINGPacket(System.nanoTime()));
if (remoteVaddr == null)
writePacketToKPL(new VADDRREQPacket());
monitor.updateSpeedSync();
flsh = true;
}
KLALBPacket kpip = IsendDequeList.poll();
if (kpip != null) {
writePacketToKPL(kpip);
flsh = true;
}
// tdb.putTime("Isend");
if (kpp == null) {
do {
kpp = sendDequeList.poll();
} while (kpp != null && kpp.isDisposed());
}
// System.out.println(sendDequeList.size());
if (kpp != null) {
kpp.getDisposeLock().lock();
try {
if (kpp.isDisposed()) {
kpp.getDisposeLock().unlock();
kpp = null;
} else {
if (congress.checkTransmit(kpp.getLength())) {
writePacketToKPL(kpp);
if (kpp instanceof DATATPacket) {
resetSleepTimer();
}
flsh = true;
kpp.getDisposeLock().unlock();
kpp = null;
}
}
} finally {
if(kpp!=null)
kpp.getDisposeLock().unlock();
}
}
if (flsh) {
flushKPL();
} else {
// tdb.putTime("send");
// tdb.putTime("flush");
if (pressure && pressureSpeed.checkTransmit(testPacket.getLength())) {
writePacketToKPL(testPacket);
flushKPL();
} else {
tlock = Thread.currentThread();
LockSupport.parkNanos(1000000);
}
// tdb.putTime("park");
}
// tdb.print();
// System.out.println(sendDequeList.isEmpty());
Thread.yield();
}
} catch (IOException | UnresolvedAddressException e) {
if (debug)
e.printStackTrace();
} finally {
try {
kplink.close();
} catch (IOException e) {
e.printStackTrace();
}
}
});
t.start();
boolean fst = true;
long stime = 5000000000L;
while ((!kplink.isClosed()) && (!closed)) {
long readStart = System.nanoTime();
KLALBPacket kpp = readPacketFromKPL();
long readTime = System.nanoTime() - readStart;
readTime *= 10;
if (readTime > stime) {
stime = readTime;
} else {
stime = (stime * 99 + readTime) / 100;
}
if (stime < 2000000000L) {
stime = 2000000000L;
}
// System.out.println(readTime);
// kplink.setSoTimeout(5000);
kplink.setSoTimeout((int) (stime / 1000000));
if (kpp == null)
break;
switch (kpp.getType()) {
case KLALBPacket.PING:
sendIPacket(new PONGPacket(((PINGPacket) kpp).getTime() ,kpp.getRcvtime(), System.nanoTime()));
sendIPacket(new BWINFPacket(monitor.getOutSpeed(), monitor.getInSpeed()));
break;
case KLALBPacket.PONG:
PONGPacket png = (PONGPacket) kpp;
long ul = png.getRcvtime() - png.getTimepingsnd()-(png.getTimepongsnd()-png.getTimepingrcv());
long DsndDelayFactor0 = png.getTimepingrcv() - png.getTimepingsnd();
long DrcvDelayFactor0 = png.getRcvtime() - png.getTimepongsnd();
long DsndDelayFactor;
long DrcvDelayFactor;
adjnc.getLock().lock();
try {
adjnc.calibrate(png.getTimepongsnd() + ul / 2 - png.getRcvtime(), ul / 2);
DsndDelayFactor = DsndDelayFactor0 - adjnc.getDelta();
DrcvDelayFactor = DrcvDelayFactor0 + adjnc.getDelta();
if (DsndDelayFactor < 0) {
if (DrcvDelayFactor < 0) {
} else {
adjnc.setDelta(adjnc.getDelta() + DsndDelayFactor);
DsndDelayFactor = DsndDelayFactor0 - adjnc.getDelta();
DrcvDelayFactor = DrcvDelayFactor0 + adjnc.getDelta();
}
} else {
if (DrcvDelayFactor < 0) {
adjnc.setDelta(adjnc.getDelta() - DrcvDelayFactor);
DsndDelayFactor = DsndDelayFactor0 - adjnc.getDelta();
DrcvDelayFactor = DrcvDelayFactor0 + adjnc.getDelta();
}
}
} finally {
adjnc.getLock().unlock();
}
if (DsndDelayFactor >= 0 && DrcvDelayFactor >= 0) {
monitor.setOutDelay(DsndDelayFactor);
monitor.setInDelay(DrcvDelayFactor);
monitor.setRecentPingNanoTime(png.getTimepingsnd());
pingInterval = monitor.getOutDelayMin();
double load = 1 - monitor.getOutDelayMin() / (double) monitor.getOutDelay();
increaceFactor =Math.max( 2.0 - load*2,1.0);
}
LockSupport.unpark(tlock);
break;
case KLALBPacket.BWINF:
BWINFPacket bwi = (BWINFPacket) kpp;
if (bwi.getDownSpeed() >= congressSpeed) {
congressSpeed = bwi.getDownSpeed();
} else {
congressSpeed = (congressSpeed * 999 + bwi.getDownSpeed()) / 1000;
}
congress.setLimitspeed((long) (congressSpeed * increaceFactor) + 32768);
// System.out.println(congressSpeed);
break;
case KLALBPacket.VADDRREQ:
sendIPacket(new VADDRPacket(localVaddrSupplier.get()));
break;
case KLALBPacket.VADDR:
Inet6Address vdr = ((VADDRPacket) kpp).getVaddr();
remoteVaddr = vdr;
sendIPacket(new VADDRACKPacket());
if (klalbController != null) {
adjnc = klalbController.getAdjustedClockByVaddr(remoteVaddr);
}
if (fst) {
// coll.resetCoolingTime();
monitor.setState(MonitorData.ONLINE);
fst = false;
}
break;
case KLALBPacket.VADDRACK:
break;
case KLALBPacket.TEST:
break;
default:
while (rec == null) {
Thread.sleep(1);
}
rec.accept(KLALBRemoteLine.this, kpp);
break;
}
Thread.yield();
}
} catch (IOException e) {
} catch (IOException | UnresolvedAddressException e) {
if (debug)
e.printStackTrace();
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
monitor.setState(MonitorData.OFFLINE);
pressure = false;
try {
kplink.close();
if (kplink != null)
kplink.close();
} catch (IOException e) {
e.printStackTrace();
}
}
});
t.start();
boolean fst=true;
long stime=5000000000L;
while ((!kplink.isClosed())&&(!closed)) {
long readStart=System.nanoTime();
KLALBPacket kpp = readPacketFromKPL();
long readTime=System.nanoTime()-readStart;
if(readTime>stime) {
stime=readTime;
}else {
stime=(stime*99+readTime)/100;
}
if(stime<3000000000L) {
stime=3000000000L;
}
//System.out.println(readTime);
kplink.setSoTimeout((int) (stime/1000000));
if(kpp==null)
break;
switch (kpp.getType()) {
case KLALBPacket.PING:
sendPacket(new PONGPacket(((PINGPacket) kpp).getTime(),kpp.getRcvtime(),System.nanoTime()), -1);
break;
case KLALBPacket.PONG:
PONGPacket png=(PONGPacket) kpp;
long ul=png.getRcvtime()-( png.getTimepingsnd()+(png.getTimepongsnd()-png.getTimepingrcv()));
//System.out.println(toString()+"\t"+ul);
monitor.updateLatency (ul);
long DsndDelayFactor=png.getTimepingrcv()- png.getTimepingsnd();
long DrcvDelayFactor=png.getRcvtime()- png.getTimepongsnd();
//sndDelayFactor=DsndDelayFactor;
//rcvDelayFactor=DrcvDelayFactor;
if(sndDelayFactor==Long.MIN_VALUE) {
sndDelayFactor=DsndDelayFactor;
}else {
sndDelayFactor=(sndDelayFactor+ DsndDelayFactor)/2;
KLALBPacket pack;
while((pack=sendDequeList.poll())!=null) {
try {
if(remoteVaddr!=null) {
klalbController.sendPacketToAddress(remoteVaddr, pack);
if(debug)
System.out.println("RETRY:"+pack);
}
if(rcvDelayFactor==Long.MIN_VALUE) {
rcvDelayFactor=DrcvDelayFactor;
}else {
rcvDelayFactor=(rcvDelayFactor+ DrcvDelayFactor)/2;
}
LockSupport.unpark(tlock);
break;
case KLALBPacket.VADDRREQ:
sendPacket(new VADDRPacket(localVaddrSupplier.get()), -1);
break;
case KLALBPacket.VADDR:
remoteVaddr=((VADDRPacket) kpp).getVaddr();
sendPacket(new VADDRACKPacket(), -1);
break;
case KLALBPacket.VADDRACK:
if(fst) {
monitor.resetCoolingTime();
monitor.setState(Monitor.ONLINE);
fst=false;
}
break;
case KLALBPacket.TEST:
break;
default:
while(rec==null) {
Thread.sleep(1);
}
rec.accept(KLALBRemoteLine.this,kpp);
break;
} catch (IOException e) {
}
}
Thread.yield();
}
} catch (IOException e) {
e.printStackTrace();
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
monitor.setState(Monitor.OFFLINE);
try {
if(kplink!=null)
kplink.close();
} catch (IOException e) {
e.printStackTrace();
if (closed) {
break;
}
// Thread.sleep(monitor.getCoolingTime());
lthd = Thread.currentThread();
LockSupport.parkNanos(coll.getCoolingTime(monitor.getReliability()));
// coll.incCoolingTime();
}
if(closed) {
break;
}
//Thread.sleep(monitor.getCoolingTime());
lthd=Thread.currentThread();
LockSupport.parkNanos(monitor.getCoolingTime()*1000000L);
monitor.incCoolingTime();
}
endRecord();
}).start();
}
private void flushKPL() throws IOException {
kplink.flush();
}
private KLALBPacket readPacketFromKPL() throws IOException {
KLALBPacket packet=kplink.readPacket();
/*if(!(packet instanceof PINGPacket) )
if(!(packet instanceof PONGPacket) )
System.out.println("RX:" + packet);*/
if(packet!=null)
monitor.getInTrafficAL().addAndGet(packet.getLength());
KLALBPacket packet = kplink.readPacket();
if (showpacket) {
if (!(packet instanceof PINGPacket))
if (!(packet instanceof PONGPacket))
System.out.println("RX:" + packet);
}
if (packet != null) {
monitor.getInTrafficAL().addAndGet(packet.getLength());
if (klalbController != null)
klalbController.getLinkMonitor().getInTrafficAL().addAndGet(packet.getLength());
}
return packet;
}
private void writePacketToKPL(KLALBPacket packet) throws IOException {
monitor.getOutTrafficAL().addAndGet(packet.getLength());
if (klalbController != null)
klalbController.getLinkMonitor().getOutTrafficAL().addAndGet(packet.getLength());
kplink.writePacket(packet);
/*if(!(packet instanceof PINGPacket) )
if(!(packet instanceof PONGPacket) )
System.err.println("TX:" + packet);*/
if (showpacket) {
if (!(packet instanceof PINGPacket))
if (!(packet instanceof PONGPacket))
System.err.println("TX:" + packet);
}
}
public void close() {
closed=true;
closed = true;
monitor.setState(MonitorData.OFFLINE);
pressure = false;
try {
if(kplink!=null)
kplink.close();
if (kplink != null)
kplink.close();
} catch (IOException e) {
e.printStackTrace();
}
}
public boolean isClosed() {
return closed;
}
private volatile long timeTurnSleep = System.nanoTime();
private volatile long time = System.nanoTime();
private volatile long pingInterval=1000000000L;
private volatile long pingInterval = 50000000L;
private volatile long pingIntervalSleep = 200000000L;
private void resetSleepTimer() {
timeTurnSleep = System.nanoTime();
}
private boolean checkPingTime() {
long cu = System.nanoTime();
if (cu - time > pingInterval) {
if (cu - time > ((System.nanoTime() - timeTurnSleep > 500000000) ? pingIntervalSleep
: Math.min(pingIntervalSleep, pingInterval))) {
time = cu;
return true;
} else {
return false;
}
}
private volatile BiConsumer<KLALBRemoteLine, KLALBPacket> rec;
//private Object lock=new Object();
private volatile Thread tlock;
private List<SumQueue<KLALBPacket>> sendDequeList =new ArrayList<SumQueue<KLALBPacket>>();
{
for(int i=0;i<12;i++) {
sendDequeList.add(new SumQueue<KLALBPacket>(new ConcurrentLinkedQueue<>()));
}
private ConcurrentLinkedQueue<KLALBPacket> IsendDequeList = new ConcurrentLinkedQueue<KLALBPacket>();
private PriorityBlockingQueue<KLALBPacket> sendDequeList = new PriorityBlockingQueue<KLALBPacket>();
public PriorityBlockingQueue<KLALBPacket> getQueue() {
return sendDequeList;
}
private KLALBPacket getNextPacket() {
for (Iterator<SumQueue<KLALBPacket>> iterator = sendDequeList.iterator(); iterator.hasNext();) {
SumQueue<KLALBPacket> queue = (SumQueue<KLALBPacket>) iterator.next();
KLALBPacket v=queue.poll();
if(v!=null)
return v;
if(monitor.getLatencyCurr()>(monitor.getLatencyMin()<<3)) {
//System.out.println(monitor.getLatencyCurr()+" "+monitor.getLatencyMin());
return null;
}
}
return null;
}
private long sumNextPacketSize() {
long i=0;
for (Iterator<SumQueue<KLALBPacket>> iterator = sendDequeList.iterator(); iterator.hasNext();) {
SumQueue<KLALBPacket> queue = (SumQueue<KLALBPacket>) iterator.next();
i+=queue.getSumValue();
}
return i;
}
private long sumNextPacketCount() {
long i=0;
for (Iterator<SumQueue<KLALBPacket>> iterator = sendDequeList.iterator(); iterator.hasNext();) {
SumQueue<KLALBPacket> queue = (SumQueue<KLALBPacket>) iterator.next();
i+=queue.size();
}
return i;
}
public void sendPacket(KLALBPacket kp) {
sendPacket(kp, 5);
}
public void sendPacket(KLALBPacket blk,int prio) {
sendDequeList.get(prio+1).add(blk);
protected void sendPacket(KLALBPacket blk) {
sendDequeList.add(blk);
LockSupport.unpark(tlock);
}
public void setPacketReceiver(BiConsumer<KLALBRemoteLine,KLALBPacket> rec) {
private void sendIPacket(KLALBPacket blk) {
IsendDequeList.add(blk);
LockSupport.unpark(tlock);
}
public void setPacketReceiver(BiConsumer<KLALBRemoteLine, KLALBPacket> rec) {
this.rec = rec;
}
public void remoeFromSendQueue(KLALBPacket klalbPacket) {
for (Iterator<SumQueue<KLALBPacket>> iterator = sendDequeList.iterator(); iterator.hasNext();) {
SumQueue<KLALBPacket> queue = iterator.next();
queue.removeIf((x)->{
return x.equals(klalbPacket);
});
}
}
public long statLengthBefore(int priority) {
AtomicLong al=new AtomicLong();
int n=0;
for (Iterator<SumQueue<KLALBPacket>> iterator = sendDequeList.iterator(); iterator.hasNext();n++) {
if(n>priority) {
break;
/*public void remoeFromSendQueue(KLALBPacket klalbPacket) {
sendDequeList.remove(klalbPacket);
}*/
public long statLengthBefore(long l) {
AtomicLong al = new AtomicLong();
for (Iterator<KLALBPacket> iterator = sendDequeList.iterator(); iterator.hasNext();) {
KLALBPacket queue = iterator.next();
if (queue.getPriority() <= l) {
al.addAndGet(queue.getLength());
}
SumQueue<KLALBPacket> queue = iterator.next();
al.addAndGet(queue.getSumValue());
}
return al.get();
}
private long sndDelayFactor=Long.MIN_VALUE;
private long rcvDelayFactor=Long.MIN_VALUE;
public long getSndDelayFactor() {
return sndDelayFactor;
}
public long getRcvDelayFactor() {
return rcvDelayFactor;
}
private volatile long predictTime;
public void runPredict(KLALBPacket curr,int priority) {
predictTime=sndDelayFactor;
//System.out.println(this+" "+sndDelayFactor);
long al=0;
al=statLengthBefore(priority);
long speed=getMonitor(). getOutSpeed();
if(speed==0) {
if(al>0) {
predictTime=Long.MAX_VALUE;
}
}else {
predictTime+=al*1000000000L/speed;
}
public void runPredict() {
//predictTime = monitor.getOutDelayPredicted();
predictTime=monitor.getOutDelay();
/*
* long al=0; al=statLengthBefore(5); long speed=getMonitor().getOutSpeedAvg();
* // spdlmt.setLimitspeed(speed*2+65536); if(speed==0) { if(al>0) {
* predictTime=Long.MAX_VALUE; } }else { predictTime+=al*1000000000L/speed; }
*/
}
@Override
public int compareTo(KLALBRemoteLine o2) {
long t1=this.predictTime;
long t2=o2.predictTime;
if(t1>t2) {
long t1 = this.predictTime;
long t2 = o2.predictTime;
if (t1 > t2) {
return 1;
}else if(t1<t2){
} else if (t1 < t2) {
return -1;
}else {
} else {
return 0;
}
}
public void reconnectImmediately() {
if(lthd!=null) {
if (lthd != null) {
LockSupport.unpark(lthd);
}
}
public void dislink() {
if (kplink != null) {
try {
kplink.close();
} catch (IOException e) {
}
}
}
private static final TESTPacket testPacket = new TESTPacket();
private volatile boolean pressure = false;
private SpeedLimiter pressureSpeed = new SpeedLimiter(32768);
public void pressureTest() {
if (pressure)
throw new IllegalStateException("Test already start");
this.pressure = true;
pressureSpeed.setLimitspeed(32768);
ThreadTool.makeVThreadIfSupport("压力测试控制线程", () -> {
while (pressure) {
try {
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
if (getMonitor().getOutSpeed() < pressureSpeed.getLimitspeed() * 2 / 3
|| getMonitor().getInDelayMin() * 30 < getMonitor().getInDelay()) {
pressure = false;
break;
}
pressureSpeed.setLimitspeed(pressureSpeed.getLimitspeed() * 400 / 399);
}
}).start();
}
public long getPredictedTime() {
return predictTime;
}
}