QuicChannelInboundHandlerAdapter class

Пакет: com.hypixel.hytale.server.core.io.transport

Файл: com/hypixel/hytale/server/core/io/transport/QUICTransport.java

extends: ChannelInboundHandlerAdapter

Поля (12)

МодификаторыТипИмя
private final QuicSslContext sslContext
Duration var2
QuicChannel var2
SSLEngine var2
ChannelHandler var3
String var3
Channel var3
Certificate[] var3
String var4
FormattedMessage var4
int var5
X509Certificate var6

Методы (15)

МодификаторыВозвратСигнатура
new ChannelInboundHandlerAdapternew ChannelInboundHandlerAdapter()
public void channelActivepublic void channelActive(@Nonnull ChannelHandlerContext var1)
public void channelActivepublic void channelActive(@Nonnull ChannelHandlerContext var1)
public void channelInactivepublic void channelInactive(@Nonnull ChannelHandlerContext var1)
public void exceptionCaughtpublic void exceptionCaught(@Nonnull ChannelHandlerContext var1, Throwable var2)
private X509Certificate extractClientCertificateprivate X509Certificate extractClientCertificate(QuicChannel var1)
if if(var2 instanceof SniCompletionEvent var3)
if if(var5 < 3)
if if(var6 == null)
if if(var2 == null)
if if(var3 != null && var3.length > 0 && var3[0] instanceof X509Certificate)
public boolean isSharablepublic boolean isSharable()
public boolean isSharablepublic boolean isSharable()
private int parseProtocolVersionprivate int parseProtocolVersion(String var1)
public void userEventTriggeredpublic void userEventTriggered(ChannelHandlerContext var1, Object var2)

Исходный код

Показать/скрыть
class="kw">package com.hypixel.hytale.server.core.io.transport;

class="kw">import com.hypixel.hytale.logger.HytaleLogger;
class="kw">import com.hypixel.hytale.protocol.FormattedMessage;
class="kw">import com.hypixel.hytale.protocol.io.ServerListener;
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 com.hypixel.hytale.server.core.HytaleServer;
class="kw">import com.hypixel.hytale.server.core.Message;
class="kw">import com.hypixel.hytale.server.core.auth.CertificateUtil;
class="kw">import com.hypixel.hytale.server.core.auth.ServerAuthManager;
class="kw">import com.hypixel.hytale.server.core.io.netty.HytaleChannelInitializer;
class="kw">import com.hypixel.hytale.server.core.io.netty.NettyUtil;
class="kw">import io.netty.bootstrap.Bootstrap;
class="kw">import io.netty.channel.Channel;
class="kw">import io.netty.channel.ChannelHandler;
class="kw">import io.netty.channel.ChannelHandlerContext;
class="kw">import io.netty.channel.ChannelInboundHandlerAdapter;
class="kw">import io.netty.channel.ChannelOption;
class="kw">import io.netty.channel.EventLoopGroup;
class="kw">import io.netty.channel.socket.DatagramChannel;
class="kw">import io.netty.channel.socket.SocketProtocolFamily;
class="kw">import io.netty.channel.socket.nio.NioChannelOption;
class="kw">import io.netty.handler.codec.quic.QLogConfiguration;
class="kw">import io.netty.handler.codec.quic.QuicChannel;
class="kw">import io.netty.handler.codec.quic.QuicChannelOption;
class="kw">import io.netty.handler.codec.quic.QuicCongestionControlAlgorithm;
class="kw">import io.netty.handler.codec.quic.QuicServerCodecBuilder;
class="kw">import io.netty.handler.codec.quic.QuicSslContext;
class="kw">import io.netty.handler.codec.quic.QuicSslContextBuilder;
class="kw">import io.netty.handler.ssl.ClientAuth;
class="kw">import io.netty.handler.ssl.SniCompletionEvent;
class="kw">import io.netty.handler.ssl.util.InsecureTrustManagerFactory;
class="kw">import io.netty.handler.ssl.util.SelfSignedCertificate;
class="kw">import io.netty.util.AttributeKey;
class="kw">import io.netty.util.Mapping;
class="kw">import java.lang.reflect.Field;
class="kw">import java.net.Inet4Address;
class="kw">import java.net.Inet6Address;
class="kw">import java.net.InetSocketAddress;
class="kw">import java.security.cert.Certificate;
class="kw">import java.security.cert.CertificateException;
class="kw">import java.security.cert.X509Certificate;
class="kw">import java.time.Duration;
class="kw">import java.util.Locale;
class="kw">import java.util.concurrent.CompletableFuture;
class="kw">import java.util.concurrent.TimeUnit;
class="kw">import java.util.logging.Level;
class="kw">import javax.annotation.Nonnull;
class="kw">import javax.annotation.Nullable;
class="kw">import javax.net.ssl.SSLEngine;
class="kw">import jdk.net.ExtendedSocketOptions;

