QuicheChannel class

Пакет: com.hypixel.hytale.lib.quiche

Файл: com/hypixel/hytale/lib/quiche/QuicheChannel.java

implements: ChannelConnection

Поля (57)

МодификаторыТипИмя
final HytaleLogger LOGGER
continue
return
return
return
return
boolean var1
var1
ToClientPacket var10
long var10
Class var11
Integer var13
MemorySegment var13
long var14
int var14
int var15
MemorySegment var2
int var2
long var2
Runnable var2
int var3
int var3
int var3
MemorySegment var3
long var3
String var3
CompletableFuture var3
long var4
int var4
var4
var4
MemorySegment var4
BooleanSupplier var4
QuicheChannel var4
PacketRegistry.PacketInfo var5
int var5
String var5
Runnable var5
String var5
QuicheChannel var5
int var6
long var6
long var6
long var6
String var6
long var7
long var7
TokenBucket var8
String var8
MemorySegment var9
var9
var9
Packet var9
long var9
var9
var9
FormattedMessage var9

Методы (42)

МодификаторыВозвратСигнатура
abstract throw new IllegalArgumentExceptionthrow new IllegalArgumentException("Invalid packet size: " + var4)
long checkTimersAndGetDeadlinelong checkTimersAndGetDeadline(long var1)
void closevoid close()
static void closeWithReasonvoid closeWithReason(@Nonnull MemorySegment var0, @Nonnull QuicApplicationErrorCode var1, @Nonnull FormattedMessage var2)
abstract closeWithReason closeWithReason(this.connection.connHandle, QuicApplicationErrorCode.NoError, var1)
abstract closeWithReason closeWithReason(this.connection.connHandle, QuicApplicationErrorCode.Timeout, var9)
abstract closeWithReason closeWithReason(this.connection.connHandle, var1, var2)
void consumePacketDatavoid consumePacketData(MemorySegment var1, int var2)
private void firePacketTimeoutvoid firePacketTimeout(long var1, long var3)
static String formatDcidString formatDcid(byte[] var0)
private void freeTempSendBuffervoid freeTempSendBuffer()
public long getStreamIdlong getStreamId()
if if(this.sendOffset < this.sendLimit)
if if(var4 < 0L)
if if(this.sendOffset < this.sendLimit)
if if(var10 == null)
if if(var13 == null)
if if(this.tempSendArena != null)
if if(!this.processingPacket)
if if(this.currentPacketSize != 0)
if if(this.currentPacketSize < 4)
if if(var2 - var3 < 4)
if if(var4 <= 0 || var4 > 1677721600)
if if(this.currentPacketOffset != this.currentPacketSize)
if if(this.tempPacketArena != null)
if if(this.handler != null)
if if(this.tempPacketArena != null)
if if(var3 < 0)
if if(var2 > 1000)
if if(var4 != null)
if if(this.loginTimingLastNs == 0L)
if if(this.stageDeadlineNs != 0L && var1 >= this.stageDeadlineNs)
if if(this.packetTimeoutNs != 0L && var1 - this.lastPacketReceivedNs >= this.packetTimeoutNs)
if if(this.packetTimeoutNs != 0L)
if if(this.stageDeadlineNs != 0L && this.stageDeadlineNs < var9)
if if(var2 != null)
private void logConnectionTimingNowvoid logConnectionTimingNow(@Nonnull String var1, @Nonnull Level var2)
private void runOnProcessingThreadvoid runOnProcessingThread(Runnable var1)
boolean tryWriteboolean tryWrite()
while while(true)
while while(var3 < var2)
void writeForClosevoid writeForClose(ToClientPacket var1)

Исходный код

Показать/скрыть
class="kw">package com.hypixel.hytale.lib.quiche;

