=headers.size()) {
- irh.setNextHeader(payload.getType());
- }else {
- irh.setNextHeader(headers.get(index).getType());
- }
- if(index-1<0) {
- setNextHeader(irh.getType());
- }else {
- headers.get(index-1).setNextHeader(irh.getType());
- }
- headers.add(index, irh);
- irh.setRoutingType(4);
- irh.setSegmentsLeft(srh.getLastEntry());
- irh.getData().put(4, srh.getRawData());
- analyse();
- return irh;
- }*/
- @Override
- protected boolean needEndPosition() {
- return false;
- }
- private AtomicInteger rerouteCounter=new AtomicInteger(0);
- public AtomicInteger getRerouteCounter() {
- return rerouteCounter;
- }
- private volatile boolean promise=false;
- public void setPromise(boolean b) {
- promise=b;
- }
+ public IPv6RoutingHeader() {
+ super(ROUTING);
+ }
- public boolean isPromise() {
- return promise;
- }
+ public IPv6RoutingHeader(int bufferLength) {
+ super(ROUTING, bufferLength);
+ }
+
+ public IPv6RoutingHeader(int bufferLength, ByteBuffer data) {
+ super(ROUTING, bufferLength, data);
+ }
+
+ public int getRoutingType() {
+ return getData().get(2) & BIT_MASK_LOW_8_BITS;
+ }
+
+ public void setRoutingType(int routingTypr) {
+ getData().put(2, (byte) routingTypr);
+ }
+
+ public int getSegmentsLeft() {
+ return getData().get(3);
+ }
+
+ public void setSegmentsLeft(int segmentsLeft) {
+ getData().put(3, (byte) segmentsLeft);
+ }
+ }
+
+ /**
+ * IPv6段路由头部类
+ */
+ /**
+ * IPv6 Segment Routing Header (SRH) 实现
+ *
+ * 根据 RFC 8754 实现,支持SRv6扩展头和TLV选项
+ *
+ *
+ * 0 1 2 3
+ * 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
+ * +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
+ * | Next Header | Hdr Ext Len | Routing Type | Segments Left |
+ * +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
+ * | Last Entry | Flags | Tag |
+ * +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
+ * | |
+ * | Segment List[0] (128-bit IPv6 address) |
+ * | |
+ * +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
+ * | ... |
+ * +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
+ * | |
+ * | Segment List[n] (128-bit IPv6 address) |
+ * | |
+ * +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
+ * // //
+ * // Optional Type-Length-Value objects (variable) //
+ * // //
+ * +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
+ *
+ */
+ public static class IPv6SegmentRoutingHeader extends IPv6RoutingHeader {
+
+ private final List addresses = new ArrayList();
+ private final List tlvs = new ArrayList();
+
+ @Override
+ public String toString() {
+ return "IPv6SegmentRoutingHeader [addresses=" + addresses + ", tlvs=" + tlvs + ", getLastEntry()="
+ + getLastEntry() + ", getFlags()=" + getFlags() + ", getTag()=" + getTag() + ", getAddresses()="
+ + getAddresses() + ", getRoutingType()=" + getRoutingType() + ", getSegmentsLeft()="
+ + getSegmentsLeft() + ", getProtocolNumber()=" + getProtocolNumber() + ", getNextHeader()="
+ + getNextHeader() + ", getExtLength()=" + getExtLength() + "]";
+ }
+
+ public IPv6SegmentRoutingHeader() {
+ super(EXT_HEADER_ALIGNMENT);
+ setRoutingType(SRV6_ROUTING_TYPE);
+ }
+
+ public IPv6SegmentRoutingHeader(ByteBuffer bbf) {
+ super(EXT_HEADER_ALIGNMENT, bbf);
+ }
+
+ public IPv6SegmentRoutingHeader(List segs) {
+ this();
+ this.addresses.addAll(segs);
+ int ln = addresses.size() - 1;
+ setLastEntry(ln);
+ setSegmentsLeft(ln);
+ }
+
+ public int getLastEntry() {
+ return getData().get(4) & BIT_MASK_LOW_8_BITS;
+ }
+
+ public void setLastEntry(int lastEntry) {
+ getData().put(4, (byte) lastEntry);
+ }
+
+ public int getFlags() {
+ return getData().get(5) & BIT_MASK_LOW_8_BITS;
+ }
+
+ public void setFlags(int flags) {
+ getData().put(5, (byte) flags);
+ }
+
+ public int getTag() {
+ return getData().getShort(6) & BIT_MASK_LOW_16_BITS;
+ }
+
+ public void setTag(int tag) {
+ getData().putShort(6, (short) tag);
+ }
+
+ public List getAddresses() {
+ return addresses;
+ }
+
+ public List getTlvs() {
+ return tlvs;
+ }
+
+ @Override
+ public void writeToChannel(WritableByteChannel dto) throws IOException {
+ int exl = calcExtLength();
+ if (exl % EXT_HEADER_ALIGNMENT != 0) {
+ throw new StreamCorruptedException("extLength % " + EXT_HEADER_ALIGNMENT + " !=0");
+ }
+
+ setExtLength(exl / EXT_HEADER_ALIGNMENT);
+ setLastEntry(addresses.size() - 1);
+ super.writeToChannel(dto);
+
+ for (Inet6Address inet6Address : addresses) {
+ dto.write(ByteBuffer.wrap(inet6Address.getAddress()));
+ }
+ for (IPv6SegmentRoutingTLV tlv : tlvs) {
+ IPv6SegmentRoutingTLV.writeIPv6SegmentRoutingTLVToChannel(dto, tlv);
+ }
+ }
+
+ private int calcExtLength() {
+ int tlvsl = 0;
+ for (IPv6SegmentRoutingTLV tlve : tlvs) {
+ tlvsl += tlve.getTotalLength();
+ }
+ return addresses.size() * IPV6_ADDRESS_LENGTH + tlvsl;
+ }
+
+ /*@Override
+ public void readFromChannel(ReadableByteChannel din, long length) throws IOException {
+ super.readFromChannel(din, length);
+ int extl = getExtLength();
+ int laste = getLastEntry();
+
+ int rl = extl * EXT_HEADER_ALIGNMENT;
+ int usdl = 0;
+
+ // System.out.println("SR length:"+(laste+1));
+ addresses.clear();
+ ByteBuffer bfr = NetworkPacket.bufferAllocator.allocateHeap(IPV6_ADDRESS_LENGTH);
+ for (int i = 0; i < (laste + 1); i++) {
+ bfr.clear();
+ while (bfr.hasRemaining()) {
+ if (din.read(bfr) == -1) {
+ throw new EOFException();
+ }
+ }
+ bfr.flip();
+ usdl += bfr.limit();
+ Inet6Address addr = (Inet6Address) Inet6Address.getByAddress(bfr.array());
+ // System.out.println("SR:"+addr);
+ addresses.add(addr);
+ }
+ tlvs.clear();
+ while (usdl < rl) {
+ IPv6SegmentRoutingTLV tlv = IPv6SegmentRoutingTLV.readIPv6SegmentRoutingTLVFromChannel(din);
+ usdl += tlv.getTotalLength();
+ tlvs.add(tlv);
+ }
+
+ return;
+ }*/
+ @Override
+ public void readFromChannel(ReadableByteChannel din, long length) throws IOException {
+ super.readFromChannel(din, length);
+
+ ByteBuffer bfr = null;
+ int extl = getExtLength();
+ int laste = getLastEntry();
+ int rl = extl * EXT_HEADER_ALIGNMENT;
+ int usdl = 0;
+
+ // 验证lastEntry的合理性
+ if (laste < 0 ) {
+ throw new StreamCorruptedException("Invalid lastEntry: " + laste);
+ }
+
+ addresses.clear();
+ bfr = NetworkPacket.bufferAllocator.allocateHeap(IPV6_ADDRESS_LENGTH);
+
+ for (int i = 0; i <= laste; i++) { // 注意:应该是 <= laste
+ bfr.clear();
+ KNEChannels.readFully(din, bfr);
+ bfr.flip();
+ usdl += bfr.limit();
+ Inet6Address addr = (Inet6Address) Inet6Address.getByAddress(bfr.array());
+ addresses.add(addr);
+ }
+
+ tlvs.clear();
+ while (usdl < rl) {
+ IPv6SegmentRoutingTLV tlv = IPv6SegmentRoutingTLV.readIPv6SegmentRoutingTLVFromChannel(din);
+ usdl += tlv.getTotalLength();
+ tlvs.add(tlv);
+ }
+
+ // 验证读取的字节数与预期一致
+ if (usdl != rl) {
+ throw new StreamCorruptedException("Length mismatch: expected " + rl + ", got " + usdl);
+ }
+
+
+ }
+
+ // 辅助方法:确保读取完整数据
+ private int readFully(ReadableByteChannel channel, ByteBuffer buffer) throws IOException {
+ int totalRead = 0;
+ while (buffer.hasRemaining()) {
+ int read = channel.read(buffer);
+ if (read == -1) {
+ break;
+ }
+ totalRead += read;
+ }
+ return totalRead;
+ }
+ @Override
+ public long getTotalLength() {
+ return calcExtLength() + EXT_HEADER_ALIGNMENT;
+ }
+
+ }
+
+ /**
+ * IPv6目标选项头部类
+ */
+ public static class IPv6DestinationHeader extends IPv6ExtHeader {
+ @Override
+ public String toString() {
+ return "IPv6DestinationHeader [protocolNumber=" + getProtocolNumber() + ", getNextHeader()="
+ + getNextHeader() + ", getExtLength()=" + getExtLength() + ", getRoutingType()=" + getRoutingType()
+ + "]";
+ }
+
+ public IPv6DestinationHeader() {
+ super(DESTINATION_OPTIONS);
+ }
+
+ public int getRoutingType() {
+ return getData().get(2) & BIT_MASK_LOW_8_BITS;
+ }
+
+ public void setRoutingType(int routingTypr) {
+ getData().put(2, (byte) routingTypr);
+ }
+
+ public int getSegmentsLeft() {
+ return getData().get(3);
+ }
+
+ public void setSegmentsLeft(int segmentsLeft) {
+ getData().put(3, (byte) segmentsLeft);
+ }
+ }
+
+ /**
+ * 获取所有扩展头部
+ * @return 扩展头部列表
+ */
+ public List getHeaders() {
+ return headers;
+ }
+
+ /**
+ * 获取IPv6逐跳头部
+ * @return IPv6段逐跳头部,如果不存在则返回null
+ */
+ public IPv6HopByHopHeader getHopByHopHeader() {
+ IPv6HopByHopHeader hoph = null;
+ List exhs = headers;
+ for (int j = 0; j < exhs.size(); j++) {
+ IPv6ExtHeader exh = exhs.get(j);
+ if (exh instanceof IPv6HopByHopHeader) {
+ hoph=(IPv6HopByHopHeader) exh;
+ break;
+ }
+ }
+ return hoph;
+ }
+
+ /**
+ * 获取SRv6段路由头部
+ * @return SRv6段路由头部,如果不存在则返回null
+ */
+ public IPv6SegmentRoutingHeader getSRHHeader() {
+ IPv6SegmentRoutingHeader srhh = null;
+ List exhs = headers;
+ for (int j = 0; j < exhs.size(); j++) {
+ IPv6ExtHeader exh = exhs.get(j);
+ if (exh instanceof IPv6SegmentRoutingHeader) {
+ if (((IPv6SegmentRoutingHeader) exh).getRoutingType() == SRV6_ROUTING_TYPE) {
+ srhh = (IPv6SegmentRoutingHeader) exh;
+ break;
+ }
+ }
+ }
+ return srhh;
+ }
+
+
+ @Override
+ protected boolean needEndPosition() {
+ return false;
+ }
+
+ // 重路由计数器
+ private AtomicInteger rerouteCounter = new AtomicInteger(0);
+
+ public AtomicInteger getRerouteCounter() {
+ return rerouteCounter;
+ }
+
+ // 承诺标志,用于某种状态跟踪
+ private volatile boolean promise = true;
+
+ public void setPromise(boolean b) {
+ promise = b;
+ }
+
+ public boolean isPromise() {
+ return promise;
+ }
-
-
-
-}
+}
\ No newline at end of file
diff --git a/src/org/kne/cloud/network/ipv6/IPv6RouteTableKey.java b/src/org/kne/cloud/network/ipv6/IPv6RouteTableKey.java
new file mode 100644
index 0000000..0914cc3
--- /dev/null
+++ b/src/org/kne/cloud/network/ipv6/IPv6RouteTableKey.java
@@ -0,0 +1,305 @@
+package org.kne.cloud.network.ipv6;
+
+import java.net.Inet6Address;
+import java.util.Objects;
+
+public class IPv6RouteTableKey {
+ private long most;
+ private long least;
+
+ public IPv6RouteTableKey(byte[] address) {
+ long[] l = inet6AddressToLongs(address);
+ this.most = l[0];
+ this.least = l[1]; // 修复:应该是 l[1] 而不是 l[0]
+ }
+
+ public IPv6RouteTableKey(long most, long least) {
+ super();
+ this.most = most;
+ this.least = least;
+ }
+
+ public IPv6RouteTableKey(Inet6Address address) {
+ this(address.getAddress());
+ }
+
+ public IPv6RouteTableKey(Inet6AddressGroup address) {
+ this(address.getAddress());
+ applyMask(address.getPrefixLength());
+ }
+
+ public long getMost() {
+ return most;
+ }
+
+ public long getLeast() {
+ return least;
+ }
+
+ @Override
+ public int hashCode() {
+ final int prime = 31;
+ int result = 1;
+ result = prime * result + (int) (least ^ (least >>> 32));
+ result = prime * result + (int) (most ^ (most >>> 32));
+ return result;
+ }
+
+ @Override
+ public boolean equals(Object obj) {
+ if (this == obj)
+ return true;
+ if (obj == null)
+ return false;
+ if (getClass() != obj.getClass())
+ return false;
+ IPv6RouteTableKey other = (IPv6RouteTableKey) obj;
+ if (least != other.least)
+ return false;
+ if (most != other.most)
+ return false;
+ return true;
+ }
+
+
+ /**
+ * 根据前缀长度应用掩码,仅保留网络部分
+ * @param prefixLength 前缀长度 (0-128)
+ * @return 应用掩码后的新 IPv6RouteTableKey 对象
+ */
+ public IPv6RouteTableKey mask(int prefixLength) {
+ if (prefixLength < 0 || prefixLength > 128) {
+ throw new IllegalArgumentException("前缀长度必须在 0 到 128 之间");
+ }
+
+ if (prefixLength == 0) {
+ return new IPv6RouteTableKey(0L, 0L); // 默认路由
+ }
+
+ long maskedMost = this.most;
+ long maskedLeast = this.least;
+
+ if (prefixLength <= 64) {
+ // 仅影响高位
+ long mask = createMask(prefixLength);
+ maskedMost &= mask;
+ maskedLeast = 0L; // 低位全部清零
+ } else {
+ // 影响高位和部分低位
+ int lowPrefixLength = prefixLength - 64;
+ long lowMask = createMask(lowPrefixLength);
+ maskedLeast &= lowMask;
+ // 高位保持不变
+ }
+
+ return new IPv6RouteTableKey(maskedMost, maskedLeast);
+ }
+
+
+ /**
+ * 根据前缀长度应用掩码,仅保留网络部分
+ * @param prefixLength 前缀长度 (0-128)
+ * @return 应用掩码后的新 IPv6RouteTableKey 对象
+ */
+ public void applyMask(int prefixLength) {
+ if (prefixLength < 0 || prefixLength > 128) {
+ throw new IllegalArgumentException("前缀长度必须在 0 到 128 之间");
+ }
+
+ if (prefixLength == 0) {
+ this.most=0;
+ this.least=0;
+ return;
+ //return new IPv6RouteTableKey(0L, 0L); // 默认路由
+ }
+
+ long maskedMost = this.most;
+ long maskedLeast = this.least;
+
+ if (prefixLength <= 64) {
+ // 仅影响高位
+ long mask = createMask(prefixLength);
+ maskedMost &= mask;
+ maskedLeast = 0L; // 低位全部清零
+ } else {
+ // 影响高位和部分低位
+ int lowPrefixLength = prefixLength - 64;
+ long lowMask = createMask(lowPrefixLength);
+ maskedLeast &= lowMask;
+ // 高位保持不变
+ }
+ this.most=maskedMost;
+ this.least=maskedLeast;
+ // return new IPv6RouteTableKey(maskedMost, maskedLeast);
+ }
+
+ /**
+ * 创建指定长度的掩码
+ * @param bits 要保留的位数 (0-64)
+ * @return 掩码值
+ */
+ private long createMask(int bits) {
+ // 使用查表法,预先计算所有可能的掩码值
+ return MASK_TABLE[bits];
+ }
+
+ // 预计算的掩码表
+ private static final long[] MASK_TABLE = new long[65]; // 0-64 共65个值
+
+ // 静态初始化块,在类加载时预计算所有掩码值
+ static {
+ for (int bits = 0; bits <= 64; bits++) {
+ if (bits == 0) {
+ MASK_TABLE[bits] = 0L;
+ } else if (bits == 64) {
+ MASK_TABLE[bits] = -1L;
+ } else {
+ MASK_TABLE[bits] = (-1L) << (64 - bits);
+ }
+ }
+ }
+
+ /**
+ * 将 byte[] 转换为两个 long 值(高位和低位)
+ * @param bytes 16字节的IPv6地址
+ * @return 包含两个long值的数组,第一个是高位,第二个是低位
+ */
+ public static long[] inet6AddressToLongs(byte[] bytes) {
+ if (bytes.length != 16) {
+ throw new IllegalArgumentException("IPv6地址必须是16字节");
+ }
+
+ long high = 0;
+ long low = 0;
+
+ // 处理前8字节(高位)
+ for (int i = 0; i < 8; i++) {
+ high = (high << 8) | (bytes[i] & 0xFF);
+ }
+
+ // 处理后8字节(低位)
+ for (int i = 8; i < 16; i++) {
+ low = (low << 8) | (bytes[i] & 0xFF);
+ }
+
+ return new long[]{high, low};
+ }
+
+ /**
+ * 将两个long值转换回byte[]
+ * @param high 高位long值
+ * @param low 低位long值
+ * @return 16字节的IPv6地址
+ */
+ public static byte[] longsToInet6Address(long high, long low) {
+ byte[] bytes = new byte[16];
+
+ // 提取高位的8个字节
+ for (int i = 0; i < 8; i++) {
+ bytes[i] = (byte) ((high >> (56 - i * 8)) & 0xFF);
+ }
+
+ // 提取低位的8个字节
+ for (int i = 0; i < 8; i++) {
+ bytes[8 + i] = (byte) ((low >> (56 - i * 8)) & 0xFF);
+ }
+
+ return bytes;
+ }
+
+ @Override
+ public String toString() {
+ return String.format("%016x:%016x", most, least);
+ }
+
+ /**
+ * 单元测试
+ */
+ public static void main(String[] args) {
+ for(int i=0;i<65;i++) {
+ long l=MASK_TABLE[i];
+ System.out.println(Long.toUnsignedString(l, 16));
+ }
+
+ System.out.println("开始 IPv6RouteTableKey 单元测试...");
+
+ // 测试1: 基本转换测试
+ System.out.println("\n1. 测试基本转换:");
+ byte[] testAddress = new byte[16];
+ // 创建测试地址: 2001:0db8:85a3::8a2e:0370:7334
+ testAddress[0] = 0x20; testAddress[1] = 0x01;
+ testAddress[2] = 0x0d; testAddress[3] = (byte) 0xb8;
+ testAddress[4] = (byte) 0x85; testAddress[5] = (byte) 0xa3;
+ // 中间部分为0
+ testAddress[12] = (byte) 0x8a; testAddress[13] = 0x2e;
+ testAddress[14] = 0x03; testAddress[15] = 0x70;
+ // testAddress[16] = 0x73; testAddress[17] = 0x34; // 注意: 数组只有16个元素
+
+ IPv6RouteTableKey key = new IPv6RouteTableKey(testAddress);
+ System.out.println("原始地址: " + key);
+
+ // 测试2: 掩码应用测试
+ System.out.println("\n2. 测试掩码应用:");
+ IPv6RouteTableKey masked64 = key.mask(64);
+ System.out.println("/64 掩码: " + masked64);
+
+ IPv6RouteTableKey masked48 = key.mask(48);
+ System.out.println("/48 掩码: " + masked48);
+
+ IPv6RouteTableKey masked128 = key.mask(128);
+ System.out.println("/128 掩码: " + masked128);
+
+ IPv6RouteTableKey masked0 = key.mask(0);
+ System.out.println("/0 掩码: " + masked0);
+
+ // 测试3: 相等性测试
+ System.out.println("\n3. 测试相等性:");
+ IPv6RouteTableKey key2 = new IPv6RouteTableKey(testAddress);
+ System.out.println("相同地址是否相等: " + key.equals(key2));
+ System.out.println("哈希码是否相同: " + (key.hashCode() == key2.hashCode()));
+
+ // 测试4: 转换函数测试
+ System.out.println("\n4. 测试转换函数:");
+ long[] longs = inet6AddressToLongs(testAddress);
+ System.out.println("转换为longs: " + Long.toHexString(longs[0]) + ":" + Long.toHexString(longs[1]));
+
+ byte[] reconverted = longsToInet6Address(longs[0], longs[1]);
+ boolean conversionOk = true;
+ for (int i = 0; i < 16; i++) {
+ if (testAddress[i] != reconverted[i]) {
+ conversionOk = false;
+ break;
+ }
+ }
+ System.out.println("转换是否可逆: " + conversionOk);
+
+ // 测试5: 边界条件测试
+ System.out.println("\n5. 测试边界条件:");
+ try {
+ key.applyMask(-1);
+ System.out.println("错误: 应该抛出异常");
+ } catch (IllegalArgumentException e) {
+ System.out.println("正确: 负前缀长度抛出异常");
+ }
+
+ try {
+ key.applyMask(129);
+ System.out.println("错误: 应该抛出异常");
+ } catch (IllegalArgumentException e) {
+ System.out.println("正确: 过大前缀长度抛出异常");
+ }
+
+ // 测试6: 全零和全一地址测试
+ System.out.println("\n6. 测试特殊地址:");
+ byte[] allZeros = new byte[16];
+ IPv6RouteTableKey zeroKey = new IPv6RouteTableKey(allZeros);
+ System.out.println("全零地址: " + zeroKey);
+
+ byte[] allOnes = new byte[16];
+ for (int i = 0; i < 16; i++) allOnes[i] = (byte) 0xFF;
+ IPv6RouteTableKey onesKey = new IPv6RouteTableKey(allOnes);
+ System.out.println("全一地址: " + onesKey);
+
+ System.out.println("\n所有测试完成!");
+ }
+}
\ No newline at end of file
diff --git a/src/org/kne/cloud/network/ipv6/IPv6TUNLoopbackNetworkLink.java b/src/org/kne/cloud/network/ipv6/IPv6TUNLoopbackNetworkLink.java
index 5372959..68d9296 100644
--- a/src/org/kne/cloud/network/ipv6/IPv6TUNLoopbackNetworkLink.java
+++ b/src/org/kne/cloud/network/ipv6/IPv6TUNLoopbackNetworkLink.java
@@ -11,10 +11,12 @@ import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CopyOnWriteArraySet;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.LockSupport;
import java.util.concurrent.locks.ReentrantLock;
@@ -40,7 +42,7 @@ import org.kne.concurrent.HighPerformanceExecutor;
import org.kne.concurrent.SpinLock;
import org.kne.io.KNEChannels;
-public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, AutoCloseable {
+public class IPv6TUNLoopbackNetworkLink extends AbstractIPv6NetworkLink implements IPv6NetworkLink, Closeable, AutoCloseable {
public static final String KLALB_DECENTRALIZED_S_RV6_NETWORK = "KLALB Decentralized SRv6 Network";
@@ -76,7 +78,7 @@ public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, A
}
Thread tb = new Thread(() -> {
while (true) {
- ByteBuffer tmp = NetworkPacket.bufferAllocator.allocate(65535);
+ ByteBuffer tmp = NetworkPacket.bufferAllocator.allocateNative(65535);
try {
tun.read(tmp);
tmp.flip();
@@ -91,8 +93,8 @@ public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, A
//NetworkPacket.databufferpool_65535.back(tmp);
if (con != null) {
- monitor.getOutPacketCounterAL().incrementAndGet();
- monitor.getOutTrafficAL().addAndGet(ipp.getLength());
+ monitor.getOutPacketCounterAL().add(1);
+ monitor.getOutTrafficAL().add(ipp.getTotalLength());
// ipp.setDisposeAfterSend(true);
// ipp.getPayload().setDisposeAfterSend(true);
ipp.setPromise(true);
@@ -152,6 +154,7 @@ public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, A
}catch(UnsatisfiedLinkError e) {
throw new IOException(e);
}
+ onOnlineStateUpdate();
}
private LinkedBlockingQueue sendQueue = new LinkedBlockingQueue();
@@ -210,14 +213,15 @@ public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, A
SRv6PacketReorder newr= new SRv6PacketReorder(new PacketConsumer() {
@Override
- public void accept(IPv6Packet packx) throws IOException {
- ByteBuffer tmp = NetworkPacket.bufferAllocator.allocate(65535);
+ public boolean accept(IPv6Packet packx) throws IOException {
+ ByteBuffer tmp = NetworkPacket.bufferAllocator.allocateNative(65535);
packx.writeToChannel(KNEChannels.newWritableChannel(tmp));
- monitor.getInPacketCounterAL().incrementAndGet();
- monitor.getInTrafficAL().addAndGet(packx.getLength());
+ monitor.getInPacketCounterAL().add(1);
+ monitor.getInTrafficAL().add(packx.getTotalLength());
tmp.flip();
sendQueue.add(tmp);
LockSupport.unpark(tr);
+ return true;
}
});
SRv6PacketReorder olr=reorder.putIfAbsent(fss,newr);
@@ -230,8 +234,8 @@ public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, A
}else {
ByteBuffer tmp = NetworkPacket.bufferAllocator.allocate(65535);
pack.writeToChannel(KNEChannels.newWritableChannel(tmp));
- monitor.getInPacketCounterAL().incrementAndGet();
- monitor.getInTrafficAL().addAndGet(pack.getLength());
+ monitor.getInPacketCounterAL().add(1);
+ monitor.getInTrafficAL().add(pack.getTotalLength());
tmp.flip();
sendQueue.add(tmp);
LockSupport.unpark(tr);
@@ -282,11 +286,7 @@ public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, A
tun.close();
tun = null;
}
- }
-
- @Override
- public boolean canSend(IPv6Packet iPv6Packet) {
- return true;
+ onOnlineStateUpdate();
}
@Override
@@ -303,8 +303,8 @@ public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, A
}
@Override
- public Inet6AddressGroup getAddressGroup() {
- return hostAddress;
+ public List< Inet6AddressGroup> getAddressGroups() {
+ return List.of(hostAddress);
}
@Override
@@ -316,8 +316,10 @@ public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, A
public List getRouteItems() {
List rlist=new ArrayList<>();
- rlist.add(new RouteItem(new Inet6AddressGroup(this.getAddressGroup().getAddress(), 128),
- this.getAddressGroup().getAddress(), this, "Direct", 0, 1, null, "D",true));
+ for(Inet6AddressGroup grp:getAddressGroups()) {
+ rlist.add(new RouteItem(new Inet6AddressGroup(grp.getAddress(), 128),
+ grp.getAddress(), this, "Direct", 0, 1, null, "D",true));
+ }
for (Iterator iteratorx = getNeighborsInfo()
.iterator(); iteratorx.hasNext();) {
@@ -339,4 +341,11 @@ public class IPv6TUNLoopbackNetworkLink implements IPv6NetworkLink, Closeable, A
return false;
}
+ @Override
+ public void setCongressCondition(Lock lock, Condition condition) {
+ // TODO 自动生成的方法存根
+
+ }
+
+
}
diff --git a/src/org/kne/cloud/network/ipv6/Inet6AddressGroup.java b/src/org/kne/cloud/network/ipv6/Inet6AddressGroup.java
index fea8b16..e5036ba 100644
--- a/src/org/kne/cloud/network/ipv6/Inet6AddressGroup.java
+++ b/src/org/kne/cloud/network/ipv6/Inet6AddressGroup.java
@@ -8,6 +8,7 @@ import java.net.InetAddress;
import java.net.UnknownHostException;
import java.util.Arrays;
import java.util.Objects;
+import java.util.UUID;
public class Inet6AddressGroup implements Comparable{
private static final byte[][] maskTransf=new byte[129][16];
@@ -121,4 +122,21 @@ public class Inet6AddressGroup implements Comparable{
prefixLength=in.read();
}
+ public Inet6AddressGroup createPrefixOnlyAddressGroup(int newPrefixLength) {
+ byte[]mask=maskTransf[newPrefixLength];
+ byte[]andm=getAndm(mask);
+ try {
+ return new Inet6AddressGroup((Inet6Address) Inet6Address.getByAddress(andm), newPrefixLength);
+ } catch (UnknownHostException e) {
+ return null;
+ }
+ }
+ public Inet6AddressGroup createPrefixOnlyAddressGroup() {
+ return createPrefixOnlyAddressGroup(prefixLength);
+ }
+ public static void main(String[] args) throws UnknownHostException {
+ Inet6AddressGroup i6ag=new Inet6AddressGroup((Inet6Address) Inet6Address.getByName("1234:1234::1234"),12);
+ System.out.println(i6ag);
+ System.out.println(i6ag.createPrefixOnlyAddressGroup());
+ }
}
diff --git a/src/org/kne/cloud/network/ipv6/LoopbackIPv6NetworkLink.java b/src/org/kne/cloud/network/ipv6/LoopbackIPv6NetworkLink.java
new file mode 100644
index 0000000..410706c
--- /dev/null
+++ b/src/org/kne/cloud/network/ipv6/LoopbackIPv6NetworkLink.java
@@ -0,0 +1,162 @@
+package org.kne.cloud.network.ipv6;
+
+import java.io.IOException;
+import java.net.Inet6Address;
+import java.util.ArrayList;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.locks.Condition;
+import java.util.concurrent.locks.Lock;
+import java.util.function.Consumer;
+
+import org.kne.cloud.network.srv6.PacketConsumer;
+import org.kne.cloud.network.srv6.SRv6Router;
+
+public class LoopbackIPv6NetworkLink extends AbstractIPv6NetworkLink implements IPv6NetworkLink {
+ private IPv6NetworkLink fallbackLink;
+
+
+ public IPv6NetworkLink getFallbackLink() {
+ return fallbackLink;
+ }
+
+ public void setFallbackLink(IPv6NetworkLink fallbackLink) {
+ this.fallbackLink = fallbackLink;
+ }
+
+ private Map protocolNumberRegister = new ConcurrentHashMap<>(256 * 2);
+
+ public Map getProtocolNumberRegister() {
+ return protocolNumberRegister;
+ }
+
+ private List addressGroups =new ArrayList<>();
+
+ //new Inet6AddressGroup(loopbackAddress, 128) new Inet6AddressGroup((Inet6Address) Inet6Address.getByName("::1"), 128)
+ public LoopbackIPv6NetworkLink(List addressGroupsx,SRv6Router router) {
+ this.addressGroups .addAll( addressGroupsx);
+ }
+
+ public Consumer getReceiveConsumer() {
+ return receiveConsumer;
+ }
+
+ private Consumer receiveConsumer;
+ private Consumer rerouteConsumer;
+ private SRv6Router router;
+
+ @Override
+ public void sendPacket(IPv6Packet pack, Inet6Address next) throws IOException {
+ PacketConsumer pcm= protocolNumberRegister.get(pack.getPayload().getProtocolNumber());
+ if (pcm != null) {
+ if(!pcm.accept(pack)) {
+ if(fallbackLink!=null) {
+ fallbackLink.sendPacket(pack, next);
+ }
+ }
+ }else {
+ if(fallbackLink!=null) {
+ fallbackLink.sendPacket(pack, next);
+ }
+ }
+ }
+
+ @Override
+ public boolean isLoopBack() {
+ return true;
+ }
+
+ @Override
+ public List getNeighborsInfo() {
+ List hs = new ArrayList<>();
+ return hs;
+ }
+
+ @Override
+ public boolean isCongress(IPv6Packet iPv6Packet, double scale) {
+ return false;
+ }
+
+ @Override
+ public String getName() {
+ return "inLoopBack";
+ }
+
+ @Override
+ public boolean isUp() {
+ return true;
+ }
+
+ @Override
+ public void setReceiveConsumer(Consumer con) {
+ this.receiveConsumer = con;
+ if(fallbackLink!=null) {
+ fallbackLink.setReceiveConsumer(con);
+ }
+ }
+
+ @Override
+ public List getAddressGroups() {
+ return addressGroups;
+ }
+
+ @Override
+ public void setRerouteConsumer(Consumer rerouteConsumer) {
+ this.rerouteConsumer=rerouteConsumer;
+ if(fallbackLink!=null) {
+ fallbackLink.setRerouteConsumer(rerouteConsumer);
+ }
+ }
+
+ public Consumer getRerouteConsumer() {
+ return rerouteConsumer;
+ }
+
+ @Override
+ public List getRouteItems() {
+ List rlist = new ArrayList<>();
+ for(Inet6AddressGroup group:addressGroups) {
+ rlist.add(new RouteItem(new Inet6AddressGroup(group.getAddress(), 128),
+ group.getAddress(), this, "Direct", 0, 0, null, "D", true));
+ }
+
+ for (Iterator iteratorx = getNeighborsInfo().iterator(); iteratorx.hasNext();) {
+ Neighbor addresses = (Neighbor) iteratorx.next();
+ RouteItem ri = new RouteItem(new Inet6AddressGroup(addresses.getAddress().getAddress(), 128),
+ addresses.getAddress().getAddress(), this, "Direct", 0, 128, addresses.getMonitor(), "D",
+ false);
+ rlist.add(ri);
+
+ RouteItem ris = new RouteItem(addresses.getLocator(),
+ (Inet6Address) addresses.getLocator().getAddress(), this, "KLALB SRv6", 13, 128,
+ addresses.getMonitor(), "D", false);
+ rlist.add(ris);
+
+ }
+ return rlist;
+ }
+
+ @Override
+ public boolean isReachSpeedLimit(IPv6Packet iPv6Packet) {
+ return false;
+ }
+
+ public SRv6Router getRouter() {
+ return router;
+ }
+
+ public void setRouter(SRv6Router router) {
+ this.router = router;
+ }
+
+ @Override
+ public void setCongressCondition(Lock lock, Condition condition) {
+ // TODO 自动生成的方法存根
+
+ }
+
+
+
+ };
\ No newline at end of file
diff --git a/src/org/kne/cloud/network/ipv6/Neighbor.java b/src/org/kne/cloud/network/ipv6/Neighbor.java
index bdd53b3..0a42976 100644
--- a/src/org/kne/cloud/network/ipv6/Neighbor.java
+++ b/src/org/kne/cloud/network/ipv6/Neighbor.java
@@ -32,9 +32,9 @@ public class Neighbor {
this.locator = locator;
this.monitor = monitor;
}
- public Neighbor(Inet6AddressGroup peerAddress, Inet6AddressGroup remoteVaddr, QueueingMonitorDataImpl monitor2,
+ public Neighbor(Inet6AddressGroup address, Inet6AddressGroup locator, MonitorData monitor2,
BandwidthDistributer bandwidthDistributer) {
- this(peerAddress,remoteVaddr,monitor2);
+ this(address,locator,monitor2);
this.bandwidthDistributer=bandwidthDistributer;
}
@Override
diff --git a/src/org/kne/cloud/network/ipv6/Pad1HopByHopTLV.java b/src/org/kne/cloud/network/ipv6/Pad1HopByHopTLV.java
new file mode 100644
index 0000000..0eb94fa
--- /dev/null
+++ b/src/org/kne/cloud/network/ipv6/Pad1HopByHopTLV.java
@@ -0,0 +1,26 @@
+package org.kne.cloud.network.ipv6;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.nio.channels.ReadableByteChannel;
+import java.nio.channels.WritableByteChannel;
+
+import org.kne.cloud.network.srv6.IPv6SegmentRoutingTLV;
+
+/**
+ * Pad1 Hop-by-Hop TLV - 没有数据部分
+ */
+public class Pad1HopByHopTLV extends IPv6HopByHopTLV {
+ public Pad1HopByHopTLV() {
+ super(IPv6HopByHopTLV.PAD1);
+ }
+
+ public Pad1HopByHopTLV(ByteBuffer klalbHeader) {
+ super(klalbHeader);
+ }
+
+ @Override
+ public String toString() {
+ return "Pad1HopByHopTLV []";
+ }
+}
\ No newline at end of file
diff --git a/src/org/kne/cloud/network/ipv6/PadNHopByHopTLV.java b/src/org/kne/cloud/network/ipv6/PadNHopByHopTLV.java
new file mode 100644
index 0000000..fa946d3
--- /dev/null
+++ b/src/org/kne/cloud/network/ipv6/PadNHopByHopTLV.java
@@ -0,0 +1,30 @@
+package org.kne.cloud.network.ipv6;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.nio.channels.ReadableByteChannel;
+import java.nio.channels.WritableByteChannel;
+
+import org.kne.cloud.network.srv6.IPv6SegmentRoutingTLV;
+
+/**
+ * PadN Hop-by-Hop TLV - 有数据部分
+ */
+public class PadNHopByHopTLV extends IPv6HopByHopTLV {
+ public PadNHopByHopTLV(ByteBuffer klalbHeader) {
+ super(klalbHeader);
+ }
+
+ public PadNHopByHopTLV(int dataLength) {
+ super(PadNHopByHopTLV.PADN);
+ getData().limit(dataLength);
+ }
+
+ @Override
+ public String toString() {
+ return "PadNHopByHopTLV [getDataLength()=" + getDataLength() + "]";
+ }
+
+
+ // 继承的readFromChannel和writeToChannel会正确处理数据部分
+}
\ No newline at end of file
diff --git a/src/org/kne/cloud/network/ipv6/TLV.java b/src/org/kne/cloud/network/ipv6/TLV.java
index 1ad6d81..2fcc4d6 100644
--- a/src/org/kne/cloud/network/ipv6/TLV.java
+++ b/src/org/kne/cloud/network/ipv6/TLV.java
@@ -11,7 +11,32 @@ import org.kne.cloud.network.NetworkPacket;
public class TLV extends NetworkPacket{
private int headerLength=2;
-
+
+public static final int PAD1=0;
+ public TLV(ByteBuffer header, boolean isDefault) {
+ this.header=header;
+ this.isDefault=isDefault;
+ if(getType()==PAD1)
+ setHeaderLength(1);
+ if(isDefault&&(getHeaderLength()>1)) {
+ data=NetworkPacket.bufferAllocator.allocate(512);
+ }
+}
+
+ public TLV(int type, boolean isDefault) {
+ if(type==PAD1) {
+ setHeaderLength(1);
+ }
+ this.header=NetworkPacket.bufferAllocator.allocate(getHeaderLength());
+ this.isDefault=isDefault;
+ header.put((byte) type);
+ header.put((byte) 0);
+ header.flip();
+ if(isDefault&&(getHeaderLength()>1)) {
+ data=NetworkPacket.bufferAllocator.allocate(512);
+ }
+ }
+
public int getHeaderLength() {
return headerLength;
}
@@ -29,30 +54,31 @@ private int headerLength=2;
}
@Override
- public long getLength() {
- return headerLength+(isDefault?data.limit():0);
+ public long getTotalLength() {
+ return getHeaderLength()+(isDefault?data.limit():0);
}
@Override
public void writeToChannel(WritableByteChannel dto) throws IOException {
- if(isDefault&&(headerLength>1)) {
- setDataLength(data.limit());
- }
- header.limit(headerLength);
- dto.write(header.slice(0,header.limit()));
- if(isDefault&&(headerLength>1)) {
- dto.write(data.slice(0, data.limit()));
- }
+ if(isDefault&&(getHeaderLength()>1)) {
+ setDataLength(data.limit());
+ }
+ header.limit(getHeaderLength());
+ dto.write(header.slice(0,header.limit()));
+ if(isDefault&&(getHeaderLength()>1)) {
+ dto.write(data.slice(0, data.limit()));
+ }
}
@Override
public void readFromChannel(ReadableByteChannel din, long length) throws IOException {
+
header.limit(1);
while (header.hasRemaining()) {
if (din.read(header) == -1) {
throw new EOFException();
}
}
- if(headerLength>1) {
+ if(getHeaderLength()>1) {
header.limit(2);
while (header.hasRemaining()) {
if (din.read(header) == -1) {
@@ -62,7 +88,7 @@ private int headerLength=2;
}
header.flip();
- if(isDefault&&(headerLength>1)) {
+ if(isDefault&&(getHeaderLength()>1)) {
data.clear();
data.limit(getDataLength());
while(data.hasRemaining()){
@@ -93,4 +119,14 @@ private int headerLength=2;
public void setDataLength(int dataLength) {
header.put(1,(byte) dataLength);
}
+
+ public void setHeaderLength(int headerLength) {
+ this.headerLength = headerLength;
+ }
+
+ @Override
+ public String toString() {
+ return "TLV [getType()=" + getType() + ", getDataLength()=" + getDataLength() + "]";
+ }
+
}
diff --git a/src/org/kne/cloud/network/klalb/ACKTPacket.java b/src/org/kne/cloud/network/klalb/ACKTPacket.java
index 8b78f56..2293934 100644
--- a/src/org/kne/cloud/network/klalb/ACKTPacket.java
+++ b/src/org/kne/cloud/network/klalb/ACKTPacket.java
@@ -33,10 +33,10 @@ public class ACKTPacket extends KLALBPacket implements PortPacket {
@Override
public String toString() {
- return "ACKT "+getSport()+"->"+getDport()+" "+getNumber()+"[] avaliable:"+getAvaliableRcvWindow();
+ return "ACKT "+getSrcPort()+"->"+getDstPort()+" "+getNumber()+"[] avaliable:"+getAvaliableRcvWindow();
}
- public int getSport() {
+ public int getSrcPort() {
return klalbHeader.getInt(1);
}
@@ -44,7 +44,7 @@ public class ACKTPacket extends KLALBPacket implements PortPacket {
return klalbHeader.get(26)!=0;
}
- public int getDport() {
+ public int getDstPort() {
return klalbHeader.getInt(5);
}
diff --git a/src/org/kne/cloud/network/klalb/ADDLINESPacket.java b/src/org/kne/cloud/network/klalb/ADDLINESPacket.java
index 1b4b8aa..389df2b 100644
--- a/src/org/kne/cloud/network/klalb/ADDLINESPacket.java
+++ b/src/org/kne/cloud/network/klalb/ADDLINESPacket.java
@@ -63,7 +63,7 @@ public class ADDLINESPacket extends KLALBPacket {
}
@Override
- public long getLength() {
+ public long getTotalLength() {
return HEADER_LENGTH+lines.getBytes(Charset.forName("UTF-8")).length;
}
diff --git a/src/org/kne/cloud/network/klalb/AbstractKLALBPacketLink.java b/src/org/kne/cloud/network/klalb/AbstractKLALBPacketLink.java
index 632cec7..2f932bd 100644
--- a/src/org/kne/cloud/network/klalb/AbstractKLALBPacketLink.java
+++ b/src/org/kne/cloud/network/klalb/AbstractKLALBPacketLink.java
@@ -3,7 +3,9 @@ package org.kne.cloud.network.klalb;
import java.io.IOException;
import java.net.SocketException;
import java.nio.ByteBuffer;
+import java.util.List;
import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.LongAdder;
import org.kne.cloud.network.NetworkPacket;
import org.kne.io.KNEChannels;
@@ -19,76 +21,83 @@ public abstract class AbstractKLALBPacketLink implements KLALBPacketLink {
AbstractKLALBPacketLink.defaultSoTimeout = defaultSoTimeout;
}
- private AtomicLong[] inputTrafficCounters;
- private AtomicLong[] outputTrafficCounters;
- private AtomicLong[] inputPacketsCounters;
- private AtomicLong[] outputPacketsCounters;
+ private LongAdder[] inputTrafficCounters;
+ private LongAdder[] outputTrafficCounters;
+ private LongAdder[] inputPacketsCounters;
+ private LongAdder[] outputPacketsCounters;
@Override
- public void setOutputPacketsCounters(AtomicLong[] outCounter) {
+ public void setOutputPacketsCounters(LongAdder[] outCounter) {
outputPacketsCounters=outCounter;
}
@Override
- public void setInputPacketsCounters(AtomicLong[] inCounter) {
+ public void setInputPacketsCounters(LongAdder[] inCounter) {
inputPacketsCounters=inCounter;
}
@Override
- public void setOutputTrafficCounters(AtomicLong[] outCounter) {
+ public void setOutputTrafficCounters(LongAdder[] outCounter) {
this.outputTrafficCounters=outCounter;
}
@Override
- public void setInputTrafficCounters(AtomicLong[] inCounter) {
+ public void setInputTrafficCounters(LongAdder[] inCounter) {
this.inputTrafficCounters=inCounter;
}
- public AtomicLong[] getInputTrafficCounters() {
+ public LongAdder[] getInputTrafficCounters() {
return inputTrafficCounters;
}
- public AtomicLong[] getOutputTrafficCounters() {
+ public LongAdder[] getOutputTrafficCounters() {
return outputTrafficCounters;
}
- public AtomicLong[] getInputPacketsCounters() {
+ public LongAdder[] getInputPacketsCounters() {
return inputPacketsCounters;
}
- public AtomicLong[] getOutputPacketsCounters() {
+ public LongAdder[] getOutputPacketsCounters() {
return outputPacketsCounters;
}
protected void incOutput(int packetLength) {
if(outputTrafficCounters!=null) {
- for(AtomicLong al:outputTrafficCounters) {
- al.addAndGet(packetLength);
+ for(LongAdder al:outputTrafficCounters) {
+ al.add(packetLength);
}
}
if(outputPacketsCounters!=null) {
- for(AtomicLong al:outputPacketsCounters) {
- al.incrementAndGet();
+ for(LongAdder al:outputPacketsCounters) {
+ al.add(1);
}
}
}
protected void incInput(int packetLength) {
if(inputTrafficCounters!=null) {
- for(AtomicLong al:inputTrafficCounters) {
- al.addAndGet(packetLength);
+ for(LongAdder al:inputTrafficCounters) {
+ al.add(packetLength);
}
}
if(inputPacketsCounters!=null) {
- for(AtomicLong al:inputPacketsCounters) {
- al.incrementAndGet();
+ for(LongAdder al:inputPacketsCounters) {
+ al.add(1);
}
}
}
+ @Override
+ public void writeKLALBPackets(List kps) throws IOException {
+ for(KLALBPacket pack:kps) {
+ writeKLALBPacket(pack);
+ }
+ }
+
@Override
public void writeKLALBPacket(KLALBPacket kp) throws IOException {
- ByteBuffer dataWrite=NetworkPacket.bufferAllocator.allocate((int) kp.getLength());
- //dataWrite.clear();
+ ByteBuffer dataWrite=NetworkPacket.bufferAllocator.allocate((int) kp.getTotalLength());
+ dataWrite.clear();
KLALBPacket.writeKLALBPacketToChannel(KNEChannels.newWritableChannel( dataWrite), kp);
dataWrite.flip();
writePacket(dataWrite);
diff --git a/src/org/kne/cloud/network/klalb/BWINFPacket.java b/src/org/kne/cloud/network/klalb/BWINFPacket.java
index 6bc02b8..03d7c82 100644
--- a/src/org/kne/cloud/network/klalb/BWINFPacket.java
+++ b/src/org/kne/cloud/network/klalb/BWINFPacket.java
@@ -20,8 +20,8 @@ public class BWINFPacket extends KLALBPacket {
}
@Override
- public long getLength() {
- return super.getLength()+16;
+ public long getTotalLength() {
+ return super.getTotalLength();
}
public BWINFPacket(long upSpeed,long downSpeed) {
diff --git a/src/org/kne/cloud/network/klalb/BindableKLALBPacketConsumer.java b/src/org/kne/cloud/network/klalb/BindableKLALBPacketConsumer.java
index 95b9a98..913d48e 100644
--- a/src/org/kne/cloud/network/klalb/BindableKLALBPacketConsumer.java
+++ b/src/org/kne/cloud/network/klalb/BindableKLALBPacketConsumer.java
@@ -3,7 +3,7 @@ package org.kne.cloud.network.klalb;
import java.net.InetAddress;
import java.net.ServerSocket;
-public interface BindableKLALBPacketConsumer extends KLALBPacketConsumer {
+public interface BindableKLALBPacketConsumer extends PacketConsumer {
public InetAddress getRemoteInetAddress() ;
public InetAddress getLocalInetAddress() ;
public int getPort() ;
diff --git a/src/org/kne/cloud/network/klalb/CONST.java b/src/org/kne/cloud/network/klalb/CONST.java
index eef23b2..31fc93a 100644
--- a/src/org/kne/cloud/network/klalb/CONST.java
+++ b/src/org/kne/cloud/network/klalb/CONST.java
@@ -2,7 +2,7 @@ package org.kne.cloud.network.klalb;
public class CONST {
public static final String klalb="KLALB";
- public static final String klalbver="3.2";
+ public static final String klalbver="3.4";
public static final int bversion=3;
public static final int sversion=1;
public static final int itemwidth = 720;
diff --git a/src/org/kne/cloud/network/klalb/DATATPacket.java b/src/org/kne/cloud/network/klalb/DATATPacket.java
index 833e786..2b321fe 100644
--- a/src/org/kne/cloud/network/klalb/DATATPacket.java
+++ b/src/org/kne/cloud/network/klalb/DATATPacket.java
@@ -5,6 +5,7 @@ import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.channels.ReadableByteChannel;
import java.nio.channels.WritableByteChannel;
+import java.util.Arrays;
import org.kne.cloud.network.NetworkPacket;
@@ -17,7 +18,7 @@ public class DATATPacket extends KLALBPacket implements PortPacket{
private static final int HEADER_LENGTH=20;
@Override
- public long getLength() {
+ public long getTotalLength() {
return HEADER_LENGTH+dataBuffer.limit();
}
volatile long resendtimer=System.nanoTime();
@@ -54,11 +55,11 @@ public class DATATPacket extends KLALBPacket implements PortPacket{
super(bb,HEADER_LENGTH);
}
- public int getSport() {
+ public int getSrcPort() {
return klalbHeader.getInt(1);
}
- public int getDport() {
+ public int getDstPort() {
return klalbHeader.getInt(5);
}
@@ -75,7 +76,9 @@ public class DATATPacket extends KLALBPacket implements PortPacket{
@Override
public String toString() {
- return "DATAT "+getSport()+"->"+getDport()+" "+getNumber()+"["+getSize()+"]";
+ byte[]b=new byte[Math.min(dataBuffer.limit(),10)];
+ dataBuffer.get(0, b);
+ return "DATAT "+getSrcPort()+"->"+getDstPort()+" "+getNumber()+"["+getSize()+"] "+Arrays.toString(b);
}
@Override
diff --git a/src/org/kne/cloud/network/klalb/IPSequence.java b/src/org/kne/cloud/network/klalb/IPSequence.java
index 89fa913..1e7433b 100644
--- a/src/org/kne/cloud/network/klalb/IPSequence.java
+++ b/src/org/kne/cloud/network/klalb/IPSequence.java
@@ -30,7 +30,10 @@ public class IPSequence {
}
@Override
public int hashCode() {
- return Objects.hash(uuid);
+ final int prime = 31;
+ int result = 1;
+ result = prime * result + ((uuid == null) ? 0 : uuid.hashCode());
+ return result;
}
@Override
public boolean equals(Object obj) {
@@ -41,7 +44,12 @@ public class IPSequence {
if (getClass() != obj.getClass())
return false;
IPSequence other = (IPSequence) obj;
- return Objects.equals(uuid, other.uuid);
+ if (uuid == null) {
+ if (other.uuid != null)
+ return false;
+ } else if (!uuid.equals(other.uuid))
+ return false;
+ return true;
}
}
diff --git a/src/org/kne/cloud/network/klalb/IPv6OverKLALBPacket.java b/src/org/kne/cloud/network/klalb/IPv6OverKLALBPacket.java
index 2c18c92..9fa0f9f 100644
--- a/src/org/kne/cloud/network/klalb/IPv6OverKLALBPacket.java
+++ b/src/org/kne/cloud/network/klalb/IPv6OverKLALBPacket.java
@@ -13,15 +13,15 @@ public class IPv6OverKLALBPacket extends KLALBPacket {
@Override
- public long getLength() {
- return HEADER_LENGTH+ipv6Packet.getLength();
+ public long getTotalLength() {
+ return HEADER_LENGTH+ipv6Packet.getTotalLength();
}
private IPv6Packet ipv6Packet;
public IPv6OverKLALBPacket(IPv6Packet ipv6Packet) {
super(IPV6OVERKLALB,HEADER_LENGTH);
- klalbHeader.putLong((int) ipv6Packet.getLength());
+ klalbHeader.putLong((int) ipv6Packet.getTotalLength());
this.ipv6Packet=ipv6Packet;
}
@@ -41,19 +41,28 @@ public class IPv6OverKLALBPacket extends KLALBPacket {
return ipv6Packet;
}
- @Override
+ /*@Override
public String toString() {
- return "IPv6 "+ipv6Packet.getSourceAddress().getHostAddress()+"->"+ipv6Packet.getDestinationAddress().getHostAddress()+" type:"+ipv6Packet.getPayload().getProtocolNumber()+"["+getSize()+"]";
- }
+ return "IPv6OverKLALBPacket "+ipv6Packet.getSourceAddress().getHostAddress()+"->"+ipv6Packet.getDestinationAddress().getHostAddress()+" type:"+ipv6Packet.getPayload().getProtocolNumber()+"["+getSize()+"]";
+ }*/
+
+
public long getSize() {
- return ipv6Packet.getLength();
+ return ipv6Packet.getTotalLength();
}
+ @Override
+ public String toString() {
+ return "IPv6OverKLALBPacket [ipv6Packet=" + ipv6Packet + "]";
+ }
+
+
+
@Override
public void writeToChannel(WritableByteChannel dto) throws IOException {
- klalbHeader.putLong(1, ipv6Packet.getLength());
+ klalbHeader.putLong(1, ipv6Packet.getTotalLength());
super.writeToChannel(dto);
ipv6Packet.writeToChannel(dto);
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBConfig.java b/src/org/kne/cloud/network/klalb/KLALBConfig.java
index 756e059..d31e91f 100644
--- a/src/org/kne/cloud/network/klalb/KLALBConfig.java
+++ b/src/org/kne/cloud/network/klalb/KLALBConfig.java
@@ -32,7 +32,7 @@ import java.io.FileNotFoundException;
public class KLALBConfig extends ArrayList{
public static void main(String[] args) throws UnknownHostException, JsonSyntaxException, JsonIOException, FileNotFoundException {
- GsonBuilder gb=new GsonBuilder();
+ GsonBuilder gb=new GsonBuilder().setPrettyPrinting();
MultipurposeSocketAddress.registerToGsonBuilder(gb);
KLALBConfigItem.registerToGsonBuilder(gb);
Gson gson=gb.create();
diff --git a/src/org/kne/cloud/network/klalb/KLALBController.java b/src/org/kne/cloud/network/klalb/KLALBController.java
index 312d9be..1baad1e 100644
--- a/src/org/kne/cloud/network/klalb/KLALBController.java
+++ b/src/org/kne/cloud/network/klalb/KLALBController.java
@@ -43,6 +43,7 @@ import java.util.function.BiConsumer;
import java.util.function.Consumer;
import org.kne.cloud.clock.AdjustedNanoClock;
+import org.kne.cloud.clock.HighAccuracyClock;
import org.kne.cloud.network.IPMulticastDiscovery;
import org.kne.cloud.network.MultipurposeSocketAddress;
import org.kne.cloud.network.PortPair;
@@ -54,11 +55,24 @@ import org.kne.cloud.network.ipv6.IPv6Packet.IPv6ExtHeader;
import org.kne.cloud.network.ipv6.IPv6Packet.IPv6Payload;
import org.kne.cloud.network.ipv6.IPv6TUNLoopbackNetworkLink;
import org.kne.cloud.network.ipv6.Inet6AddressGroup;
+import org.kne.cloud.network.ipv6.LoopbackIPv6NetworkLink;
+import org.kne.cloud.network.ipv6.Neighbor;
import org.kne.cloud.network.monitor.MonitorData;
import org.kne.cloud.network.monitor.SpeedAndTrafficMonitorDataImpl;
+import org.kne.cloud.network.monitor.TimestampMonitor;
+import org.kne.cloud.network.ntp.NTPContext;
+import org.kne.cloud.network.ntp.NTPv4Packet;
+import org.kne.cloud.network.ntp.NTPv4Protocol;
+import org.kne.cloud.network.ntp.NTPv4Protocol.NTPPeer;
+import org.kne.cloud.network.srv6.JsonDataPacket;
import org.kne.cloud.network.srv6.KLALBRoutingProtocol;
+import org.kne.cloud.network.srv6.KLALBRoutingProtocolAPIClient;
+import org.kne.cloud.network.srv6.KLALBRoutingProtocolAPIServer;
+import org.kne.cloud.network.srv6.KLALBRoutingProtocolJsonData;
import org.kne.cloud.network.srv6.PacketConsumer;
import org.kne.cloud.network.srv6.SRv6Router;
+import org.kne.cloud.network.srv6.SRv6RouterListener;
+import org.kne.cloud.network.tcp.UDPPacket;
import org.kne.cloud.network.te.BandwidthDistributer;
import org.kne.cloud.network.te.DWRRLoadingBalanceAlgorithm;
import org.kne.concurrent.DisruptorExecutor;
@@ -66,239 +80,314 @@ import org.kne.concurrent.HighPerformanceExecutor;
import org.kne.io.KNEChannels;
import org.pcap4j.packet.IpV6Packet.IpV6Header;
+import com.google.gson.JsonArray;
import com.google.gson.JsonElement;
public class KLALBController {
-
- private SpeedAndTrafficMonitorDataImpl linkMonitor=new SpeedAndTrafficMonitorDataImpl();
+ //TimestampMonitor<>
- private SpeedAndTrafficMonitorDataImpl datatMonitor=new SpeedAndTrafficMonitorDataImpl();
+ private SpeedAndTrafficMonitorDataImpl linkMonitor = new SpeedAndTrafficMonitorDataImpl();
+
+ private SpeedAndTrafficMonitorDataImpl datatMonitor = new SpeedAndTrafficMonitorDataImpl();
+
+ private List selflineTable = new ArrayList<>();
- private ListselflineTable=new ArrayList<>();
-
private static final boolean showpacket = false;
- private static final int PREFIX = 128;//112
- private static final int DISCOVERY_PORT=4569;
-
- private static List ipmd=new ArrayList<>();
-
- private ListdnsAddresses=new ArrayList<>();
-
+ private static final int PREFIX = 128;// 112
+ private static final int DISCOVERY_PORT = 4569;
+
+ private static List ipmd = new ArrayList<>();
+
+ private List networkInterfaceExcept = new ArrayList<>();
+
+ public List getNetworkInterfaceExcept() {
+ return networkInterfaceExcept;
+ }
+
+ private List dnsAddresses = new ArrayList<>();
+
public List getDnsAddresses() {
return dnsAddresses;
}
- private Timer twk=new Timer("网卡检测扫描计时器", true);
- {
-
- twk.schedule(new TimerTask() {
-
+ private Timer twk = new Timer("网卡检测扫描计时器", true);
+
+
+ private TimerTask tsk1=new TimerTask() {
+
@Override
public void run() {
-
+
try {
-
-
-
- Enumerationeu= NetworkInterface.getNetworkInterfaces();
+
+ Enumeration eu = NetworkInterface.getNetworkInterfaces();
while (eu.hasMoreElements()) {
NetworkInterface networkInterface = (NetworkInterface) eu.nextElement();
- if(networkInterface.getDisplayName().startsWith(IPv6TUNLoopbackNetworkLink.KLALB_DECENTRALIZED_S_RV6_NETWORK)) {
+
+ if (networkInterface.getDisplayName()
+ .startsWith(IPv6TUNLoopbackNetworkLink.KLALB_DECENTRALIZED_S_RV6_NETWORK)) {
continue;
}
- if(networkInterface.isUp()) {
- //System.out.println(networkInterface+" "+networkInterface.isUp());
- Enumerationei= networkInterface.getInetAddresses();
- while (ei.hasMoreElements()) {
- InetAddress inetAddress = (InetAddress) ei.nextElement();
- if(!inetAddress.isLoopbackAddress())
- for (Iterator iterator = listens.iterator(); iterator.hasNext();) {
- MultipurposeSocketAddress tcpl = (MultipurposeSocketAddress) iterator.next();
-
- try {
- if(tcpl.getInetAddress().isAnyLocalAddress()||tcpl.getInetAddress().equals(inetAddress)) {
- MultipurposeSocketAddress bind=new MultipurposeSocketAddress(tcpl.getType(),inetAddress.getHostAddress(),tcpl.getPort());
- //System.out.println(bind);
- synchronized (selflineTable) {
- if(!selflineTable.contains(bind)) {
- selflineTable.add(bind);
- }
-}
- }
- } catch (UnknownHostException e) {
- // TODO 自动生成的 catch 块
- e.printStackTrace();
- }
-
-
- }
+ if (networkInterfaceExcept.contains(networkInterface)) {
+ continue;
}
-
-
+ if (networkInterface.isUp()) {
+ // System.out.println(networkInterface+" "+networkInterface.isUp());
+ Enumeration ei = networkInterface.getInetAddresses();
+ while (ei.hasMoreElements()) {
+ InetAddress inetAddress = (InetAddress) ei.nextElement();
+ if (!inetAddress.isLoopbackAddress())
+ for (Iterator iterator = listens.iterator(); iterator
+ .hasNext();) {
+ MultipurposeSocketAddress tcpl = (MultipurposeSocketAddress) iterator.next();
+
+ try {
+ if (tcpl.getInetAddress().isAnyLocalAddress()
+ || tcpl.getInetAddress().equals(inetAddress)) {
+ MultipurposeSocketAddress bind = new MultipurposeSocketAddress(
+ tcpl.getType(), inetAddress.getHostAddress(), tcpl.getPort());
+ // System.out.println(bind);
+ synchronized (selflineTable) {
+ if (!selflineTable.contains(bind)) {
+ selflineTable.add(bind);
+ }
+ }
+ }
+ } catch (UnknownHostException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }
+
+ }
+ }
+
}
}
} catch (SocketException e) {
}
- lineslock.writeLock().lock();
+ lineslock.writeLock().lock();
try {
-
- Listlocaladdress=new ArrayList<>();
- Enumerationeu= NetworkInterface.getNetworkInterfaces();
+
+ List localaddress = new ArrayList<>();
+ Enumeration eu = NetworkInterface.getNetworkInterfaces();
while (eu.hasMoreElements()) {
NetworkInterface networkInterface = (NetworkInterface) eu.nextElement();
- if(networkInterface.isUp()) {
- //System.out.println(networkInterface+" "+networkInterface.isUp());
- Enumerationei= networkInterface.getInetAddresses();
- while (ei.hasMoreElements()) {
- InetAddress inetAddress = (InetAddress) ei.nextElement();
- MultipurposeSocketAddress bind=new MultipurposeSocketAddress(inetAddress.getHostAddress(),0);
- localaddress.add(bind);
-
- }
+ if (networkInterface.isUp()) {
+ // System.out.println(networkInterface+" "+networkInterface.isUp());
+ Enumeration ei = networkInterface.getInetAddresses();
+ while (ei.hasMoreElements()) {
+ InetAddress inetAddress = (InetAddress) ei.nextElement();
+ MultipurposeSocketAddress bind = new MultipurposeSocketAddress(
+ inetAddress.getHostAddress(), 0);
+ localaddress.add(bind);
+
+ }
}
}
- Setst=new HashSet<>();
- for (Iterator iterator = lines.iterator(); iterator.hasNext();) {
- KLALBRemoteLine multipurposeSocketAddress = (KLALBRemoteLine) iterator.next();
- if(multipurposeSocketAddress.getSocketAddress()!=null)
- st.add(multipurposeSocketAddress.getSocketAddress());
+ Set st = new HashSet<>();
+ for (Iterator iterator = srv6Router.getLinkTabel().iterator(); iterator.hasNext();) {
+ IPv6NetworkLink link=iterator.next();
+ if(link instanceof KLALBRemoteLine) {
+ KLALBRemoteLine multipurposeSocketAddress = (KLALBRemoteLine) link;
+ if (multipurposeSocketAddress.getSocketAddress() != null)
+ st.add(multipurposeSocketAddress.getSocketAddress());
+ }
}
-
-
+
for (Iterator iterator = st.iterator(); iterator.hasNext();) {
- MultipurposeSocketAddress target = (MultipurposeSocketAddress) iterator
- .next();
+ MultipurposeSocketAddress target = (MultipurposeSocketAddress) iterator.next();
addRemoteLines(target);
- /*for (Iterator iterator2 = localaddress.iterator(); iterator2.hasNext();) {
- MultipurposeSocketAddress bind = (MultipurposeSocketAddress) iterator2
- .next();
- try {
- if(target.getInetAddress().isLoopbackAddress() &&(!bind.getInetAddress().isLoopbackAddress())) {
- continue;
- }
- if((!target.getInetAddress().isLoopbackAddress()) &&bind.getInetAddress().isLoopbackAddress()) {
- continue;
- }
- if(target.getInetAddress()instanceof Inet4Address&&bind.getInetAddress() instanceof Inet6Address) {
- continue;
- }
- if(target.getInetAddress() instanceof Inet6Address &&bind.getInetAddress() instanceof Inet4Address) {
- continue;
- }
- } catch (UnknownHostException e) {
- if(e.getMessage().toLowerCase().contains("no scope_id found"))
- }
- if(!checkContainsTargetAndBind(target,bind)) {
- //System.out.println(target+" "+bind);
- addRemoteLine( new KLALBRemoteLine(target,bind));
- }
- }*/
-
+ /*
+ * for (Iterator iterator2 = localaddress.iterator();
+ * iterator2.hasNext();) { MultipurposeSocketAddress bind =
+ * (MultipurposeSocketAddress) iterator2 .next(); try {
+ * if(target.getInetAddress().isLoopbackAddress()
+ * &&(!bind.getInetAddress().isLoopbackAddress())) { continue; }
+ * if((!target.getInetAddress().isLoopbackAddress())
+ * &&bind.getInetAddress().isLoopbackAddress()) { continue; }
+ * if(target.getInetAddress()instanceof Inet4Address&&bind.getInetAddress()
+ * instanceof Inet6Address) { continue; } if(target.getInetAddress() instanceof
+ * Inet6Address &&bind.getInetAddress() instanceof Inet4Address) { continue; } }
+ * catch (UnknownHostException e) {
+ * if(e.getMessage().toLowerCase().contains("no scope_id found")) }
+ * if(!checkContainsTargetAndBind(target,bind)) {
+ * //System.out.println(target+" "+bind); addRemoteLine( new
+ * KLALBRemoteLine(target,bind)); } }
+ */
+
}
- for (Iterator iterator = lines.iterator(); iterator.hasNext();) {
- KLALBRemoteLine reml = (KLALBRemoteLine) iterator
- .next();
- if(reml.getSocketAddress()!=null) {
- if(reml.getBindAddress()!=null||reml.getRemoteVaddr()!=null)
- if(!localaddress.contains(reml.getBindAddress())) {
- reml.close();
- lines.remove(reml);
+ for (Iterator iterator = srv6Router.getLinkTabel().iterator(); iterator.hasNext();) {
+ IPv6NetworkLink link=iterator.next();
+ if(link instanceof KLALBRemoteLine) {
+ KLALBRemoteLine reml = (KLALBRemoteLine) link;
+ if (reml.getSocketAddress() != null) {
+ if (reml.getBindAddress() != null || reml.getRemoteVaddr() != null)
+ if (!localaddress.contains(reml.getBindAddress())) {
+ reml.close();
+ srv6Router.getLinkTabel().remove(reml);
+ }
}
}
}
-
+
for (Iterator iterator = ipmd.iterator(); iterator.hasNext();) {
IPMulticastDiscovery ipMulticastDiscovery = (IPMulticastDiscovery) iterator.next();
- if(ipMulticastDiscovery.isClosed()||(!ipMulticastDiscovery.getNinterface().isUp())) {
+ if (ipMulticastDiscovery.isClosed() || (!ipMulticastDiscovery.getNinterface().isUp())) {
iterator.remove();
try {
ipMulticastDiscovery.close();
} catch (IOException e) {
e.printStackTrace();
}
- System.out.println("移除网卡:"+ipMulticastDiscovery.getNinterface());
+ System.out.println("移除网卡:" + ipMulticastDiscovery.getNinterface());
}
}
- Enumerationeu2= NetworkInterface.getNetworkInterfaces();
+ Enumeration eu2 = NetworkInterface.getNetworkInterfaces();
while (eu2.hasMoreElements()) {
NetworkInterface networkInterface = (NetworkInterface) eu2.nextElement();
- if(networkInterface.getDisplayName().startsWith(IPv6TUNLoopbackNetworkLink.KLALB_DECENTRALIZED_S_RV6_NETWORK)) {
+ if (networkInterface.getDisplayName()
+ .startsWith(IPv6TUNLoopbackNetworkLink.KLALB_DECENTRALIZED_S_RV6_NETWORK)) {
continue;
}
- if(networkInterface.isUp()) {
-
- Enumerationei= networkInterface.getInetAddresses();
- loop:while (ei.hasMoreElements()) {
- InetAddress bidr = (InetAddress) ei.nextElement();
-
- for (Iterator iterator = ipmd.iterator(); iterator.hasNext();) {
- IPMulticastDiscovery ipMulticastDiscovery = (IPMulticastDiscovery) iterator.next();
- if(networkInterface.equals(ipMulticastDiscovery.getNinterface())&&bidr.equals(ipMulticastDiscovery.getBind().getAddress())) {
- continue loop;
- }
- }
-
-
-
- try {
- //InetAddress bidr=InetAddress.getByName("::0");
- if(bidr instanceof Inet6Address) {
- IPMulticastDiscovery ipd=new IPMulticastDiscovery(new InetSocketAddress(bidr,DISCOVERY_PORT), new InetSocketAddress(InetAddress.getByName("ff02::2486"),DISCOVERY_PORT), networkInterface,selflineTable,10000L);
- ipd.setCon((mpa)->{
- //System.out.println("添加本地IPv6链路:"+mpa);
- try {
- if(!checkIsSelf(mpa))
- addRemoteLines(mpa);
- } catch (UnknownHostException e) {
+ if (networkInterfaceExcept.contains(networkInterface)) {
+ continue;
+ }
+
+ if (networkInterface.isUp()) {
+
+ Enumeration ei = networkInterface.getInetAddresses();
+ loop: while (ei.hasMoreElements()) {
+ InetAddress bidr = (InetAddress) ei.nextElement();
+
+ for (Iterator iterator = ipmd.iterator(); iterator.hasNext();) {
+ IPMulticastDiscovery ipMulticastDiscovery = (IPMulticastDiscovery) iterator.next();
+ if (networkInterface.equals(ipMulticastDiscovery.getNinterface())
+ && bidr.equals(ipMulticastDiscovery.getBind().getAddress())) {
+ continue loop;
}
- });
- ipd.start();
- ipmd.add(ipd);
- }else if(bidr instanceof Inet4Address) {
- //bidr=InetAddress.getByName("0.0.0.0");
- IPMulticastDiscovery ipd2=new IPMulticastDiscovery(new InetSocketAddress(bidr,DISCOVERY_PORT), new InetSocketAddress(InetAddress.getByName("224.0.0.86"),DISCOVERY_PORT), networkInterface,selflineTable,10000L);
- ipd2.setCon((mpa)->{
- //System.out.println("添加本地IPv4链路:"+mpa);
- try {
- if(!checkIsSelf(mpa))
- addRemoteLines(mpa);
- } catch (UnknownHostException e) {
+ }
+
+ try {
+ // InetAddress bidr=InetAddress.getByName("::0");
+
+ if (bidr instanceof Inet6Address) {
+ IPMulticastDiscovery ipd = new IPMulticastDiscovery(
+ new InetSocketAddress(bidr, DISCOVERY_PORT),
+ new InetSocketAddress(InetAddress.getByName("ff02::2486"),
+ DISCOVERY_PORT),
+ networkInterface, selflineTable, 10000L);
+ ipd.setCon((mpa) -> {
+ // System.out.println("添加本地IPv6链路:"+mpa);
+ try {
+ if (!checkIsSelf(mpa))
+ addRemoteLines(mpa);
+ } catch (UnknownHostException e) {
+ }
+ });
+ ipd.start();
+ ipmd.add(ipd);
+ } else if (bidr instanceof Inet4Address) {
+ // bidr=InetAddress.getByName("0.0.0.0");
+ IPMulticastDiscovery ipd2 = new IPMulticastDiscovery(
+ new InetSocketAddress(bidr, DISCOVERY_PORT),
+ new InetSocketAddress(InetAddress.getByName("224.0.0.86"),
+ DISCOVERY_PORT),
+ networkInterface, selflineTable, 10000L);
+ ipd2.setCon((mpa) -> {
+ // System.out.println("添加本地IPv4链路:"+mpa);
+ try {
+ if (!checkIsSelf(mpa))
+ addRemoteLines(mpa);
+ } catch (UnknownHostException e) {
+ }
+ });
+ ipd2.start();
+ ipmd.add(ipd2);
}
- });
- ipd2.start();
- ipmd.add(ipd2);
+ System.out.println("添加网卡:" + networkInterface);
+ } catch (BindException e) {
+ // e.printStackTrace();
+ } catch (UnknownHostException e) {
+ e.printStackTrace();
+ } catch (IOException e) {
+ e.printStackTrace();
}
- System.out.println("添加网卡:"+networkInterface);
- }catch(BindException e) {
- //e.printStackTrace();
- } catch (UnknownHostException e) {
- e.printStackTrace();
- } catch (IOException e) {
- e.printStackTrace();
- }
-
-
- }
- //}
+
+ }
+ // }
}
}
-
+
} catch (SocketException e) {
- }finally {
+ } finally {
lineslock.writeLock().unlock();
}
-
}
private boolean checkIsSelf(MultipurposeSocketAddress inetAddress) throws UnknownHostException {
-
- return inetAddress.getInetAddress().isAnyLocalAddress()||inetAddress.getInetAddress().isLoopbackAddress()||selflineTable.contains(inetAddress);
+
+ return inetAddress.getInetAddress().isAnyLocalAddress()
+ || inetAddress.getInetAddress().isLoopbackAddress() || selflineTable.contains(inetAddress);
}
- }, 2000, 2000);
- }
+ };
+
+ private TimerTask tsk2=new TimerTask() {
+
+ @Override
+ public void run() {
+
+ List nt=ntptable;
+ //System.out.println("srv6 ip"+nb);
+ if(nvc2!=null) {
+ Set sp=nvc2.getPeers();
+ for (MultipurposeSocketAddress server : nt) {
+ sp.add(new NTPPeer(server, NTPv4Packet.NTP_CLIENT));
+
+ }
+ sp.removeIf((p)->{
+ for (MultipurposeSocketAddress server : nt) {
+ if(server.equals(p.getAddress())) {
+ return false;
+ }
+ }
+ return true;
+ });
+ }
+ }
+ };
+
+
+ private TimerTask tsk3=new TimerTask() {
+
+ @Override
+ public void run() {
+ if(srv6Router!=null&&routingProtocol!=null) {
+ Set nb=srv6Router.getLocators();
+ for (Iterator iterator = nb.iterator(); iterator.hasNext();) {
+ Inet6AddressGroup neighbor = (Inet6AddressGroup) iterator.next();
+ try {
+ InetSocketAddress iaddr= new InetSocketAddress( neighbor.getAddress(),KLALBRoutingProtocol.DEFAULT_PORT);
+ apiClient.requestOpenLines(iaddr,(v)->{
+ ThreadTool.makeVDaemonThreadIfSupport("线路添加任务", () -> {
+ for (MultipurposeSocketAddress msa:v) {
+ // System.out.print(msa);
+ addRemoteLines(msa);
+ }
+ }).start();
+ });
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+
+ }
+ }
+ }
+ };
+
+
public SpeedAndTrafficMonitorDataImpl getLinkMonitor() {
return linkMonitor;
}
@@ -307,246 +396,295 @@ public class KLALBController {
return datatMonitor;
}
- private static ConcurrentHashMapclocks=new ConcurrentHashMap<>();
+ private HighAccuracyClock clock;
- protected AdjustedNanoClock getAdjustedClockByVaddr(Inet6Address vaddr) {
- AdjustedNanoClock clk=clocks.get(vaddr);
- if(clk==null) {
- clk=new AdjustedNanoClock();
- clocks.put(vaddr, clk);
+ public HighAccuracyClock getClock() {
+ return clock;
+ }
+
+ private SocketType streamSocketType = new KLALBStreamSocketType();
+
+ public class KLALBStreamSocketType extends SocketType {
+
+ public KLALBStreamSocketType() {
+ super(new KLALBVirtualSocketFactory(KLALBController.this),
+ new KLALBVirtualServerSocketFactory(KLALBController.this),
+ new KLALBVirtualSocketChannelFactory(KLALBController.this),
+ new KLALBVirtualServerSocketChannelFactory(KLALBController.this));
}
- return clk;
+
}
+ public SocketType getStreamSocketType() {
+ return streamSocketType;
+ }
-
-
+ private SocketType datagramSocketType = new KLALBDatagramSocketType();
- private SocketType socketType= new KLALBSocketType();
-
- public class KLALBSocketType extends SocketType{
+ public class KLALBDatagramSocketType extends SocketType {
- public KLALBSocketType() {
- super(new KLALBVirtualSocketFactory(KLALBController.this), new KLALBVirtualServerSocketFactory(KLALBController.this),new KLALBVirtualSocketChannelFactory(KLALBController.this), new KLALBVirtualServerSocketChannelFactory(KLALBController.this));
+ public KLALBDatagramSocketType() {
+ super(new KLALBVirtualDatagramSocketFactory(KLALBController.this),
+ new KLALBVirtualDatagramServerSocketFactory(KLALBController.this));
+ //new KLALBVirtualDatagramSocketChannelFactory(KLALBController.this),
+ //new KLALBVirtualDatagramServerSocketChannelFactory(KLALBController.this));
}
-
- }
- public SocketType getSocketType() {
- return socketType;
+
}
-
-
- //private Supplier selflineTableSupplier=()->{return null;};
-
-
- /*public Supplier getSelflineTableSupplier() {
- return selflineTableSupplier;
+ public SocketType getDatagramSocketType() {
+ return datagramSocketType;
}
- public void setSelflineTableSupplier(Supplier selflineTableSupplier) {
- this.selflineTableSupplier = selflineTableSupplier;
- }*/
+ // private Supplier selflineTableSupplier=()->{return null;};
+
+ /*
+ * public Supplier getSelflineTableSupplier() { return
+ * selflineTableSupplier; }
+ *
+ * public void setSelflineTableSupplier(Supplier selflineTableSupplier)
+ * { this.selflineTableSupplier = selflineTableSupplier; }
+ */
public List getSelflineTable() {
return selflineTable;
}
-
- //private Inet6AddressGroup self;
+ // private Inet6AddressGroup self;
public Inet6AddressGroup getSelf() {
return srv6Router.getLocator();
}
- private List lines=new CopyOnWriteArrayList<>();
- private ReadWriteLock lineslock=new ReentrantReadWriteLock();
-
- public List getLines() {
- return lines;
+
+ private ReadWriteLock lineslock = new ReentrantReadWriteLock();
+
+ public List getLines() {
+ return srv6Router.getLinkTabel();
}
- private PortBinder streamPortBinder=new PortBinder(this);
+ private PortBinder streamPortBinder = new PortBinder(this);
- private PortBinder rawPortBinder=new PortBinder(this);
+ private PortBinder datagramPortBinder = new PortBinder(this);
+ private PortBinder rawPortBinder = new PortBinder(this);
+
protected PortBinder getStreamPortBinder() {
return streamPortBinder;
}
- public PortBinder getRawPortBinder() {
+ protected PortBinder getDatagramPortBinder() {
+ return datagramPortBinder;
+ }
+
+ protected PortBinder getRawPortBinder() {
return rawPortBinder;
}
public void reconnectImmediately() {
- for (Iterator iterator = lines.iterator(); iterator.hasNext();) {
- KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next();
+ for (Iterator iterator = srv6Router.getLinkTabel().iterator(); iterator.hasNext();) {
+ IPv6NetworkLink val=iterator.next();
+ if(val instanceof KLALBRemoteLine) {
+
+ KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) val;
klalbRemoteLine.reconnectImmediately();
}
+ }
}
-
- private PacketReceiver prc=new PacketReceiver();
- private class PacketReceiver implements Consumer{
+
+
+ private NTPv4Protocol nvc2;
+
+ private KLALBRoutingProtocolAPIServer apiServer;
+
+ private KLALBRoutingProtocolAPIClient apiClient;
+
+ /*private PacketReceiver prc = new PacketReceiver();
+ private class PacketReceiver implements Consumer {
@Override
- public void accept( KLALBPacket rec) {
+ public void accept(KLALBPacket rec) {
switch (rec.getType()) {
case KLALBPacket.ADDLINES:
- ADDLINESPacket lpt=(ADDLINESPacket) rec;
- String s=lpt.getLines();
- ThreadTool.makeVDaemonThreadIfSupport("线路添加任务", ()->{
- Scanner scn=new Scanner(s);
- while(scn.hasNext()) {
- String sn=scn.nextLine();
- MultipurposeSocketAddress msa= new MultipurposeSocketAddress(sn);
+ ADDLINESPacket lpt = (ADDLINESPacket) rec;
+ String s = lpt.getLines();
+ ThreadTool.makeVDaemonThreadIfSupport("线路添加任务", () -> {
+ Scanner scn = new Scanner(s);
+ while (scn.hasNext()) {
+ String sn = scn.nextLine();
+ MultipurposeSocketAddress msa = new MultipurposeSocketAddress(sn);
addRemoteLines(msa);
-
- }
+
+ }
}).start();
break;
default:
break;
-
+
}
}
-
+ }*/
+ public void addRemoteLines(List select) {
+ for(MultipurposeSocketAddress target:select) {
+ addRemoteLines(target);
+ }
}
-
- public List addRemoteLines(MultipurposeSocketAddress target) {
+ public List addRemoteLines(MultipurposeSocketAddress target) {
lineslock.writeLock().lock();
- try {
- Listadded=new ArrayList<>();
- try {
- Enumerationeu= NetworkInterface.getNetworkInterfaces();
- while (eu.hasMoreElements()) {
- NetworkInterface networkInterface = (NetworkInterface) eu.nextElement();
- if(networkInterface.isUp()) {
- //System.out.println(networkInterface+" "+networkInterface.isUp());
- Enumerationei= networkInterface.getInetAddresses();
+ try {
+ List added = new ArrayList<>();
+ try {
+ Enumeration eu = NetworkInterface.getNetworkInterfaces();
+ while (eu.hasMoreElements()) {
+ NetworkInterface networkInterface = (NetworkInterface) eu.nextElement();
+ if (networkInterface.isUp()) {
+ // System.out.println(networkInterface+" "+networkInterface.isUp());
+ Enumeration ei = networkInterface.getInetAddresses();
while (ei.hasMoreElements()) {
InetAddress inetAddress = (InetAddress) ei.nextElement();
try {
- MultipurposeSocketAddress bind=new MultipurposeSocketAddress(inetAddress.getHostAddress(),0);
- try {
- if(target.getInetAddress().isLoopbackAddress() &&(!bind.getInetAddress().isLoopbackAddress())) {
- continue;
+ MultipurposeSocketAddress bind = new MultipurposeSocketAddress(
+ inetAddress.getHostAddress(), 0);
+ try {
+ if (target.getInetAddress().isLoopbackAddress()
+ && (!bind.getInetAddress().isLoopbackAddress())) {
+ continue;
+ }
+ if ((!target.getInetAddress().isLoopbackAddress())
+ && bind.getInetAddress().isLoopbackAddress()) {
+ continue;
+ }
+ if (target.getInetAddress() instanceof Inet4Address
+ && bind.getInetAddress() instanceof Inet6Address) {
+ continue;
+ }
+ if (target.getInetAddress() instanceof Inet6Address
+ && bind.getInetAddress() instanceof Inet4Address) {
+ continue;
+ }
+ // System.out.println(target+"->"+bind);
+ } catch (UnknownHostException e) {
+ if (e.getMessage().trim().toLowerCase().contains("no scope_id found")) {
+ throw e;
+ }
+ // e.printStackTrace();
}
- if((!target.getInetAddress().isLoopbackAddress()) &&bind.getInetAddress().isLoopbackAddress()) {
- continue;
+ if (!checkContainsTargetAndBind(target, bind)) {
+ KLALBRemoteLine line = new KLALBRemoteLine(target, bind);
+ addRemoteLine(line);
+ added.add(line);
}
- if(target.getInetAddress()instanceof Inet4Address&&bind.getInetAddress() instanceof Inet6Address) {
- continue;
- }
- if(target.getInetAddress() instanceof Inet6Address &&bind.getInetAddress() instanceof Inet4Address) {
- continue;
- }
- //System.out.println(target+"->"+bind);
} catch (UnknownHostException e) {
- if(e.getMessage().trim().toLowerCase().contains("no scope_id found")) {
- throw e;
- }
- //e.printStackTrace();
+
}
- if(!checkContainsTargetAndBind(target,bind)) {
- KLALBRemoteLine line=new KLALBRemoteLine(target,bind);
- addRemoteLine(line );
- added.add(line);
- }
- } catch (UnknownHostException e) {
-
- }
- }
}
}
- } catch (SocketException e) {
- if(!checkContainsTarget(target)) {
- KLALBRemoteLine line= new KLALBRemoteLine(target);
+ }
+ } catch (SocketException e) {
+ if (!checkContainsTarget(target)) {
+ KLALBRemoteLine line = new KLALBRemoteLine(target);
addRemoteLine(line);
added.add(line);
- }
- //throw e;
}
-
- return added;
- }finally {
- lineslock.writeLock().unlock();
- }
+ // throw e;
+ }
+
+ return added;
+ } finally {
+ lineslock.writeLock().unlock();
+ }
}
+
public List removeRemoteLines(MultipurposeSocketAddress mpsa) {
- Listrmved=new ArrayList<>();
- for (Iterator iterator = lines.iterator(); iterator.hasNext();) {
- KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next();
- if(mpsa.equals( klalbRemoteLine.getSocketAddress())){
+ List rmved = new ArrayList<>();
+ for (Iterator iterator = srv6Router.getLinkTabel().iterator(); iterator.hasNext();) {
+ IPv6NetworkLink link= iterator.next();
+ if(link instanceof KLALBRemoteLine) {
+ KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine)link;
+ if (mpsa.equals(klalbRemoteLine.getSocketAddress())) {
klalbRemoteLine.close();
rmved.add(klalbRemoteLine);
}
+ }
}
return rmved;
}
+
public Inet6Address getRemoteVaddrBySocketAddress(MultipurposeSocketAddress target) throws SocketTimeoutException {
- KLALBRemoteLine kr=null;
- for (Iterator iterator = lines.iterator(); iterator.hasNext();) {
- KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next();
- if(target.equals(klalbRemoteLine.getSocketAddress())&&klalbRemoteLine.getBindAddress()==null) {
- kr=klalbRemoteLine;
+ KLALBRemoteLine kr = null;
+ for (Iterator iterator = srv6Router.getLinkTabel().iterator(); iterator.hasNext();) {
+ IPv6NetworkLink link= iterator.next();
+ if(link instanceof KLALBRemoteLine) {
+ KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) link;
+ if (target.equals(klalbRemoteLine.getSocketAddress()) && klalbRemoteLine.getBindAddress() == null) {
+ kr = klalbRemoteLine;
break;
}
+ }
}
-
- if(kr==null) {
- kr=new KLALBRemoteLine(target);
- addRemoteLine(kr);
- kr.waitForRemoteVaddrAvaliable(20000);
- }else {
+
+ if (kr == null) {
+ kr = new KLALBRemoteLine(target);
+ addRemoteLine(kr);
+ kr.waitForRemoteVaddrAvaliable(20000);
+ } else {
kr.reconnectImmediately();
kr.waitForRemoteVaddrAvaliable(20000);
}
return kr.getRemoteVaddr().getAddress();
}
- private boolean checkContainsTargetAndBind(MultipurposeSocketAddress target,MultipurposeSocketAddress bind) {
- boolean b=false;
- for (Iterator iterator = lines.iterator(); iterator.hasNext();) {
- KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next();
- if(bind.equals(klalbRemoteLine.getBindAddress())&&target.equals(klalbRemoteLine.getSocketAddress())) {
- b=true;
+
+ private boolean checkContainsTargetAndBind(MultipurposeSocketAddress target, MultipurposeSocketAddress bind) {
+ boolean b = false;
+ for (Iterator iterator = srv6Router.getLinkTabel().iterator(); iterator.hasNext();) {
+ IPv6NetworkLink link= iterator.next();
+ if(link instanceof KLALBRemoteLine) {
+ KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) link;
+ if (bind.equals(klalbRemoteLine.getBindAddress()) && target.equals(klalbRemoteLine.getSocketAddress())) {
+ b = true;
break;
}
+ }
}
return b;
}
private boolean checkContainsTarget(MultipurposeSocketAddress target) {
- boolean b=false;
- for (Iterator iterator = lines.iterator(); iterator.hasNext();) {
- KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next();
- if(target.equals(klalbRemoteLine.getSocketAddress())) {
- b=true;
+ boolean b = false;
+ for (Iterator iterator = srv6Router.getLinkTabel().iterator(); iterator.hasNext();) {
+ IPv6NetworkLink link= iterator.next();
+ if(link instanceof KLALBRemoteLine) {
+ KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) link;
+ if (target.equals(klalbRemoteLine.getSocketAddress())) {
+ b = true;
break;
}
}
+ }
return b;
}
- public void addRemoteLine(KLALBRemoteLine krs) {
+ public void addRemoteLine(KLALBRemoteLine krs) {
lineslock.writeLock().lock();
try {
- krs.setPacketReceiver(prc);
- krs.setKlalbController(this);
- krs.setLocalVaddrSupplier(()->{return getSelf();});
- krs.startIO();
- String selflineTable=generateSelfLineTable();
- if(selflineTable!=null&&!selflineTable.equals(""))
- krs.sendPacket(new ADDLINESPacket(selflineTable));
- lines.add(krs);
- }finally {
+ //krs.setPacketReceiver(prc);
+ krs.setKlalbController(this);
+ krs.setLocalVaddrSupplier(() -> {
+ return getSelf();
+ });
+ krs.startIO();
+ String selflineTable = generateSelfLineTable();
+ if (selflineTable != null && !selflineTable.equals(""))
+ krs.sendPacket(new ADDLINESPacket(selflineTable));
+ srv6Router.getLinkTabel().add(krs);
+ } finally {
lineslock.writeLock().unlock();
}
}
-
private String generateSelfLineTable() {
- StringBuilder sbd=new StringBuilder();
+ StringBuilder sbd = new StringBuilder();
for (Iterator iterator = selflineTable.iterator(); iterator.hasNext();) {
MultipurposeSocketAddress klalbRemoteLine = (MultipurposeSocketAddress) iterator.next();
sbd.append(klalbRemoteLine.toString());
@@ -554,348 +692,395 @@ public class KLALBController {
}
return sbd.toString();
}
-
- private class KLALBProtocolPacketConsumer implements PacketConsumer{
+
+ private class UDPProtocolPacketConsumer implements PacketConsumer {
@Override
- public void accept(IPv6Packet packx) throws IOException {
- //HighPerformanceExecutor.defaultExecutor.execute(()->{
-
- /*System.out.println("RCV:"+packx.getPayload().getData());
- byte[]b=new byte[packx.getPayload().getData().limit()];
- packx.getPayload().getData().get(0, b);
- System.out.println(Arrays.toString(b));*/
- IPv6Payload pl=packx.getPayload();
+ public boolean accept(IPv6Packet packx) throws IOException {
+
+ IPv6Payload pl = packx.getPayload();
packx.putTimePassport("unpacked");
packx.printPassport();
- if(pl instanceof KLALBPacket) {
- KLALBPacket rec=(KLALBPacket) pl;
- //KLALBPacket rec=KLALBPacket.createByBuffer(packx.getPayload().getData());
- if(showpacket)
- System.out.println("KLALB_RX:"+rec);
- if(rec instanceof PortPacket) {
- rec.setCE(packx.isCE());
- Inet6Address srcA=packx.getSourceAddress();
- PortPacket pt=(PortPacket) rec;
- if(!streamPortBinder.distributePacketToConsumer( srcA,pt)) {
- if(!(pt instanceof RSTPacket))
- try {
- sendPacketToAddress(srcA,0, new RSTPacket(pt.getDport(), pt.getSport()),2);
- } catch (IOException e) {
- e.printStackTrace();
- }
+ if (pl instanceof UDPPacket) {
+ UDPPacket rec = (UDPPacket) pl;
+ if (showpacket)
+ System.out.println("UDP_RX:" + rec);
+ Inet6Address srcA = packx.getSourceAddress();
+ PortPacket pt = (PortPacket) rec;
+ if(datagramPortBinder.distributePacketToConsumer(srcA, pt)) {
+ return true;
}
}
+ return false;
+ // packx.dispose();
+ // });
+ }
+
+ }
+
+ private class KLALBProtocolPacketConsumer implements PacketConsumer {
+
+ @Override
+ public boolean accept(IPv6Packet packx) throws IOException {
+ IPv6Payload pl = packx.getPayload();
+ packx.putTimePassport("unpacked");
+ packx.printPassport();
+ if (pl instanceof KLALBPacket) {
+ KLALBPacket rec = (KLALBPacket) pl;
+ if (showpacket)
+ System.out.println("KLALB_RX:" + rec);
+ if (rec instanceof PortPacket) {
+ rec.setCE(packx.isCE());
+ Inet6Address srcA = packx.getSourceAddress();
+ PortPacket pt = (PortPacket) rec;
+ if (streamPortBinder.distributePacketToConsumer(srcA, pt)) {
+ }else {
+ if (!(pt instanceof RSTPacket))
+ HighPerformanceExecutor.defaultExecutor.execute(() -> {
+ sendPacketToAddress(srcA, 0, new RSTPacket(pt.getDstPort(), pt.getSrcPort()), 2);
+ });
+ }
+ }
}
- //packx.dispose();
- //});
+ return true;
+ // packx.dispose();
}
-
+
}
- private void loadSRv6ProtocolStack(Inet6AddressGroup selfx,boolean enableVirtualAdapter) {
- srv6Router=new SRv6Router(selfx);
- if(enableVirtualAdapter)
- try {
- IPv6TUNLoopbackNetworkLink tunlink=new IPv6TUNLoopbackNetworkLink(new Inet6AddressGroup( srv6Router.getLocator().getAddress(),32),SRv6Router.MTU, dnsAddresses);
- tunlink.setMonitor(datatMonitor);
- srv6Router.getLinkTabel().add(tunlink);
- } catch (Exception e) {
- e.printStackTrace();
- }
+
+ private void loadSRv6ProtocolStack(Inet6AddressGroup selfx, boolean enableVirtualAdapter) {
+ srv6Router = new SRv6Router(selfx);
srv6Router.runKLALBRouteProtocol();
- srv6Router.getProtocolNumberRegister().put(KLALBPacket.KLALB_PROTOCOL_NUMBER,new KLALBProtocolPacketConsumer());
+ routingProtocol = srv6Router.getKlalbRouteProtol();
+ routingProtocol.addReceiver((addr,packet)->{
+
+
+ });
+
System.out.println("SRv6协议栈已加载");
- }
- public KLALBController() {
- this(true,null);
- }
- public KLALBController (Inet6Address self) {
- this(self,true,null);
- }
- public KLALBController(Inet6Address self,boolean enableVirtualAdapter,List dnsaddr) {
- this.dnsAddresses=dnsaddr;
- Inet6AddressGroup selfg =new Inet6AddressGroup(self, PREFIX);
- loadSRv6ProtocolStack(selfg,enableVirtualAdapter);
- }
- public KLALBController(boolean enableVirtualAdapter,List dnsaddr) {
- this.dnsAddresses=dnsaddr;
- SecureRandom sc=new SecureRandom();
- byte[]v=new byte[16];
- sc.nextBytes(v);
- v[0]=(byte)0x24;
- v[1]=(byte) 0x86;
- v[2]=0;
- v[3]=1;
- v[14]=0;
- v[15]=1;
- Inet6AddressGroup selfx=null;
+
+ if (enableVirtualAdapter)
+ try {
+ IPv6TUNLoopbackNetworkLink tunlink = new IPv6TUNLoopbackNetworkLink(
+ new Inet6AddressGroup(srv6Router.getLocator().getAddress(), 32), SRv6Router.MTU, dnsAddresses);
+ tunlink.setMonitor(datatMonitor);
+ //srv6Router.getLinkTabel().add(tunlink);
+ srv6Router.getinLoopback().setFallbackLink(tunlink);
+ System.out.println("虚拟网卡已挂载");
+ } catch (Exception e) {
+ System.err.println("挂载虚拟网卡失败,可能是无管理员权限?请尝试使用管理员权限运行软件");
+ e.printStackTrace();
+ }
+
+
+ srv6Router.getinLoopback(). getProtocolNumberRegister().put(UDPPacket.UDP_PROTOCOL_NUMBER, new UDPProtocolPacketConsumer());
+ System.out.println("UDP协议已加载");
+ srv6Router.getinLoopback().getProtocolNumberRegister().put(KLALBPacket.KLALB_PROTOCOL_NUMBER,
+ new KLALBProtocolPacketConsumer());
+ System.out.println("KLALB协议已加载");
+ String iproxyname="KLALB_"+UUID.randomUUID();
+ registerToProxyTypeAs(iproxyname);
+
+ clock = new HighAccuracyClock();
+ NTPContext context = new NTPContext(clock);
+ context.syncToSystem();
+ //System.out.println(context);
try {
- selfx=new Inet6AddressGroup( (Inet6Address) InetAddress.getByAddress(v),PREFIX);
+ NTPv4Protocol nvc = new NTPv4Protocol(context, new MultipurposeSocketAddress("{"+iproxyname+"_Datagram}0.0.0.0:123"));
+ srv6Router.addSRv6RouterListener(new SRv6RouterListener() {
+
+ @Override
+ public void onLinkChanged(SRv6Router router, IPv6NetworkLink link) {
+
+ List nb=router.getNeighbors();
+ //System.out.println("srv6 ip"+nb);
+ Set sp=nvc.getPeers();
+ for (Neighbor neighbor : nb) {
+ if(neighbor.getLocator()!=null)
+ sp.add(new NTPPeer(new MultipurposeSocketAddress(iproxyname+"_Datagram",neighbor.getLocator().getAddress().getHostAddress(),123), NTPv4Packet.NTP_SYMMETRIC_ACTIVE));
+
+ }
+ sp.removeIf((p)->{
+ for (Neighbor neighbor : nb) {
+ try {
+ if(neighbor.getLocator()!=null)
+ if(neighbor.getLocator().getAddress().equals(p.getAddress().getInetAddress())) {
+ return false;
+ }
+ } catch (UnknownHostException e) {
+ e.printStackTrace();
+ }
+ }
+ return true;
+ });
+ }
+ });
+ nvc2 = new NTPv4Protocol(context, new MultipurposeSocketAddress("{UDP}0.0.0.0:0"));// 106.55.184.199
+
+ System.out.println("NTP时钟同步协议已加载");
+ } catch (UnknownHostException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ } catch (IOException e) {
+ // TODO 自动生成的 catch 块
+ e.printStackTrace();
+ }// 106.55.184.199
+ }
+
+
+ private void loadController() {
+ this.apiServer=new KLALBRoutingProtocolAPIServer(routingProtocol, this);
+ this.apiClient=new KLALBRoutingProtocolAPIClient(routingProtocol);
+ twk.schedule(tsk1, 2000, 2000);
+ twk.schedule(tsk2, 2000, 2000);
+ twk.schedule(tsk3, 5000, 5000);
+ }
+
+ public KLALBController() {
+ this(true, null);
+ }
+
+ public KLALBController(Inet6Address self) {
+ this(self, true, null);
+ }
+
+ public KLALBController(Inet6Address self, boolean enableVirtualAdapter, List dnsaddr) {
+ this.dnsAddresses = dnsaddr;
+ Inet6AddressGroup selfg = new Inet6AddressGroup(self, PREFIX);
+ initKLALB(enableVirtualAdapter, selfg);
+ }
+
+ protected void initKLALB(boolean enableVirtualAdapter, Inet6AddressGroup selfg) {
+ loadSRv6ProtocolStack(selfg, enableVirtualAdapter);
+ loadController();
+ }
+
+
+ public KLALBController(boolean enableVirtualAdapter, List dnsaddr) {
+ this.dnsAddresses = dnsaddr;
+ SecureRandom sc = new SecureRandom();
+ byte[] v = new byte[16];
+ sc.nextBytes(v);
+ v[0] = (byte) 0x24;
+ v[1] = (byte) 0x86;
+ v[2] = 0;
+ v[3] = 1;
+ v[14] = 0;
+ v[15] = 1;
+ Inet6AddressGroup selfx = null;
+ try {
+ selfx = new Inet6AddressGroup((Inet6Address) InetAddress.getByAddress(v), PREFIX);
} catch (UnknownHostException e) {
e.printStackTrace();
}
- loadSRv6ProtocolStack(selfx,enableVirtualAdapter);
+ initKLALB(enableVirtualAdapter, selfx);
}
public KLALBController(List daddr) {
- this(true,daddr);
+ this(true, daddr);
}
public KLALBController(Inet6Address self, List daddr) {
- this(self,true,daddr);
+ this(self, true, daddr);
}
public KLALBController(boolean enableVirtualAdapter) {
- this(enableVirtualAdapter,null);
+ this(enableVirtualAdapter, null);
}
- protected KLALBVirtualSocketImpl createVirtualImpl() {
+ protected KLALBVirtualSocketImpl createVirtualSocketImpl() {
return new KLALBVirtualSocketImpl(this);
}
+
+ protected KLALBVirtualDatagramSocketImpl createVirtualDatagramSocketImpl() {
+ return new KLALBVirtualDatagramSocketImpl(this);
+ }
private SRv6Router srv6Router;
-
-
+
public SRv6Router getIpv6Router() {
return srv6Router;
}
-private Timer trtr=new Timer("路由表刷新计时器", true);
-{
- trtr.scheduleAtFixedRate(new TimerTask() {
-
- @Override
- public void run() {
- //Map> lines2 =new ConcurrentHashMap<>();
- for (Iterator iterator = lines.iterator(); iterator.hasNext();) {
- KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next();
- if(klalbRemoteLine.isClosed()) {
- lineslock.writeLock().lock();
- try {
- lines.remove(klalbRemoteLine);
- }finally{
- lineslock.writeLock().unlock();
- }
- }else {
- if(klalbRemoteLine.getRemoteVaddr()!=null&&klalbRemoteLine.getMonitor().getState()==MonitorData.ONLINE) {
- /* if(lines2.containsKey(klalbRemoteLine.getRemoteVaddr())) {
- lines2.get(klalbRemoteLine.getRemoteVaddr()).getEntries().add(klalbRemoteLine);
- }else {
- ArrayListal1=new ArrayList<>();
- al1.add(klalbRemoteLine);
- lines2.put(klalbRemoteLine.getRemoteVaddr().getAddress(), new DWRRLoadingBalanceAlgorithm<>(al1));
- }*/
- if(!srv6Router.getLinkTabel().contains(klalbRemoteLine)) {
- srv6Router.getLinkTabel().add(klalbRemoteLine);
- }
- }
- }
- }
- /*for (Iterator>> iterator = lines2.entrySet().iterator(); iterator.hasNext();) {
- Entry> klalbRemoteLine = (Entry>) iterator.next();
- List lineList=klalbRemoteLine.getValue().getEntries();
- lineList.forEach((r)->{
- r.runPredict();
- });
- Collections.sort(lineList);
- for (int i = 0; i < lineList.size(); i++) {
- KLALBRemoteLine krl=lineList.get(i);
- krl.noticeRank(i);
- }
-
- }
- KLALBController.this.lines2=lines2;*/
- if(srv6Router!=null) {
- srv6Router.getLinkTabel().removeIf((v)->{
- return (v instanceof KLALBRemoteLine)&&(!v.isUp());
- });
- srv6Router.updateRouteTabel();
- }
- }
- }, 2, 2);
- }
- /*private void updateLines2(Inet6Address addr) throws SocketTimeoutException {
- List l=new ArrayList();
- for (int i = 0; i < lines.size(); i++) {
- KLALBRemoteLine klalbRemoteLine = lines.get(i);
- if(klalbRemoteLine.isClosed()) {
- lines.remove(i);
- i--;
- continue;
- }
- if(addr.equals(klalbRemoteLine.getRemoteVaddr())&&klalbRemoteLine.getMonitor().getState()==MonitorData.ONLINE) {
- l.add(klalbRemoteLine);
- }
- }
-
- if(l.isEmpty()) {
- lines2.remove(addr);
- }else {
- l.forEach((r)->{
- r.runPredict();
+ private Timer trtr = new Timer("路由表刷新计时器", true);
+ {
+ trtr.scheduleAtFixedRate(new TimerTask() {
+
+ @Override
+ public void run() {
+ if (srv6Router != null) {
+ srv6Router.getLinkTabel().removeIf((v) -> {
+ return (v instanceof KLALBRemoteLine) && ((KLALBRemoteLine) v).isClosed();
});
- Collections.sort(l);
- lines2.put(addr, l);
+ srv6Router.updateRouteTabel();
}
- }*/
- //private Map> lines2 = new ConcurrentHashMap<>();
+ }
+ }, 100, 100);
+ }
-
- protected void sendPacketToAddress(Inet6Address addr,int flowlabel, KLALBPacket packet)
- throws IOException {
- sendPacketToAddress(addr,flowlabel, packet, 1);
+ protected void sendPacketToAddress(Inet6Address addr, int flowlabel, IPv6Packet.IPv6Payload packet) {
+ sendPacketToAddress(addr, flowlabel, packet, 1);
}
-
- protected void sendPacketToAddress(Inet6Address addr, int flowlabel,KLALBPacket packet, int count)
- throws IOException {
- HighPerformanceExecutor.defaultExecutor.execute(()->{
-
- IPv6Packet ipv=new IPv6Packet();
-
-
- if(showpacket)
- System.out.println("KLALB_TX:"+packet);
-
- ipv.setVersion(6);
- ipv.setTrafficClass(0);
- ipv.setFlowLabel(flowlabel);
- ipv.setHopLimit(255);
- ipv.setSourceAddress(getSelf().getAddress());
- ipv.setDestinationAddress(addr);
- ipv.setPriority(packet.getPriority());
- if(packet instanceof DATATPacket) {
+
+ protected void sendPacketToAddress(Inet6Address addr, int flowlabel, IPv6Packet.IPv6Payload packet, int count)
+ {
+ //HighPerformanceExecutor.defaultExecutor.execute(() -> {
+
+ IPv6Packet ipv = new IPv6Packet();
+
+ if (showpacket)
+ System.out.println("KLALB_TX:" + packet);
+
+ ipv.setVersion(6);
+ ipv.setTrafficClass(0);
+ ipv.setFlowLabel(flowlabel);
+ ipv.setHopLimit(255);
+ ipv.setSourceAddress(getSelf().getAddress());
+ ipv.setDestinationAddress(addr);
+ ipv.setPriority(packet.getPriority());
+ // if(packet instanceof DATATPacket) {
ipv.setPromise(true);
- }
- ipv.enableECN();
- /* if(b)
- ipv.markCE();*/
- ipv.setPayload(packet);
- //System.out.println(ipv.getPayload().getProtocolNumber());
- /*System.out.println("SND:"+ipv.getPayload().getData());
- byte[]b=new byte[ipv.getPayload().getData().limit()];
- ipv.getPayload().getData().get(0, b);
- System.out.println(Arrays.toString(b));*/
- //packet.getSendRecord().add(null);
- packet.incSendCounter();
-
- packet.putTimePassport("packedInIPv6");
- packet.printPassport();
- ipv.putTimePassport("packed");
- srv6Router.insertSRHandRoutePacket(ipv);
- //srv6Router.putProtocolNumberPacketAndInsertSRHAsync(ipv);
- });
+ // }
+ ipv.enableECN();
+ /*
+ * if(b) ipv.markCE();
+ */
+ ipv.setPayload(packet);
+ packet.setParent(ipv);
+ // System.out.println(ipv.getPayload().getProtocolNumber());
+ /*
+ * System.out.println("SND:"+ipv.getPayload().getData()); byte[]b=new
+ * byte[ipv.getPayload().getData().limit()]; ipv.getPayload().getData().get(0,
+ * b); System.out.println(Arrays.toString(b));
+ */
+ // packet.getSendRecord().add(null);
+ if(packet instanceof KLALBPacket)
+ ((KLALBPacket) packet).incSendCounter();
+
+ packet.putTimePassport("packedInIPv6");
+ packet.printPassport();
+ ipv.putTimePassport("packed");
+ srv6Router.insertSRHandRoutePacket(ipv);
+ // srv6Router.putProtocolNumberPacketAndInsertSRHAsync(ipv);
+ //});
}
-
-
- private Map> bandwidthDistrmap=new ConcurrentHashMap<>();
-
- private Lock bdmLock=new ReentrantLock();
-
+
+ private Map> bandwidthDistrmap = new ConcurrentHashMap<>();
+
+ private Lock bdmLock = new ReentrantLock();
+
protected Map> getBandwidthDistrmap() {
return bandwidthDistrmap;
}
-
- protected void registerDistUpdateConsumer(Inet6Address targetaAddress,PortPair portp, Consumer updateConsumer) {
- if(updateConsumer==null) {
- System.out.println("连接"+targetaAddress.getHostAddress()+" "+portp+" 释放带宽!");
- bdmLock.lock();
- try {
- BandwidthDistributer bdr= bandwidthDistrmap.get(targetaAddress);
- if(bdr!=null) {
- bdr.setDistrUpdateConsumer(portp, updateConsumer);
- bdr.setBandwidthRequest(portp, 0L);
- if(bdr.getDistrUpdateConsumerMap().isEmpty()) {
- bandwidthDistrmap.remove(targetaAddress);
- }
- }
- }finally {
- bdmLock.unlock();
- }
- }else {
- System.out.println("连接"+targetaAddress.getHostAddress()+" "+portp+" 申请带宽!");
- bdmLock.lock();
- try {
- BandwidthDistributer bdr= bandwidthDistrmap.get(targetaAddress);
- if(bdr==null) {
- bdr=new BandwidthDistributer<>(1024*20000L*1024);
- bdr.setTotalRequestUpdateConsumer((treq)->{
- KLALBRoutingProtocol krp= srv6Router.getKlalbRouteProtol();
- krp.updateTotalRequestBandwidth(targetaAddress,treq);});
- bandwidthDistrmap.put(targetaAddress, bdr);
- }
-
- bdr.setDistrUpdateConsumer(portp, updateConsumer);
- }finally {
- bdmLock.unlock();
- }
- }
- }
-
- protected void updateBandwidthRequest(Inet6Address targetaAddress,PortPair portp,long bandwidth) {
- if(bandwidth<0) {
- throw new IllegalArgumentException(bandwidth+"<0");
- }
- //System.out.println("连接"+targetaAddress.getHostAddress()+" "+portp+" 调整带宽到"+bandwidth/1024 +"KB/s!");
- BandwidthDistributer bdr= bandwidthDistrmap.get(targetaAddress);
-
- if(bdr!=null) {
- bdr.setBandwidthRequest(portp, bandwidth);
- }else {
- throw new NullPointerException("连接"+targetaAddress.getHostAddress()+" "+portp+" 未申请带宽!");
- }
-
- }
-
-private Timer t=new Timer("数据包重传计时器", true);
+ protected void registerDistUpdateConsumer(Inet6Address targetaAddress, PortPair portp,
+ Consumer updateConsumer) {
+ if (updateConsumer == null) {
+ System.out.println("连接" + targetaAddress.getHostAddress() + " " + portp + " 释放带宽!");
+ bdmLock.lock();
+ try {
+ BandwidthDistributer bdr = bandwidthDistrmap.get(targetaAddress);
+ if (bdr != null) {
+ bdr.setDistrUpdateConsumer(portp, updateConsumer);
+ bdr.setBandwidthRequest(portp, 0L);
+ if (bdr.getDistrUpdateConsumerMap().isEmpty()) {
+ bandwidthDistrmap.remove(targetaAddress);
+ }
+ }
+ } finally {
+ bdmLock.unlock();
+ }
+ } else {
+ System.out.println("连接" + targetaAddress.getHostAddress() + " " + portp + " 申请带宽!");
+ bdmLock.lock();
+ try {
+ BandwidthDistributer bdr = bandwidthDistrmap.get(targetaAddress);
+ if (bdr == null) {
+ bdr = new BandwidthDistributer<>(1024 * 20000L * 1024);
+ bdr.setTotalRequestUpdateConsumer((treq) -> {
+ KLALBRoutingProtocol krp = srv6Router.getKlalbRouteProtol();
+ krp.updateTotalRequestBandwidth(targetaAddress, treq);
+ });
+ bandwidthDistrmap.put(targetaAddress, bdr);
+ }
+
+ bdr.setDistrUpdateConsumer(portp, updateConsumer);
+ } finally {
+ bdmLock.unlock();
+ }
+ }
+ }
+
+ protected void updateBandwidthRequest(Inet6Address targetaAddress, PortPair portp, long bandwidth) {
+ if (bandwidth < 0) {
+ throw new IllegalArgumentException(bandwidth + "<0");
+ }
+ // System.out.println("连接"+targetaAddress.getHostAddress()+" "+portp+"
+ // 调整带宽到"+bandwidth/1024 +"KB/s!");
+ BandwidthDistributer bdr = bandwidthDistrmap.get(targetaAddress);
+
+ if (bdr != null) {
+ bdr.setBandwidthRequest(portp, bandwidth);
+ } else {
+ throw new NullPointerException("连接" + targetaAddress.getHostAddress() + " " + portp + " 未申请带宽!");
+ }
+
+ }
+
+ private Timer t = new Timer("数据包重传计时器", true);
+
protected Timer getResendTimer() {
return t;
}
-private Timer t2=new Timer("数据包粘包计时器", true);
-private List listens=new CopyOnWriteArrayList<>();
+ private Timer t2 = new Timer("数据包粘包计时器", true);
+
+ private List listens = new CopyOnWriteArrayList<>();
public Timer getNagleTimer() {
return t2;
}
-
-
-
+
public void registerToProxyTypeAs(String proxyname) {
- MultipurposeSocketAddress.getSocketTypeRegister().put(proxyname+"_Stream",socketType);
- //ProxyProfileEntry.getRegister().put(proxyname, this);
+ MultipurposeSocketAddress.getSocketTypeRegister().put(proxyname + "_Stream", streamSocketType);
+ MultipurposeSocketAddress.getSocketTypeRegister().put(proxyname + "_Datagram", datagramSocketType);
+ // ProxyProfileEntry.getRegister().put(proxyname, this);
}
-
public List getListenSocketAddress() {
return listens;
}
+ private Listntptable=new CopyOnWriteArrayList();
- /*@Override
- public void listen(MultipurposeSocketAddress msa) {
- // TODO 自动生成的方法存根
-
+ private KLALBRoutingProtocol routingProtocol;
+ public List getNTPTable() {
+ return ntptable;//nvc2.getPeers().add(new NTPPeer("{UDP}106.55.184.199:123", NTPv4Packet.NTP_CLIENT));
}
- @Override
- public void unlisten(MultipurposeSocketAddress msc) {
- // TODO 自动生成的方法存根
-
- }
-
- @Override
- public void connect(MultipurposeSocketAddress msa) {
- // TODO 自动生成的方法存根
-
- }
-
- @Override
- public void unconnect(MultipurposeSocketAddress msc) {
- // TODO 自动生成的方法存根
-
- }*/
-
+
+ /*
+ * @Override public void listen(MultipurposeSocketAddress msa) { // TODO
+ * 自动生成的方法存根
+ *
+ * }
+ *
+ * @Override public void unlisten(MultipurposeSocketAddress msc) { // TODO
+ * 自动生成的方法存根
+ *
+ * }
+ *
+ * @Override public void connect(MultipurposeSocketAddress msa) { // TODO
+ * 自动生成的方法存根
+ *
+ * }
+ *
+ * @Override public void unconnect(MultipurposeSocketAddress msc) { // TODO
+ * 自动生成的方法存根
+ *
+ * }
+ */
+
}
diff --git a/src/org/kne/cloud/network/klalb/KLALBControllerConfigItem.java b/src/org/kne/cloud/network/klalb/KLALBControllerConfigItem.java
index eb1199b..3d3c9c3 100644
--- a/src/org/kne/cloud/network/klalb/KLALBControllerConfigItem.java
+++ b/src/org/kne/cloud/network/klalb/KLALBControllerConfigItem.java
@@ -22,9 +22,17 @@ public class KLALBControllerConfigItem extends KLALBConfigItem {
private String VirtualSocketName;
private ListLineTable=new ArrayList<>();
private List