class="kw">public class QUICTransport class="kw">implements Transport {
   class="kw">private class="kw">static class="kw">final HytaleLogger LOGGER = HytaleLogger.forEnclosingClass();
   class="kw">public class="kw">static class="kw">final AttributeKey<X509Certificate> CLIENT_CERTIFICATE_ATTR = AttributeKey.valueOf("CLIENT_CERTIFICATE");
   class="kw">public class="kw">static class="kw">final AttributeKey<QuicApplicationErrorCode> ALPN_REJECT_ERROR_CODE_ATTR = AttributeKey.valueOf("ALPN_REJECT_ERROR_CODE");
   class="kw">public class="kw">static class="kw">final AttributeKey<String> SNI_HOSTNAME_ATTR = AttributeKey.valueOf("SNI_HOSTNAME");
   class="kw">private class="kw">static class="kw">final int SOCKET_RECEIVE_BUFFER_SIZE = 33554432;
   class="kw">private class="kw">static class="kw">final boolean IS_LINUX = System.getProperty("os.name", "").toLowerCase(Locale.ROOT).contains("linux");
   class="kw">private class="kw">static class="kw">final int AUTOTUNE_CONNECTION_WINDOW = 25165824;
   @Nonnull
   class="kw">private class="kw">final EventLoopGroup workerGroup = NettyUtil.getEventLoopGroup("ServerWorkerGroup");
   class="kw">private class="kw">final Bootstrap bootstrapIpv4;
   class="kw">private class="kw">final Bootstrap bootstrapIpv6;

   class="kw">public QUICTransport() class="kw">throws InterruptedException {
      SelfSignedCertificate var1 = null;

      try {
         var1 = new SelfSignedCertificate("localhost");
      } catch (CertificateException var7) {
         throw new RuntimeException(var7);
      }

      ServerAuthManager.getInstance().setServerCertificate(var1.cert());
      LOGGER.at(Level.INFO).log("Server certificate registered for mutual auth, fingerprint: %s", CertificateUtil.computeCertificateFingerprint(var1.cert()));
      QuicSslContext var2 = QuicSslContextBuilder.forServer(var1.key(), null, new X509Certificate[]{var1.cert()})
         .applicationProtocols(new String[]{"hytale/3", "hytale/2"})
         .earlyData(false)
         .clientAuth(ClientAuth.REQUIRE)
         .trustManager(InsecureTrustManagerFactory.INSTANCE)
         .build();

      QuicSslContext var3;
      try {
         QuicSslContextBuilder var4 = QuicSslContextBuilder.forServer(var1.key(), null, new X509Certificate[]{var1.cert()})
            .earlyData(false)
            .clientAuth(ClientAuth.REQUIRE);
         Field var5 = var4.getClass().getDeclaredField("mapping");
         var5.setAccessible(true);
         var5.set(var4, (Mapping)var1x -> var2);
         var3 = var4.build();
      } catch (NoSuchFieldException | IllegalAccessException var6) {
         ((HytaleLogger.Api)LOGGER.at(Level.WARNING).withCause(var6)).log("Failed to set SNI mapping via reflection, SNI support disabled");
         var3 = var2;
      }

      NettyUtil.ReflectiveChannelFactory var9 = NettyUtil.getDatagramChannelFactory(SocketProtocolFamily.INET);
      LOGGER.at(Level.INFO).log("Using IPv4 Datagram Channel: %s...", var9.getSimpleName());
      this.bootstrapIpv4 = ((Bootstrap)((Bootstrap)((Bootstrap)((Bootstrap)((Bootstrap)((Bootstrap)new Bootstrap().group(this.workerGroup))
                        .channelFactory(var9))
                     .option(ChannelOption.SO_REUSEADDR, true))
                  .option(ChannelOption.SO_RCVBUF, 33554432))
               .option(NioChannelOption.of(ExtendedSocketOptions.IP_DONTFRAGMENT), true))
            .handler(new QUICTransport.QuicChannelInboundHandlerAdapter(var3)))
         .validate();
      NettyUtil.ReflectiveChannelFactory var10 = NettyUtil.getDatagramChannelFactory(SocketProtocolFamily.INET6);
      LOGGER.at(Level.INFO).log("Using IPv6 Datagram Channel: %s...", var10.getSimpleName());
      this.bootstrapIpv6 = ((Bootstrap)((Bootstrap)((Bootstrap)((Bootstrap)((Bootstrap)((Bootstrap)new Bootstrap().group(this.workerGroup))
                        .channelFactory(var10))
                     .option(ChannelOption.SO_REUSEADDR, true))
                  .option(ChannelOption.SO_RCVBUF, 33554432))
               .option(NioChannelOption.of(ExtendedSocketOptions.IP_DONTFRAGMENT), true))
            .handler(new QUICTransport.QuicChannelInboundHandlerAdapter(var3)))
         .validate();
      this.bootstrapIpv4.register().sync();
      this.bootstrapIpv6.register().sync();
   }