class="kw">import com.hypixel.hytale.logger.HytaleLogger;
class="kw">import com.hypixel.hytale.protocol.CachedPacket;
class="kw">import com.hypixel.hytale.protocol.FormattedMessage;
class="kw">import com.hypixel.hytale.protocol.NetworkChannel;
class="kw">import com.hypixel.hytale.protocol.Packet;
class="kw">import com.hypixel.hytale.protocol.PacketRegistry;
class="kw">import com.hypixel.hytale.protocol.ToClientPacket;
class="kw">import com.hypixel.hytale.protocol.ToServerPacket;
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.ConnectionOrigin;
class="kw">import com.hypixel.hytale.protocol.io.PacketIO;
class="kw">import com.hypixel.hytale.protocol.io.PacketStatsRecorder;
class="kw">import com.hypixel.hytale.protocol.io.ProtocolException;
class="kw">import com.hypixel.hytale.protocol.io.ZstdNative;
class="kw">import com.hypixel.hytale.protocol.packets.connection.DisconnectType;
class="kw">import com.hypixel.hytale.protocol.packets.connection.QuicApplicationErrorCode;
class="kw">import com.hypixel.hytale.protocol.packets.connection.ServerDisconnect;
class="kw">import java.lang.foreign.Arena;
class="kw">import java.lang.foreign.MemorySegment;
class="kw">import java.lang.foreign.ValueLayout;
class="kw">import java.lang.foreign.ValueLayout.OfInt;
class="kw">import java.net.SocketAddress;
class="kw">import java.nio.ByteOrder;
class="kw">import java.security.cert.X509Certificate;
class="kw">import java.time.Duration;
class="kw">import java.util.Arrays;
class="kw">import java.util.HexFormat;
class="kw">import java.util.concurrent.CompletableFuture;
class="kw">import java.util.concurrent.ConcurrentLinkedQueue;
class="kw">import java.util.concurrent.TimeUnit;
class="kw">import java.util.concurrent.atomic.AtomicBoolean;
class="kw">import java.util.function.BiConsumer;
class="kw">import java.util.function.BooleanSupplier;
class="kw">import java.util.logging.Level;
class="kw">import javax.annotation.Nonnull;
class="kw">import javax.annotation.Nullable;

class="kw">public class QuicheChannel class="kw">implements ChannelConnection {
   class="kw">private class="kw">static class="kw">final HytaleLogger LOGGER = HytaleLogger.forEnclosingClass();
   class="kw">private class="kw">static class="kw">final OfInt PACKET_SIZE_LAYOUT = ValueLayout.JAVA_INT_UNALIGNED.withOrder(ByteOrder.LITTLE_ENDIAN);
   class="kw">static class="kw">final int MAX_CLOSE_REASON_BYTES = 1000;
   class="kw">private class="kw">final QuicheConnection connection;
   class="kw">private class="kw">final long streamId;
   class="kw">private class="kw">final NetworkChannel networkChannel;
   boolean isNew = true;
   ConnectionHandler handler;
   class="kw">private class="kw">final PacketStatsRecorder packetStatsRecorder = PacketStatsRecorder.NOOP;
   class="kw">private class="kw">final ConcurrentLinkedQueue<ToClientPacket> packetQueue = new ConcurrentLinkedQueue<>();
   class="kw">private class="kw">final Arena arena = Arena.ofConfined();
   class="kw">private class="kw">final MemorySegment packetBuffer = this.arena.allocate(102400L);
   class="kw">private boolean processingPacket = false;
   class="kw">private int currentPacketOffset = 0;
   class="kw">private int currentPacketSize = 0;
   @Nullable
   class="kw">private Arena tempPacketArena = null;
   @Nullable
   class="kw">private MemorySegment tempPacketBuffer = null;
   class="kw">private class="kw">final MemorySegment sendBuffer = this.arena.allocate(102400L);
   class="kw">private class="kw">final MemorySegment sendErrorCode = this.arena.allocate(QuicheNative.C_LONG);
   class="kw">private int sendOffset = 0;
   class="kw">private int sendLimit = 0;
   @Nullable
   class="kw">private Arena tempSendArena = null;
   @Nullable
   class="kw">private MemorySegment tempSendBuffer = null;
   class="kw">private class="kw">final AtomicBoolean writable = new AtomicBoolean(true);
   class="kw">private class="kw">final AtomicBoolean closed = new AtomicBoolean();
   class="kw">private long packetTimeoutNs;
   class="kw">private long lastPacketReceivedNs;
   class="kw">private long stageDeadlineNs;
   @Nullable
   class="kw">private BooleanSupplier stageCondition;
   @Nullable
   class="kw">private Runnable stageOnTimeout;
   @Nullable
   class="kw">private String currentStage;
   class="kw">private long connectionStartNs;
   @Nullable
   class="kw">private String playerIdentifier;
   class="kw">private long loginTimingLastNs;

