DatagramBuffer class
Пакет: com.hypixel.hytale.lib.quiche
Файл: com/hypixel/hytale/lib/quiche/QuicheListener.java
Поля (3)
| Модификаторы | Тип | Имя |
|---|---|---|
|
final ByteBuffer |
buffer |
|
final MemorySegment |
memory |
|
InetSocketAddress |
source |
Исходный код
Показать/скрыть
class="kw">package com.hypixel.hytale.lib.quiche;
class="kw">import com.hypixel.hytale.logger.HytaleLogger;
class="kw">import com.hypixel.hytale.protocol.io.ChannelConnection;
class="kw">import com.hypixel.hytale.protocol.io.ConnectionHandler;
class="kw">import com.hypixel.hytale.protocol.io.ServerListener;
class="kw">import it.unimi.dsi.fastutil.Hash.Strategy;
class="kw">import it.unimi.dsi.fastutil.objects.Object2ObjectOpenCustomHashMap;
class="kw">import it.unimi.dsi.fastutil.objects.ObjectArrayList;
class="kw">import java.io.IOException;
class="kw">import java.lang.foreign.Arena;
class="kw">import java.lang.foreign.MemorySegment;
class="kw">import java.lang.foreign.SegmentAllocator;
class="kw">import java.lang.foreign.ValueLayout;
class="kw">import java.net.InetSocketAddress;
class="kw">import java.net.ProtocolFamily;
class="kw">import java.net.SocketAddress;
class="kw">import java.net.StandardSocketOptions;
class="kw">import java.nio.ByteBuffer;
class="kw">import java.nio.channels.ClosedChannelException;
class="kw">import java.nio.channels.DatagramChannel;
class="kw">import java.nio.charset.StandardCharsets;
class="kw">import java.security.cert.X509Certificate;
class="kw">import java.time.Duration;
class="kw">import java.util.Arrays;
class="kw">import java.util.Locale;
class="kw">import java.util.Set;
class="kw">import java.util.concurrent.ArrayBlockingQueue;
class="kw">import java.util.concurrent.CompletableFuture;
class="kw">import java.util.concurrent.ConcurrentHashMap;
class="kw">import java.util.concurrent.TimeUnit;
class="kw">import java.util.concurrent.atomic.AtomicBoolean;
class="kw">import java.util.concurrent.atomic.AtomicLong;
class="kw">import java.util.function.BiConsumer;
class="kw">import java.util.function.BiFunction;
class="kw">import java.util.function.Function;
class="kw">import java.util.logging.Level;
class="kw">import javax.annotation.Nonnull;
class="kw">import javax.annotation.Nullable;
class="kw">import jdk.net.ExtendedSocketOptions;
class="kw">public class QuicheListener class="kw">implements ServerListener {
class="kw">private class="kw">static class="kw">final HytaleLogger LOGGER = HytaleLogger.forEnclosingClass();
class="kw">private class="kw">static class="kw">final int QUICHE_PROTOCOL_VERSION = 1;
class="kw">private class="kw">static class="kw">final int LOCAL_CONN_ID_LEN = 16;
class="kw">private class="kw">static class="kw">final boolean IS_LINUX = System.getProperty("os.name", "").toLowerCase(Locale.ROOT).contains("linux");
class="kw">static class="kw">final int MAX_UDP_DATAGRAM_SIZE = 65535;
class="kw">private class="kw">static class="kw">final int MAX_SEND_UDP_PAYLOAD_SIZE = 9000;
class="kw">static class="kw">final int BUFFER_POOL_SIZE = 256;
class="kw">private class="kw">static class="kw">final int STUN_SEND_QUEUE_SIZE = 256;
class="kw">private class="kw">static class="kw">final int SOCKET_RECEIVE_BUFFER_SIZE = 16777216;
class="kw">private class="kw">static class="kw">final int INITIAL_CONNECTION_WINDOW = 524288;
class="kw">private class="kw">static class="kw">final int INITIAL_STREAM_WINDOW = 131072;
class="kw">private class="kw">static class="kw">final int MAX_CONNECTION_WINDOW = 1048576;
class="kw">private class="kw">static class="kw">final int MAX_STREAM_WINDOW = 524288;
class="kw">private class="kw">static class="kw">final Duration SHUTDOWN_TIMEOUT = Duration.ofMillis(250L);
class="kw">private class="kw">final Arena arena = Arena.ofShared();
class="kw">final ArrayBlockingQueue<QuicheListener.DatagramBuffer> buffers = new ArrayBlockingQueue<>(256);
class="kw">final Set<QuicheConnection> liveConnections = ConcurrentHashMap.newKeySet();
class="kw">private class="kw">final AtomicBoolean closed = new AtomicBoolean(false);
class="kw">private class="kw">volatile boolean stopRequested;
class="kw">private class="kw">final Thread readerThread;
class="kw">private class="kw">final Thread stunWriterThread;
class="kw">private class="kw">final ArrayBlockingQueue<QuicheListener.OutboundDatagram> stunSendQueue = new ArrayBlockingQueue<>(256);
class="kw">final AtomicLong wireBytesReceived = new AtomicLong();
class="kw">final AtomicLong wireBytesSent = new AtomicLong();
@Nullable
class="kw">private class="kw">volatile BiConsumer<SocketAddress, ByteBuffer> stunHandler;
class="kw">private InetSocketAddress address;
class="kw">final DatagramChannel channel;
class="kw">private class="kw">final MemorySegment config;
class="kw">final QuicheConfig quicheConfig;
class="kw">final Function<ChannelConnection, ConnectionHandler> handler;
@Nullable
class="kw">final BiFunction<ConnectionHandler, ChannelConnection, ConnectionHandler> auxStreamHandler;
class="kw">private class="kw">final X509Certificate serverCertificate;
class="kw">final MemorySegment localAddr;
class="kw">final int localAddrLen;
class="kw">public void setStunHandler(@Nullable BiConsumer<SocketAddress, ByteBuffer> var1) {
this.stunHandler = var1;
}
class="kw">public void sendDatagram(@Nonnull ByteBuffer var1, @Nonnull SocketAddress var2) {
byte[] var3 = new byte[var1.remaining()];
var1.duplicate().get(var3);
if (!this.stunSendQueue.offer(new QuicheListener.OutboundDatagram(var3, var2))) {
((HytaleLogger.Api)((HytaleLogger.Api)LOGGER.atWarning()).atMostEvery(1, TimeUnit.MINUTES))
.log("STUN send queue is full; dropping a datagram to %s", var2);
}
}
class="kw">private class="kw">static boolean isStunDatagram(@Nonnull ByteBuffer var0) {
int var1 = var0.limit();
if (var1 < 20) {
class="kw">return false;
}
if ((var0.get(0) & 192) != 0) {
class="kw">return false;
}
if (var0.getInt(4) != 554869826) {
class="kw">return false;
}
int var2 = (var0.get(2) & 255) << 8 | var0.get(3) & 255;
class="kw">return var2 % 4 == 0 && 20 + var2 == var1;
}
class="kw">public QuicheListener(
ProtocolFamily var1,
InetSocketAddress var2,
@Nonnull QuicheServerCredentials var3,
@Nonnull QuicheConfig var4,
Function<ChannelConnection, ConnectionHandler> var5
) class="kw">throws IOException {
this(var1, var2, var3, var4, var5, null);
}
class="kw">public QuicheListener(
ProtocolFamily var1,
InetSocketAddress var2,
@Nonnull QuicheServerCredentials var3,
@Nonnull QuicheConfig var4,
Function<ChannelConnection, ConnectionHandler> var5,
@Nullable BiFunction<ConnectionHandler, ChannelConnection, ConnectionHandler> var6
) class="kw">throws IOException {
this.address = var2;
this.quicheConfig = var4;
this.handler = var5;
this.auxStreamHandler = var6;
this.serverCertificate = var3.certificate();
for (int var7 = 0; var7 < 256; var7++) {
MemorySegment var8 = this.arena.allocate(65535L);
this.buffers.add(new QuicheListener.DatagramBuffer(var8, var8.asByteBuffer()));
}
this.channel = DatagramChannel.open(var1);
try {
this.channel.setOption(StandardSocketOptions.SO_REUSEADDR, Boolean.TRUE);
this.channel.setOption(StandardSocketOptions.SO_RCVBUF, 16777216);
int var12 = this.channel.getOption(StandardSocketOptions.SO_RCVBUF);
int var13 = IS_LINUX ? var12 / 2 : var12;
if (var13 < 1048576) {
LOGGER.at(Level.WARNING)
.log(
"UDP receive buffer is %d bytes (effective %d), below the %d-byte connection window. Raise net.core.rmem_max so inbound bursts are not dropped",
var12,
var13,
1048576
);
}
this.channel.setOption(ExtendedSocketOptions.IP_DONTFRAGMENT, Boolean.TRUE);
this.channel.bind(var2);
this.address = (InetSocketAddress)this.channel.getLocalAddress();
this.localAddr = this.arena.allocate(QuicheNative.SOCKADDR_STORAGE_LAYOUT);
this.localAddrLen = QuicheUtil.toSocketAddrStorage(this.localAddr, this.address);
this.config = QuicheNative.configNew(1);
try {
initConfig(this.config, var3, var4);
} catch (Exception var10) {
QuicheNative.configFree(this.config);
throw var10;
}
} catch (IOException | RuntimeException var11) {
this.arena.close();
this.channel.close();
throw var11;
}
this.stunWriterThread = Thread.ofVirtual().name("quiche-stun-writer-" + var2).start(this::stunWriter);
this.readerThread = Thread.ofVirtual().name("quiche-listener-" + var2).start(this::reader);
}
class="kw">public X509Certificate getServerCertificate() {
class="kw">return this.serverCertificate;
}
class="kw">private class="kw">static void initConfig(MemorySegment var0, QuicheServerCredentials var1, QuicheConfig var2) {
try (Arena var3 = Arena.ofConfined()) {
String var4 = "hytale/3";
String var5 = "hytale/2";
MemorySegment var6 = buildAlpnWireFormat(var3, var4, var5);
int var7 = QuicheNative.configSetApplicationProtos(var0, var6, var6.byteSize());
if (var7 < 0) {
throw new RuntimeException("Failed to set ALPN protocols: " + QuicheNative.QuicheError.fromCode(var7));
}
MemorySegment var8 = var3.allocateFrom(ValueLayout.JAVA_BYTE, var1.certificateDer());
MemorySegment var9 = var3.allocateFrom(ValueLayout.JAVA_BYTE, var1.privateKeyDer());
var7 = QuicheNative.configLoadCert(var0, var8, var8.byteSize());
if (var7 < 0) {
throw new RuntimeException("Failed to load certificate: " + QuicheNative.QuicheError.fromCode(var7));
}
var7 = QuicheNative.configLoadPrivKey(var0, var9, var9.byteSize());
if (var7 < 0) {
throw new RuntimeException("Failed to load class="kw">private key: " + QuicheNative.QuicheError.fromCode(var7));
}
QuicheNative.configVerifyPeerOptional(var0);
QuicheNative.configSetMaxIdleTimeout(var0, var2.idleTimeout().toMillis());
applyTransportConfig(var0);
QuicheNative.configSetInitialMaxStreamsBidi(var0, 8L);
QuicheNative.configSetInitialMaxStreamsUni(var0, 0L);
}
}
class="kw">static void applyTransportConfig(MemorySegment var0) {
QuicheNative.configSetMaxSendUdpPayloadSize(var0, 9000L);
QuicheNative.configSetInitialMaxData(var0, 524288L);
QuicheNative.configSetInitialMaxStreamDataBidiLocal(var0, 131072L);
QuicheNative.configSetInitialMaxStreamDataBidiRemote(var0, 131072L);
QuicheNative.configSetInitialMaxStreamDataUni(var0, 131072L);
QuicheNative.configSetMaxConnectionWindow(var0, 1048576L);
QuicheNative.configSetMaxStreamWindow(var0, 524288L);
QuicheNative.configSetDisableActiveMigration(var0, true);
try (Arena var1 = Arena.ofConfined()) {
MemorySegment var2 = var1.allocateFrom("bbr2_gcongestion", StandardCharsets.US_ASCII);
int var3 = QuicheNative.configSetCcAlgorithmName(var0, var2);
if (var3 < 0) {
throw new RuntimeException("Failed to select congestion control: " + QuicheNative.QuicheError.fromCode(var3));
}
}
QuicheNative.configDiscoverPmtu(var0, true);
}
class="kw">static MemorySegment buildAlpnWireFormat(SegmentAllocator var0, String... var1) {
int var2 = 0;
for (String var6 : var1) {
var2 += 1 + var6.getBytes(StandardCharsets.UTF_8).length;
}
MemorySegment var10 = var0.allocate(var2);
int var11 = 0;
for (String var8 : var1) {
byte[] var9 = var8.getBytes(StandardCharsets.UTF_8);
var10.set(ValueLayout.JAVA_BYTE, var11, (byte)var9.length);
MemorySegment.copy(var9, 0, var10, ValueLayout.JAVA_BYTE, ++var11, var9.length);
var11 += var9.length;
}
class="kw">return var10;
}
class="kw">private void reader() {
Object2ObjectOpenCustomHashMap var1 = new Object2ObjectOpenCustomHashMap(new QuicheListener.DcidHash());
ObjectArrayList var2 = new ObjectArrayList();
try {
try (Arena var3 = Arena.ofConfined()) {
MemorySegment var4 = var3.allocate(QuicheNative.C_BYTE);
MemorySegment var5 = var3.allocate(QuicheNative.C_INT);
MemorySegment var6 = var3.allocate(20L);
MemorySegment var7 = var3.allocate(QuicheNative.C_SIZE_T);
MemorySegment var8 = var3.allocate(20L);
MemorySegment var9 = var3.allocate(QuicheNative.C_SIZE_T);
MemorySegment var10 = MemorySegment.NULL;
MemorySegment var11 = var3.allocate(QuicheNative.C_SIZE_T);
var11.fill((byte)0);
MemorySegment var12 = var3.allocate(QuicheNative.SOCKADDR_STORAGE_LAYOUT);
while (!this.stopRequested) {
for (int var13 = 0; var13 < var2.size(); var13++) {
byte[] var14 = (byte[])var2.get(var13);
QuicheConnection var15 = (QuicheConnection)var1.get(var14);
if (!var15.isActive()) {
var1.remove(var14);
var2.remove(var13);
var13--;
}
}
QuicheListener.DatagramBuffer var46;
try {
var46 = this.buffers.take();
} catch (InterruptedException var41) {
Thread.currentThread().interrupt();
class="kw">return;
}
var46.buffer.position(0);
var46.buffer.limit(var46.buffer.capacity());
InetSocketAddress var47 = (InetSocketAddress)this.channel.receive(var46.buffer);
var46.buffer.flip();
this.wireBytesReceived.addAndGet(var46.buffer.limit());
BiConsumer var48 = this.stunHandler;
if (var48 != null && isStunDatagram(var46.buffer)) {
try {
var48.accept(var47, var46.buffer.asReadOnlyBuffer());
} catch (Throwable var39) {
((HytaleLogger.Api)((HytaleLogger.Api)LOGGER.atWarning()).withCause(var39)).log("STUN handler threw; dropping the datagram");
} class="kw">finally {
this.buffers.add(var46);
}
} else {
QuicheNative.setSizeT(var7, 0L, 20L);
QuicheNative.setSizeT(var9, 0L, 20L);
QuicheNative.setSizeT(var11, 0L, 0L);
var8.fill((byte)0);
int var16 = QuicheNative.headerInfo(var46.memory, var46.buffer.limit(), 16L, var5, var4, var6, var7, var8, var9, var10, var11);
if (var16 < 0) {
this.buffers.add(var46);
} else {
QuicheConnection var17 = (QuicheConnection)var1.get(var8);
if (var17 == null) {
if (var4.get(ValueLayout.JAVA_BYTE, 0L) != 1) {
this.buffers.add(var46);
class="kw">continue;
}
byte[] var18 = var8.toArray(ValueLayout.JAVA_BYTE);
((HytaleLogger.Api)LOGGER.atInfo()).log("New connection: %s", Arrays.toString(var18));
int var19 = QuicheUtil.toSocketAddrStorage(var12, var47);
MemorySegment var20 = QuicheNative.accept(
var8, 20L, MemorySegment.NULL, 0L, this.localAddr, this.localAddrLen, var12, var19, this.config
);
if (var20.address() == 0L) {
this.buffers.add(var46);
class="kw">continue;
}
var17 = new QuicheConnection(this, var18, var20);
var1.put(var18, var17);
var2.add(var18);
this.liveConnections.add(var17);
var17.processingThread.start();
}
var17.submit(var47, var46);
}
}
}
} catch (ClosedChannelException var43) {
((HytaleLogger.Api)LOGGER.atInfo()).log("Reader channel closed, shutting down");
} catch (IOException var44) {
((HytaleLogger.Api)LOGGER.atSevere()).log("Fatal I/O error in reader thread: %s", var44.getMessage());
}
} class="kw">finally {
this.stopConnections();
}
}
class="kw">private boolean stopConnections() {
for (QuicheConnection var2 : this.liveConnections) {
var2.stop();
}
boolean var10 = true;
long var11 = System.nanoTime() + SHUTDOWN_TIMEOUT.toNanos();
for (QuicheConnection var5 : this.liveConnections) {
var5.stop();
long var6 = var11 - System.nanoTime();
try {
boolean var8 = var6 > 0L ? var5.processingThread.join(Duration.ofNanos(var6)) : !var5.processingThread.isAlive();
if (!var8) {
((HytaleLogger.Api)LOGGER.atWarning())
.log("Connection %s is still running after the %d ms shutdown budget", QuicheChannel.formatDcid(var5.dcid), SHUTDOWN_TIMEOUT.toMillis());
var10 = false;
}
} catch (InterruptedException var9) {
Thread.currentThread().interrupt();
class="kw">return false;
}
}
class="kw">return var10;
}
class="kw">private void stunWriter() {
try {
while (!Thread.currentThread().isInterrupted()) {
QuicheListener.OutboundDatagram var1 = this.stunSendQueue.take();
try {
int var2 = this.channel.send(ByteBuffer.wrap(var1.bytes()), var1.target());
this.wireBytesSent.addAndGet(var2);
} catch (ClosedChannelException var3) {
if (!this.closed.get()) {
((HytaleLogger.Api)((HytaleLogger.Api)LOGGER.atWarning()).withCause(var3)).log("STUN writer channel closed unexpectedly");
}
class="kw">return;
} catch (IOException var4) {
if (!this.closed.get()) {
((HytaleLogger.Api)((HytaleLogger.Api)((HytaleLogger.Api)LOGGER.atWarning()).atMostEvery(1, TimeUnit.MINUTES)).withCause(var4))
.log("STUN writer failed to send a datagram");
}
}
}
} catch (InterruptedException var5) {
Thread.currentThread().interrupt();
}
}
class="kw">public CompletableFuture<Void> close() {
class="kw">return !this.closed.compareAndSet(false, true)
? CompletableFuture.completedFuture(null)
: CompletableFuture.runAsync(
() -> {
this.stopRequested = true;
this.stopConnections();
try {
this.channel.close();
} catch (IOException var3) {
((HytaleLogger.Api)LOGGER.atWarning()).log("Failed to close datagram channel: %s", var3.getMessage());
}
this.readerThread.interrupt();
this.stunWriterThread.interrupt();
boolean var1 = true;
try {
if (!this.readerThread.join(SHUTDOWN_TIMEOUT)) {
((HytaleLogger.Api)LOGGER.atWarning()).log("Reader thread did not exit within %d ms", SHUTDOWN_TIMEOUT.toMillis());
var1 = false;
}
} catch (InterruptedException var4) {
Thread.currentThread().interrupt();
throw new RuntimeException(var4);
}
if (this.stopConnections() && var1) {
QuicheNative.configFree(this.config);
this.arena.close();
} else {
((HytaleLogger.Api)LOGGER.atSevere())
.log("Listener %s still has running threads, so its quiche config and buffer arena are leaked rather than freed", this.address);
}
}
);
}
@Override
class="kw">public SocketAddress localAddress() {
class="kw">return this.address;
}
@Override
class="kw">public long wireBytesReceived() {
class="kw">return this.wireBytesReceived.get();
}
@Override
class="kw">public long wireBytesSent() {
class="kw">return this.wireBytesSent.get();
}
@Override
class="kw">public String toString() {
class="kw">return "QuicheListener{address=" + this.address + ", channel=" + this.channel + "}";
}
void returnBuffer(QuicheListener.DatagramBuffer var1) {
this.buffers.add(var1);
}
class="kw">static class="kw">final class DatagramBuffer {
class="kw">final MemorySegment memory;
class="kw">final ByteBuffer buffer;
InetSocketAddress source;
DatagramBuffer(MemorySegment var1, ByteBuffer var2) {
this.memory = var1;
this.buffer = var2;
}
}
class="kw">private class="kw">static class DcidHash class="kw">implements Strategy<Object> {
class="kw">private DcidHash() {
}
class="kw">public int hashCode(Object var1) {
class="kw">return class="kw">switch (var1) {
case null -> 0;
case byte[] var4 -> {
int var10 = 0;
for (byte var9 : var4) {
var10 = 31 * var10 + var9;
}
yield var10;
}
case MemorySegment var5 -> {
int var6 = 0;
for (int var7 = 0; var7 < var5.byteSize(); var7++) {
var6 = 31 * var6 + var5.get(ValueLayout.JAVA_BYTE, var7);
}
yield var6;
}
class="kw">default -> throw new UnsupportedOperationException();
};
}
class="kw">public boolean equals(Object var1, Object var2) {
if (var1 class="kw">instanceof byte[] var7) {
if (var2 class="kw">instanceof byte[] var8) {
class="kw">return Arrays.equals(var7, var8);
} else if (var2 class="kw">instanceof MemorySegment var9) {
if (var7.length != var9.byteSize()) {
class="kw">return false;
}
for (int var10 = 0; var10 < var7.length; var10++) {
if (var7[var10] != var9.get(ValueLayout.JAVA_BYTE, var10)) {
class="kw">return false;
}
}
class="kw">return true;
} else {
class="kw">return false;
}
} else if (!(var1 class="kw">instanceof MemorySegment var3)) {
class="kw">return var1 == var2;
} else if (var2 class="kw">instanceof byte[] var4) {
if (var3.byteSize() != var4.length) {
class="kw">return false;
}
for (int var6 = 0; var6 < var4.length; var6++) {
if (var3.get(ValueLayout.JAVA_BYTE, var6) != var4[var6]) {
class="kw">return false;
}
}
class="kw">return true;
} else {
class="kw">return !(var2 class="kw">instanceof MemorySegment var5) ? false : var3.byteSize() == var5.byteSize() && var3.mismatch(var5) == -1L;
}
}
}
class="kw">private record OutboundDatagram(byte[] bytes, SocketAddress target) {
class="kw">private OutboundDatagram {
}
}
}