   @Nonnull
   @Override
   class="kw">public TransportType getType() {
      class="kw">return TransportType.QUIC;
   }

   @Override
   class="kw">public CompletableFuture<ServerListener> bind(@Nonnull InetSocketAddress var1) {
      if (var1.getAddress() class="kw">instanceof Inet4Address) {
         class="kw">return NettyUtil.wrapChannelFuture(this.bootstrapIpv4.bind(var1), var0 -> {
            warnIfReceiveBufferClamped(var0.channel());
            class="kw">return new NettyUtil.NettyChannelServerListener(var0.channel());
         });
      } else if (var1.getAddress() class="kw">instanceof Inet6Address) {
         class="kw">return NettyUtil.wrapChannelFuture(this.bootstrapIpv6.bind(var1), var0 -> {
            warnIfReceiveBufferClamped(var0.channel());
            class="kw">return new NettyUtil.NettyChannelServerListener(var0.channel());
         });
      } else {
         throw new UnsupportedOperationException("Unsupported address type: " + var1.getAddress().getClass());
      }
   }

   class="kw">private class="kw">static void warnIfReceiveBufferClamped(Channel var0) {
      Integer var1 = (Integer)var0.config().getOption(ChannelOption.SO_RCVBUF);
      if (var1 != null) {
         int var2 = IS_LINUX ? var1 / 2 : var1;
         if (var2 < 25165824) {
            LOGGER.at(Level.WARNING)
               .log(
                  "UDP receive buffer is %d bytes (effective %d), below the %d-byte autotune window. Raise net.core.rmem_max so inbound bursts are not dropped",
                  var1,
                  var2,
                  25165824
               );
         }
      }
   }

   @Override
   class="kw">public void shutdown() {
      LOGGER.at(Level.INFO).log("Shutting down workerGroup...");

      try {
         this.workerGroup.shutdownGracefully(0L, 1L, TimeUnit.SECONDS).await(1L, TimeUnit.SECONDS);
      } catch (InterruptedException var2) {
         ((HytaleLogger.Api)LOGGER.at(Level.SEVERE).withCause(var2)).log("Failed to await for listener to close!");
         Thread.currentThread().interrupt();
      }
   }

   class="kw">private class="kw">static class QuicChannelInboundHandlerAdapter class="kw">extends ChannelInboundHandlerAdapter {
      class="kw">private class="kw">final QuicSslContext sslContext;

      class="kw">public QuicChannelInboundHandlerAdapter(QuicSslContext var1) {
         this.sslContext = var1;
      }

      class="kw">public boolean isSharable() {
         class="kw">return true;
      }