   QuicheChannel(QuicheConnection var1, long var2, NetworkChannel var4) {
      this.connection = var1;
      this.streamId = var2;
      this.networkChannel = var4;
      this.lastPacketReceivedNs = System.nanoTime();
   }

   class="kw">public long getStreamId() {
      class="kw">return this.streamId;
   }

   boolean tryWrite() {
      boolean var1 = false;

      while (true) {
         if (this.sendOffset < this.sendLimit) {
            MemorySegment var2 = this.tempSendBuffer != null ? this.tempSendBuffer : this.sendBuffer;
            int var3 = this.sendLimit - this.sendOffset;
            long var4 = QuicheNative.connStreamSend(
               this.connection.connHandle, this.streamId, var2.asSlice(this.sendOffset, var3), var3, false, this.sendErrorCode
            );
            if (var4 == QuicheNative.QuicheError.DONE.code()) {
               this.writable.set(false);
               class="kw">return var1;
            }

            if (var4 < 0L) {
               long var14 = this.sendErrorCode.get(ValueLayout.JAVA_LONG, 0L);
               ((HytaleLogger.Api)LOGGER.atSevere())
                  .log("Error writing to stream %d: %s (%d) appErr=%d", this.streamId, QuicheNative.QuicheError.fromCode((int)var4), var4, var14);
               class="kw">return var1;
            }

            var1 = true;
            this.sendOffset += (int)var4;
            if (this.sendOffset < this.sendLimit) {
               this.writable.set(false);
               class="kw">return true;
            }

            this.sendOffset = 0;
            this.sendLimit = 0;
            this.freeTempSendBuffer();
         }

         ToClientPacket var10 = this.packetQueue.poll();
         this.writable.set(true);
         if (var10 == null) {
            class="kw">return var1;
         }

         Class var11 = (Class<? class="kw">extends ToClientPacket>)(var10 class="kw">instanceof CachedPacket var12 ? var12.getPacketType() : var10.getClass());
         Integer var13 = PacketRegistry.getId(var11);
         if (var13 == null) {
            throw new ProtocolException("Unknown packet type: " + var11.getName());
         }

         PacketRegistry.PacketInfo var5 = PacketRegistry.getToClientPacketById(var13);
         int var6 = var10.computeSize();
         long var7 = (var5.compressed() ? ZstdNative.compressBound(var6) : var6) + 8L;
         MemorySegment var9;
         if (var7 > this.sendBuffer.byteSize()) {
            this.tempSendArena = Arena.ofConfined();
            this.tempSendBuffer = this.tempSendArena.allocate(var7);
            var9 = this.tempSendBuffer;
         } else {
            var9 = this.sendBuffer;
         }

         this.sendLimit = PacketIO.writeFramedPacket(this.connection.compressionContext, var10, var5, var9, 0, var6, this.packetStatsRecorder);
         this.sendOffset = 0;
      }
   }

   class="kw">private void freeTempSendBuffer() {
      if (this.tempSendArena != null) {
         this.tempSendArena.close();
         this.tempSendArena = null;
         this.tempSendBuffer = null;
      }
   }

