forked from KNEMC/KLALB
KLALB 3.6.0开发一半的状态
This commit is contained in:
Generated
+10
@@ -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
|
||||
Generated
+6
@@ -0,0 +1,6 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project version="4">
|
||||
<component name="ProjectRootManager" version="2">
|
||||
<output url="file://$PROJECT_DIR$/classes" />
|
||||
</component>
|
||||
</project>
|
||||
Generated
+8
@@ -0,0 +1,8 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project version="4">
|
||||
<component name="ProjectModuleManager">
|
||||
<modules>
|
||||
<module fileurl="file://$PROJECT_DIR$/KLALB.iml" filepath="$PROJECT_DIR$/KLALB.iml" />
|
||||
</modules>
|
||||
</component>
|
||||
</project>
|
||||
Generated
+6
@@ -0,0 +1,6 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project version="4">
|
||||
<component name="VcsDirectoryMappings">
|
||||
<mapping directory="" vcs="Git" />
|
||||
</component>
|
||||
</project>
|
||||
@@ -0,0 +1,253 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<module type="JAVA_MODULE" version="4">
|
||||
<component name="EclipseModuleManager">
|
||||
<libelement value="jar://$MODULE_DIR$/lib/gson-2.1.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/javassist.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/ini4j.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/JFreeChart1.5.2.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/toml4j-0.7.1.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/pcap4j-core-1.8.2.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/pcap4j-packetfactory-propertiesbased-1.8.2.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/pcap4j-packetfactory-static-1.8.2.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/pcap4j-sample-1.8.2.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/slf4j-api-2.0.9.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/slf4j-api-2.0.9-javadoc.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/slf4j-api-2.0.9-sources.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/slf4j-api-2.0.9-tests.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/disruptor-4.0.0.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/disruptor-4.0.0-javadoc.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/disruptor-4.0.0-sources.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/KNElib1.0.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/KNEOptimize.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/JavaTUN0.3.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/uuid-creator-6.1.1.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/jctools-core-4.0.6.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/jctools-core-4.0.6-javadoc.jar!/" />
|
||||
<libelement value="jar://$MODULE_DIR$/lib/jctools-core-4.0.6-sources.jar!/" />
|
||||
<src_description expected_position="0">
|
||||
<src_folder value="file://$MODULE_DIR$/src" expected_position="0" />
|
||||
</src_description>
|
||||
</component>
|
||||
<component name="NewModuleRootManager">
|
||||
<output url="file://$MODULE_DIR$/bin" />
|
||||
<exclude-output />
|
||||
<content url="file://$MODULE_DIR$">
|
||||
<sourceFolder url="file://$MODULE_DIR$/src" isTestSource="false" />
|
||||
</content>
|
||||
<orderEntry type="sourceFolder" forTests="false" />
|
||||
<orderEntry type="module-library">
|
||||
<library name="gson-2.1.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/gson-2.1.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="javassist.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/javassist.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES>
|
||||
<root url="file://$MODULE_DIR$/../../../../Users/ADMINI~1/AppData/Local/Temp/1/.org.sf.feeling.decompiler1686963856070/source/javassist-3.29.2-GA-sources.jar" />
|
||||
</SOURCES>
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="ini4j.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/ini4j.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="JFreeChart1.5.2.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/JFreeChart1.5.2.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="toml4j-0.7.1.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/toml4j-0.7.1.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="pcap4j-core-1.8.2.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/pcap4j-core-1.8.2.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES>
|
||||
<root url="file://$MODULE_DIR$/../../../../Users/ADMINI~1/AppData/Local/Temp/1/.org.sf.feeling.decompiler1725520274095/source/pcap4j-core-1.8.2-sources.jar" />
|
||||
</SOURCES>
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="pcap4j-packetfactory-propertiesbased-1.8.2.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/pcap4j-packetfactory-propertiesbased-1.8.2.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="pcap4j-packetfactory-static-1.8.2.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/pcap4j-packetfactory-static-1.8.2.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="pcap4j-sample-1.8.2.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/pcap4j-sample-1.8.2.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="slf4j-api-2.0.9.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/slf4j-api-2.0.9.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="slf4j-api-2.0.9-javadoc.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/slf4j-api-2.0.9-javadoc.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="slf4j-api-2.0.9-sources.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/slf4j-api-2.0.9-sources.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="slf4j-api-2.0.9-tests.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/slf4j-api-2.0.9-tests.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="disruptor-4.0.0.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/disruptor-4.0.0.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES>
|
||||
<root url="jar://$MODULE_DIR$/lib/disruptor-4.0.0-sources.jar!/" />
|
||||
</SOURCES>
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="disruptor-4.0.0-javadoc.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/disruptor-4.0.0-javadoc.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="disruptor-4.0.0-sources.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/disruptor-4.0.0-sources.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="KNElib1.0.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/KNElib1.0.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="KNEOptimize.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/KNEOptimize.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="JavaTUN0.3.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/JavaTUN0.3.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="uuid-creator-6.1.1.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/uuid-creator-6.1.1.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="jctools-core-4.0.6.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/jctools-core-4.0.6.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="jctools-core-4.0.6-javadoc.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/jctools-core-4.0.6-javadoc.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="module-library">
|
||||
<library name="jctools-core-4.0.6-sources.jar">
|
||||
<CLASSES>
|
||||
<root url="jar://$MODULE_DIR$/lib/jctools-core-4.0.6-sources.jar!/" />
|
||||
</CLASSES>
|
||||
<JAVADOC />
|
||||
<SOURCES />
|
||||
</library>
|
||||
</orderEntry>
|
||||
<orderEntry type="jdk" jdkName="jdk-26.0.1" jdkType="JavaSDK" />
|
||||
</component>
|
||||
</module>
|
||||
Binary file not shown.
Binary file not shown.
Binary file not shown.
@@ -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();
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
public interface CongestionAlgorithmFactory {
|
||||
public CongestionAlgorithm create();
|
||||
|
||||
public String getName();
|
||||
}
|
||||
@@ -0,0 +1,18 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
public class CongestionAlgorithms {
|
||||
private static Map<String,CongestionAlgorithmFactory>regs=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();
|
||||
}
|
||||
}
|
||||
@@ -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<T> implements MessageBatcher<T> {
|
||||
private static final Cleaner clr = Cleaner.create();
|
||||
|
||||
private final MpscMessageBatcher0<T> impl;
|
||||
private final MpscMessageBatcherReleaser<T> 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<List<T>> consumer) { impl.setConsumer(consumer); }
|
||||
|
||||
@Override
|
||||
public void putMessage(T message) { impl.putMessage(message); }
|
||||
|
||||
@Override
|
||||
public void putMessages(List<T> 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<Integer> 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<T> extends Releaser<MpscMessageBatcher0<T>> {
|
||||
public MpscMessageBatcherReleaser(MpscMessageBatcher0<T> resource) {
|
||||
super(resource);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void release(MpscMessageBatcher0<T> resource) {
|
||||
resource.close();
|
||||
}
|
||||
}
|
||||
|
||||
class MpscMessageBatcher0<T> implements Runnable, MessageBatcher<T> {
|
||||
|
||||
private final MpscArrayQueue<T> queue;
|
||||
private final int batchSize;
|
||||
private final long maxDelay;
|
||||
private volatile Consumer<List<T>> 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<List<T>> 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<T> 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) {
|
||||
ArrayList<T>batch=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;
|
||||
}
|
||||
}
|
||||
@@ -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<K, V extends NetworkPacket> implements Closeable, AutoCloseable{
|
||||
private ThreadParker tp=new ThreadParker();
|
||||
|
||||
private Map<K,SendItem<V>> recvMap = new ConcurrentHashMap<K,SendItem<V>>();
|
||||
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<V> newitem = new SendItem<V>(packet,headerCalibrate);
|
||||
SendItem<V> 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<V> 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<K, SendItem<V>> 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 + "]";
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
}
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
package org.kne.cloud.network.ipv6;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
public interface IPv6ProtocolRegister {
|
||||
public boolean onaccept(IPv6Packet packx)throws IOException;
|
||||
}
|
||||
@@ -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<IPv6Packet> 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<IPv6Packet> 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);
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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 ReceivePacketSlidingWindow<Long, KLTPPacket>recvMap=new ReceivePacketSlidingWindow<Long, KLTPPacket>(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<offset+length;i++) {
|
||||
lth+=read(dsts[i]);
|
||||
if(dsts[i].hasRemaining()) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
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_DATA:
|
||||
//ackSequenceBatcher.putMessage(kseq);
|
||||
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);
|
||||
});
|
||||
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 + "]";
|
||||
}
|
||||
|
||||
|
||||
|
||||
}
|
||||
@@ -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 SendPacketSlidingWindow<Long,KLTPPacket>sendMap=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<IOException> 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;i<offset+length;i++) {
|
||||
lth+=writeWithoutFlush(srcs[i]);
|
||||
}
|
||||
if(dataPack!=null&& dataPack.getKLTPData().position()>0) {
|
||||
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 + "]";
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
}
|
||||
@@ -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();
|
||||
}
|
||||
@@ -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<KLTPSessionPacket> 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<UUID, KLTPPacketConsumer> recvRegisterMap=new ConcurrentHashMap<UUID, KLTPPacketConsumer>();
|
||||
private ConcurrentHashMap<UUID, KLTPPacketConsumer> sendRegisterMap=new ConcurrentHashMap<UUID, KLTPPacketConsumer>();
|
||||
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<KLTPSessionPacket> 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;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
@@ -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<Long> owdSupplier(DelayMonitorData monitor) {
|
||||
return () -> monitor.getOutDelay();
|
||||
}
|
||||
|
||||
/**
|
||||
* 静态成本 Supplier(用于测试或静态路由)
|
||||
*
|
||||
* @param cost 固定的成本值
|
||||
* @return 返回固定值的 Supplier
|
||||
*/
|
||||
public static Supplier<Long> staticSupplier(long cost) {
|
||||
return () -> cost;
|
||||
}
|
||||
/**
|
||||
* 使用你设计的“概率期望延迟”公式:OWD + RTO × (1 - Reliability)
|
||||
* @param monitor 延迟监控数据(提供 OWD)
|
||||
* @param linkStatus 链路状态(提供 Reliability)
|
||||
* @param rtoNanos 超时重传时间(纳秒)
|
||||
* @return 返回期望延迟的 Supplier
|
||||
*/
|
||||
public static Supplier<Long> expectedDelaySupplier(DelayMonitorData monitor, LinkStatus linkStatus, Supplier<Long> 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<Long> maxSupplier(Supplier<Long> a, Supplier<Long> b) {
|
||||
return () -> Math.max(a.get(), b.get());
|
||||
}
|
||||
|
||||
/**
|
||||
* 组合两个 Supplier,取最小值(可用于 ECMP 场景下的乐观调度)
|
||||
*/
|
||||
public static Supplier<Long> minSupplier(Supplier<Long> a, Supplier<Long> b) {
|
||||
return () -> Math.min(a.get(), b.get());
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,93 @@
|
||||
package org.kne.cloud.network.monitor;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
|
||||
/**
|
||||
* 链路状态监控类,负责维护链路的在线状态和在线率。
|
||||
* 设计理念:
|
||||
* 1. 状态变化时通过回调通知监听者。
|
||||
* 2. 在线率采用指数加权移动平均 (EWMA) 算法,平滑且对历史数据有衰减记忆。
|
||||
* 3. 自身不启动任何后台线程,状态的更新由外部(例如收到心跳包时)主动触发。
|
||||
* 4. 不包含链路名称,名称由外部管理(如 Map<String, LinkStatus>),实现关注点分离。
|
||||
*/
|
||||
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<LinkStatus> 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<LinkStatus> 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);
|
||||
}
|
||||
}
|
||||
@@ -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<IPv6Packet> 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<IPv6Packet> cons;
|
||||
if((cons=distributePacketToConsumer(srcA, pt))!=null) {
|
||||
cons.accept(packx);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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<Runnable> queue=new MpscArrayQueue<Runnable>(2048);
|
||||
private AtomicInteger size=new AtomicInteger();
|
||||
|
||||
private volatile ThreadParker parker=new ThreadParker();
|
||||
public Queue<Runnable> 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<Boolean> 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<Boolean> command,int limit) {
|
||||
for(int i=0;i<threads.length;i++) {
|
||||
ThreadElement te= threads[i];
|
||||
boolean b=te.size()>limit;
|
||||
if(!b) {
|
||||
if(te.putTask( ()->{command.accept(false);})){
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for(int i=0;i<threads.length;i++) {
|
||||
ThreadElement te= threads[i];
|
||||
if(te.putTask( ()->{command.accept(true);})){
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
public void executeWithCongestionReport(Consumer<Boolean> command) {
|
||||
if(execute0(command,1000)) {
|
||||
return;
|
||||
}
|
||||
System.out.println("loss!");
|
||||
}
|
||||
public void executeWithCongestionReport(Consumer<Boolean> command,int limit) {
|
||||
if(execute0(command,limit)) {
|
||||
return;
|
||||
}
|
||||
System.out.println("loss!");
|
||||
|
||||
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user