      class="kw">public void channelActive(@Nonnull ChannelHandlerContext var1) {
         Duration var2 = HytaleServer.get().getConfig().getConnectionTimeouts().getPlay();
         ChannelHandler var3 = ((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)((QuicServerCodecBuilder)new QuicServerCodecBuilder()
                                                .sslContext(this.sslContext))
                                             .tokenHandler(null)
                                             .activeMigration(false))
                                          .maxIdleTimeout(var2.toMillis(), TimeUnit.MILLISECONDS))
                                       .ackDelayExponent(3L))
                                    .initialMaxData(524288L))
                                 .initialMaxStreamDataUnidirectional(0L))
                              .initialMaxStreamsUnidirectional(0L))
                           .initialMaxStreamDataBidirectionalLocal(131072L))
                        .initialMaxStreamDataBidirectionalRemote(131072L))
                     .initialMaxStreamsBidirectional(8L))
                  .discoverPmtu(true))
               .congestionControlAlgorithm(QuicCongestionControlAlgorithm.BBR))
            .option(QuicChannelOption.QLOG, System.getProperty("hytale.qlog") != null ? new QLogConfiguration(".", "hytale-server-quic-qlogs", "") : null)
            .handler(
               new ChannelInboundHandlerAdapter() {
                  class="kw">public boolean isSharable() {
                     class="kw">return true;
                  }

                  class="kw">public void userEventTriggered(ChannelHandlerContext var1, Object var2) class="kw">throws Exception {
                     if (var2 class="kw">instanceof SniCompletionEvent var3) {
                        var1.channel().attr(QUICTransport.SNI_HOSTNAME_ATTR).set(var3.hostname());
                     }

                     super.userEventTriggered(var1, var2);
                  }

                  class="kw">public void channelActive(@Nonnull ChannelHandlerContext var1) {
                     QuicChannel var2 = (QuicChannel)var1.channel();
                     String var3 = (String)var2.attr(QUICTransport.SNI_HOSTNAME_ATTR).get();
                     QUICTransport.LOGGER
                        .at(Level.INFO)
                        .log("Received connection from %s to %s (SNI: %s)", NettyUtil.formatRemoteAddress(var2), NettyUtil.formatLocalAddress(var2), var3);
                     String var4 = var2.sslEngine().getApplicationProtocol();
                     int var5 = this.parseProtocolVersion(var4);
                     if (var5 < 3) {
                        QUICTransport.LOGGER
                           .at(Level.INFO)
                           .log("Marking connection from %s (SNI: %s) for rejection: ALPN %s < required %d", NettyUtil.formatRemoteAddress(var2), var3, var4, 3);
                        var2.attr(QUICTransport.ALPN_REJECT_ERROR_CODE_ATTR).set(QuicApplicationErrorCode.ClientOutdated);
                     }

                     X509Certificate var6 = QuicChannelInboundHandlerAdapter.this.extractClientCertificate(var2);
                     if (var6 == null) {
                        QUICTransport.LOGGER
                           .at(Level.WARNING)
                           .log("Connection rejected: no client certificate from %s (SNI: %s)", NettyUtil.formatRemoteAddress(var2), var3);
                        NettyUtil.closeConnection(var2);
                     } else {
                        var2.attr(QUICTransport.CLIENT_CERTIFICATE_ATTR).set(var6);
                        QUICTransport.LOGGER.at(Level.FINE).log("Client certificate: %s", var6.getSubjectX500Principal().getName());
                     }
                  }

                  class="kw">private int parseProtocolVersion(String var1) {
                     if (var1 != null && var1.startsWith("hytale/")) {
                        try {
                           class="kw">return Integer.parseInt(var1.substring(7));
                        } catch (NumberFormatException var3) {
                           class="kw">return 0;
                        }
                     } else {
                        class="kw">return 0;
                     }
                  }

                  class="kw">public void channelInactive(@Nonnull ChannelHandlerContext var1) {
                     ((QuicChannel)var1.channel()).collectStats().addListener(var0 -> {
                        if (var0.isSuccess()) {
                           QUICTransport.LOGGER.at(Level.INFO).log("Connection closed: %s", var0.getNow());
                        }
                     });
                  }

                  class="kw">public void exceptionCaught(@Nonnull ChannelHandlerContext var1, Throwable var2) {
                     ((HytaleLogger.Api)QUICTransport.LOGGER.at(Level.WARNING).withCause(var2)).log("Got exception from netty pipeline in ChannelInitializer!");
                     Channel var3 = var1.channel();
                     if (var3.isWritable()) {
                        FormattedMessage var4 = Message.translation("client.general.disconnect.internalServerError").getFormattedMessage();
                        var3.writeAndFlush(new ServerDisconnect(var4, DisconnectType.Crash))
                           .addListener(var2x -> NettyUtil.closeApplicationConnection(var3, QuicApplicationErrorCode.Crash, var4));
                     } else {
                        NettyUtil.closeApplicationConnection(var3);
                     }
                  }
               }
            )
            .streamHandler(new HytaleChannelInitializer())
            .build();
         var1.channel().pipeline().addLast(new ChannelHandler[]{var3});
      }

      @Nullable
      class="kw">private X509Certificate extractClientCertificate(QuicChannel var1) {
         try {
            SSLEngine var2 = var1.sslEngine();
            if (var2 == null) {
               class="kw">return null;
            }

            Certificate[] var3 = var2.getSession().getPeerCertificates();
            if (var3 != null && var3.length > 0 && var3[0] class="kw">instanceof X509Certificate) {
               class="kw">return (X509Certificate)var3[0];
            }
         } catch (Exception var4) {
            QUICTransport.LOGGER.at(Level.FINEST).log("No peer certificate available: %s", var4.getMessage());
         }

         class="kw">return null;
      }
   }
}