   void consumePacketData(MemorySegment var1, int var2) {
      int var3 = 0;

      while (var3 < var2) {
         if (!this.processingPacket) {
            int var4;
            if (this.currentPacketSize != 0) {
               int var5 = Math.min(var2 - var3, 4 - this.currentPacketSize);
               MemorySegment.copy(var1, var3, this.packetBuffer, this.currentPacketOffset, var5);
               this.currentPacketSize += var5;
               this.currentPacketOffset += var5;
               if (this.currentPacketSize < 4) {
                  class="kw">return;
               }

               var3 += var5;
               var4 = this.packetBuffer.get(PACKET_SIZE_LAYOUT, 0L) + 4;
            } else {
               if (var2 - var3 < 4) {
                  int var15 = var2 - var3;
                  MemorySegment.copy(var1, var3, this.packetBuffer, 0L, var15);
                  this.currentPacketSize += var15;
                  this.currentPacketOffset += var15;
                  class="kw">return;
               }

               var4 = var1.get(PACKET_SIZE_LAYOUT, var3) + 4;
               var3 += 4;
            }

            if (var4 <= 0 || var4 > 1677721600) {
               throw new IllegalArgumentException("Invalid packet size: " + var4);
            }

            if (var4 > this.packetBuffer.byteSize()) {
               ((HytaleLogger.Api)LOGGER.atWarning()).log("Packet size exceeds buffer capacity, allocating new buffer: %d", var4);
               this.tempPacketArena = Arena.ofConfined();
               this.tempPacketBuffer = this.tempPacketArena.allocate(var4);
            }

            this.currentPacketOffset = 0;
            this.currentPacketSize = var4;
         }

         this.processingPacket = true;
         MemorySegment var13 = this.tempPacketBuffer != null ? this.tempPacketBuffer : this.packetBuffer;
         int var14 = Math.min(var2 - var3, this.currentPacketSize - this.currentPacketOffset);
         MemorySegment.copy(var1, var3, var13, this.currentPacketOffset, var14);
         var3 += var14;
         this.currentPacketOffset += var14;
         if (this.currentPacketOffset != this.currentPacketSize) {
            class="kw">return;
         }

         try {
            long var6 = System.nanoTime();
            TokenBucket var8 = this.connection.rateLimitBucket;
            if (var8 == null || var8.tryConsume(var6)) {
               this.lastPacketReceivedNs = var6;
               Packet var9 = PacketIO.readFramedPacket(this.connection.decompressionContext, var13, this.currentPacketSize - 4, this.packetStatsRecorder);
               this.handler.handle((ToServerPacket)var9);
               class="kw">continue;
            }

            LOGGER.at(Level.WARNING).log("Rate limit exceeded for %s, disconnecting", this.formatRemoteAddress());
            if (!QuicheNative.connIsClosed(this.connection.connHandle)) {
               QuicheNative.connClose(this.connection.connHandle, true, QuicApplicationErrorCode.RateLimited.getValue(), MemorySegment.NULL, 0L);
            }
         } class="kw">finally {
            this.processingPacket = false;
            this.currentPacketSize = 0;
            this.currentPacketOffset = 0;
            if (this.tempPacketArena != null) {
               this.tempPacketArena.close();
               this.tempPacketArena = null;
               this.tempPacketBuffer = null;
            }
         }

         class="kw">return;
      }
   }

   void writeForClose(ToClientPacket var1) {
      this.packetQueue.add(var1);
      this.tryWrite();
   }

   void close() {
      if (this.closed.compareAndSet(false, true)) {
         if (this.handler != null) {
            this.handler.closed(this.networkChannel);
         }

         this.freeTempSendBuffer();
         if (this.tempPacketArena != null) {
            this.tempPacketArena.close();
         }

         this.arena.close();
      }
   }

   @Override
   class="kw">public void updateStreamPriority(int var1, boolean var2) {
      this.connection.execute(() -> {
         int var3 = QuicheNative.connStreamPriority(this.connection.connHandle, this.streamId, (byte)var1, var2);
         if (var3 < 0) {
            ((HytaleLogger.Api)LOGGER.atWarning()).log("Failed to set priority for stream %d: %d", this.streamId, var3);
         }
      });
   }

   @Override
   class="kw">public void flush() {
      this.connection.wake();
   }

   @Override
   class="kw">public void write(ToClientPacket var1) {
      this.packetQueue.add(var1);
   }

   @Override
   class="kw">public void writeAndFlush(ToClientPacket var1) {
      this.packetQueue.add(var1);
      this.connection.wake();
   }

   @Override
   class="kw">public void write(ToClientPacket[] var1) {
      this.packetQueue.addAll(Arrays.asList(var1));
   }

   @Override
   class="kw">public void writeAndFlush(ToClientPacket[] var1) {
      this.packetQueue.addAll(Arrays.asList(var1));
      this.connection.wake();
   }

   @Override
   class="kw">public boolean isActive() {
      class="kw">return !this.closed.get() && this.connection.isActive();
   }

   @Override
   class="kw">public boolean isWritable() {
      class="kw">return this.writable.get();
   }

   @Override
   class="kw">public SocketAddress remoteAddress() {
      class="kw">return this.connection.address;
   }

   @Override
   class="kw">public String formatRemoteAddress() {
      class="kw">return this.connection.address + "(" + formatDcid(this.connection.dcid) + ")";
   }

   class="kw">static String formatDcid(byte[] var0) {
      class="kw">return HexFormat.of().formatHex(var0);
   }

   @Nullable
   class="kw">static MemorySegment encodeCloseReason(@Nonnull FormattedMessage var0, @Nonnull Arena var1) {
      int var2 = var0.computeSize();
      if (var2 > 1000) {
         class="kw">return null;
      }

      MemorySegment var3 = var1.allocate(var2);
      var0.serialize(var3, 0);
      class="kw">return var3;
   }

   @Override
   class="kw">public void disconnect(@Nonnull FormattedMessage var1) {
      this.connection.execute(() -> {
         if (!QuicheNative.connIsClosed(this.connection.connHandle)) {
            this.packetQueue.add(new ServerDisconnect(var1, DisconnectType.Disconnect));
            this.tryWrite();
            closeWithReason(this.connection.connHandle, QuicApplicationErrorCode.NoError, var1);
         }
      });
   }

   class="kw">static void closeWithReason(@Nonnull MemorySegment var0, @Nonnull QuicApplicationErrorCode var1, @Nonnull FormattedMessage var2) {
      try (Arena var3 = Arena.ofConfined()) {
         MemorySegment var4 = encodeCloseReason(var2, var3);
         if (var4 != null) {
            QuicheNative.connClose(var0, true, var1.getValue(), var4, var4.byteSize());
         } else {
            QuicheNative.connClose(var0, true, var1.getValue(), MemorySegment.NULL, 0L);
         }
      }
   }

   @Nullable
   @Override
   class="kw">public PacketStatsRecorder getPacketStatsRecorder() {
      class="kw">return this.packetStatsRecorder;
   }

   @Nullable
   @Override
   class="kw">public String getSniHostname() {
      class="kw">return this.connection.serverName;
   }

   @Override
   class="kw">public boolean isFromSameOrigin(ChannelConnection var1) {
      class="kw">return ConnectionOrigin.isFromSameOrigin(this.remoteAddress(), var1.remoteAddress());
   }

   @Override
   class="kw">public void execute(Runnable var1) {
      this.connection.execute(var1);
   }

   @Override
   class="kw">public X509Certificate getClientCertificate() {
      class="kw">return this.connection.certificate;
   }

   @Override
   class="kw">public void initTimeoutContext(@Nonnull String var1, @Nonnull String var2) {
      this.runOnProcessingThread(() -> {
         this.currentStage = var1;
         this.playerIdentifier = var2;
         this.connectionStartNs = System.nanoTime();
         this.loginTimingLastNs = 0L;
      });
   }

   @Override
   class="kw">public void updateTimeoutContext(@Nonnull String var1, @Nonnull String var2) {
      this.runOnProcessingThread(() -> {
         this.currentStage = var1;
         this.playerIdentifier = var2;
      });
   }

   @Override
   class="kw">public void updateTimeoutContext(@Nonnull String var1) {
      this.runOnProcessingThread(() -> this.currentStage = var1);
   }

   @Override
   class="kw">public void setPacketTimeout(@Nonnull Duration var1) {
      long var2 = var1.toNanos();
      this.runOnProcessingThread(() -> {
         this.packetTimeoutNs = var2;
         this.lastPacketReceivedNs = System.nanoTime();
      });
   }

   @Override
   class="kw">public void clearPacketTimeout() {
      this.runOnProcessingThread(() -> this.packetTimeoutNs = 0L);
   }

   @Override
   class="kw">public void setStageTimeout(@Nonnull String var1, @Nonnull Duration var2, @Nonnull BooleanSupplier var3, @Nonnull Runnable var4) {
      this.runOnProcessingThread(() -> {
         this.currentStage = var1;
         this.stageDeadlineNs = System.nanoTime() + var2.toNanos();
         this.stageCondition = var3;
         this.stageOnTimeout = var4;
         this.logConnectionTimingNow("Entering stage '" + var1 + "'", Level.FINEST);
      });
   }

   @Override
   class="kw">public void clearStageTimeout() {
      this.runOnProcessingThread(() -> {
         this.stageDeadlineNs = 0L;
         this.stageCondition = null;
         this.stageOnTimeout = null;
      });
   }

   @Override
   class="kw">public void logConnectionTimings(@Nonnull String var1, @Nonnull Level var2) {
      this.runOnProcessingThread(() -> this.logConnectionTimingNow(var1, var2));
   }

   class="kw">private void logConnectionTimingNow(@Nonnull String var1, @Nonnull Level var2) {
      long var3 = System.nanoTime();
      String var5 = this.playerIdentifier != null ? this.playerIdentifier : this.formatRemoteAddress();
      if (this.loginTimingLastNs == 0L) {
         LOGGER.at(var2).log("[%s] %s", var5, var1);
      } else {
         long var6 = var3 - this.loginTimingLastNs;
         LOGGER.at(var2).log("[%s] %s took %d ms", var5, var1, TimeUnit.NANOSECONDS.toMillis(var6));
      }

      this.loginTimingLastNs = var3;
   }

   class="kw">private void runOnProcessingThread(Runnable var1) {
      if (Thread.currentThread() == this.connection.processingThread) {
         var1.run();
      } else {
         this.connection.execute(var1);
      }
   }

   long checkTimersAndGetDeadline(long var1) {
      if (this.closed.get()) {
         class="kw">return Long.MAX_VALUE;
      }

      if (this.stageDeadlineNs != 0L && var1 >= this.stageDeadlineNs) {
         String var3 = this.currentStage;
         BooleanSupplier var4 = this.stageCondition;
         Runnable var5 = this.stageOnTimeout;
         this.stageDeadlineNs = 0L;
         this.stageCondition = null;
         this.stageOnTimeout = null;
         if (!var4.getAsBoolean()) {
            long var6 = this.connectionStartNs > 0L ? TimeUnit.NANOSECONDS.toMillis(var1 - this.connectionStartNs) : -1L;
            String var8 = this.playerIdentifier != null ? this.playerIdentifier : this.formatRemoteAddress();
            LOGGER.at(Level.WARNING).log("Stage timeout for %s at stage '%s' after %d ms connected", var8, var3, var6);
            var5.run();
            class="kw">return Long.MAX_VALUE;
         }
      }

      if (this.packetTimeoutNs != 0L && var1 - this.lastPacketReceivedNs >= this.packetTimeoutNs) {
         long var10 = this.packetTimeoutNs;
         this.packetTimeoutNs = 0L;
         this.firePacketTimeout(var10, var1);
         class="kw">return Long.MAX_VALUE;
      }

      long var9 = Long.MAX_VALUE;
      if (this.packetTimeoutNs != 0L) {
         var9 = this.lastPacketReceivedNs + this.packetTimeoutNs;
      }

      if (this.stageDeadlineNs != 0L && this.stageDeadlineNs < var9) {
         var9 = this.stageDeadlineNs;
      }

      class="kw">return var9;
   }

   class="kw">private void firePacketTimeout(long var1, long var3) {
      String var5 = this.playerIdentifier != null ? this.playerIdentifier : this.formatRemoteAddress();
      String var6 = this.currentStage != null ? this.currentStage : "unknown";
      long var7 = this.connectionStartNs > 0L ? TimeUnit.NANOSECONDS.toMillis(var3 - this.connectionStartNs) : -1L;
      LOGGER.at(Level.INFO)
         .log("Read timeout for %s at stage '%s' after %d ms connected (timeout=%d ms)", var5, var6, var7, TimeUnit.NANOSECONDS.toMillis(var1));
      if (!QuicheNative.connIsClosed(this.connection.connHandle)) {
         FormattedMessage var9 = new FormattedMessage();
         var9.messageId = "client.general.disconnect.timeout.read";
         this.writeForClose(new ServerDisconnect(var9, DisconnectType.Disconnect));
         closeWithReason(this.connection.connHandle, QuicApplicationErrorCode.Timeout, var9);
      }
   }

   @Nonnull
   @Override
   class="kw">public CompletableFuture<Void> setupAuxiliaryChannels(@Nonnull ConnectionHandler var1, @Nonnull BiConsumer<NetworkChannel, ChannelConnection> var2) {
      CompletableFuture var3 = new CompletableFuture<>();
      this.execute(() -> {
         QuicheNative.connStreamPriority(this.connection.connHandle, this.streamId, (byte)0, true);
         QuicheChannel var4 = this.connection.createChannel(NetworkChannel.Chunks);
         var4.handler = var1;
         QuicheNative.connStreamPriority(this.connection.connHandle, var4.getStreamId(), (byte)0, true);
         var2.accept(NetworkChannel.Chunks, var4);
         QuicheChannel var5 = this.connection.createChannel(NetworkChannel.WorldMap);
         var5.handler = var1;
         QuicheNative.connStreamPriority(this.connection.connHandle, var5.getStreamId(), (byte)1, true);
         var2.accept(NetworkChannel.WorldMap, var5);
         var3.complete(null);
      });
      class="kw">return var3;
   }

   @Override
   class="kw">public void setChannelHandler(@Nonnull ConnectionHandler var1) {
      Runnable var2 = () -> {
         ConnectionHandler var2 = this.handler;
         if (var2 != null) {
            var2.unregistered(var1);
         }

         this.handler = var1;
         var1.registered(var2);
      };
      if (Thread.currentThread() == this.connection.processingThread) {
         var2.run();
      } else {
         this.execute(var2);
      }
   }

   @Override
   class="kw">public void closeConnection() {
      this.execute(() -> {
         QuicheNative.connStreamShutdown(this.connection.connHandle, this.streamId, 0, 0L);
         QuicheNative.connStreamShutdown(this.connection.connHandle, this.streamId, 1, 0L);
      });
   }

   @Override
   class="kw">public void closeApplicationConnection() {
      this.closeApplicationConnection(QuicApplicationErrorCode.NoError);
   }

   @Override
   class="kw">public void closeApplicationConnection(@Nonnull QuicApplicationErrorCode var1) {
      this.connection.execute(() -> {
         if (!QuicheNative.connIsClosed(this.connection.connHandle)) {
            QuicheNative.connClose(this.connection.connHandle, true, var1.getValue(), MemorySegment.NULL, 0L);
         }
      });
   }

   @Override
   class="kw">public void closeApplicationConnection(@Nonnull QuicApplicationErrorCode var1, @Nonnull FormattedMessage var2) {
      this.connection.execute(() -> {
         if (!QuicheNative.connIsClosed(this.connection.connHandle)) {
            closeWithReason(this.connection.connHandle, var1, var2);
         }
      });
   }
}