AuxiliaryStreamExceptionHandler class

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

Файл: com/hypixel/hytale/server/core/io/netty/HytaleChannelInitializer.java

extends: ChannelInboundHandlerAdapter

Поля (2)

МодификаторыТипИмя
private static final HytaleLogger LOGGER
private final String identifier

Методы (1)

МодификаторыВозвратСигнатура
public void exceptionCaughtpublic void exceptionCaught(@Nonnull ChannelHandlerContext var1, Throwable var2)

Исходный код

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

class="kw">import com.hypixel.hytale.common.util.FormatUtil;
class="kw">import com.hypixel.hytale.logger.HytaleLogger;
class="kw">import com.hypixel.hytale.protocol.FormattedMessage;
class="kw">import com.hypixel.hytale.protocol.io.PacketStatsRecorder;
class="kw">import com.hypixel.hytale.protocol.io.netty.PacketDecoder;
class="kw">import com.hypixel.hytale.protocol.io.netty.PacketEncoder;
class="kw">import com.hypixel.hytale.protocol.io.netty.ProtocolUtil;
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.config.RateLimitConfig;
class="kw">import com.hypixel.hytale.server.core.io.PacketHandler;
class="kw">import com.hypixel.hytale.server.core.io.PacketStatsRecorderImpl;
class="kw">import com.hypixel.hytale.server.core.io.handlers.InitialPacketHandler;
class="kw">import com.hypixel.hytale.server.core.io.stream.PendingStreamHandler;
class="kw">import com.hypixel.hytale.server.core.io.stream.StreamManager;
class="kw">import com.hypixel.hytale.server.core.io.transport.QUICTransport;
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.ChannelInitializer;
class="kw">import io.netty.handler.codec.quic.QuicChannel;
class="kw">import io.netty.handler.codec.quic.QuicStreamChannel;
class="kw">import io.netty.handler.timeout.ReadTimeoutException;
class="kw">import io.netty.handler.timeout.ReadTimeoutHandler;
class="kw">import io.netty.handler.timeout.TimeoutException;
class="kw">import io.netty.handler.timeout.WriteTimeoutException;
class="kw">import io.netty.util.AttributeKey;
class="kw">import java.io.IOException;
class="kw">import java.nio.channels.ClosedChannelException;
class="kw">import java.security.cert.X509Certificate;
class="kw">import java.time.Duration;
class="kw">import java.util.concurrent.TimeUnit;
class="kw">import java.util.concurrent.atomic.AtomicBoolean;
class="kw">import java.util.logging.Level;
class="kw">import javax.annotation.Nonnull;

class="kw">public class HytaleChannelInitializer class="kw">extends ChannelInitializer<Channel> {
   class="kw">public class="kw">static class="kw">final AttributeKey<PacketHandler> GAME_PACKET_HANDLER_ATTR = AttributeKey.valueOf("GAME_PACKET_HANDLER");

   class="kw">public HytaleChannelInitializer() {
   }

   class="kw">protected void initChannel(Channel var1) {
      if (var1 class="kw">instanceof QuicStreamChannel var2) {
         QuicChannel var3 = var2.parent();
         HytaleLogger.getLogger()
            .at(Level.INFO)
            .log("Received stream %d from %s to %s", var2.streamId(), NettyUtil.formatRemoteAddress(var1), NettyUtil.formatLocalAddress(var1));
         QuicApplicationErrorCode var4 = (QuicApplicationErrorCode)var3.attr(QUICTransport.ALPN_REJECT_ERROR_CODE_ATTR).get();
         if (var4 != null) {
            HytaleLogger.getLogger().at(Level.INFO).log("Rejecting stream from %s: client outdated (ALPN mismatch)", NettyUtil.formatRemoteAddress(var1));
            var1.config().setAutoRead(false);
            var1.pipeline().addLast("packetEncoder", new PacketEncoder());
            FormattedMessage var11 = Message.translation("client.general.disconnect.clientOutdated").getFormattedMessage();
            var1.writeAndFlush(new ServerDisconnect(var11, DisconnectType.Disconnect))
               .addListener(var3x -> var1.eventLoop().schedule(() -> NettyUtil.closeApplicationConnection(var1, var4, var11), 100L, TimeUnit.MILLISECONDS));
            class="kw">return;
         }

         X509Certificate var5 = (X509Certificate)var3.attr(QUICTransport.CLIENT_CERTIFICATE_ATTR).get();
         if (var5 != null) {
            var1.attr(QUICTransport.CLIENT_CERTIFICATE_ATTR).set(var5);
            HytaleLogger.getLogger().at(Level.FINE).log("Copied client certificate to stream: %s", var5.getSubjectX500Principal().getName());
         }

         PacketHandler var6 = (PacketHandler)var3.attr(GAME_PACKET_HANDLER_ATTR).get();
         if (var6 != null) {
            HytaleLogger.getLogger().at(Level.INFO).log("Setting up auxiliary stream %d for %s", var2.streamId(), var6.getIdentifier());
            this.initAuxiliaryStream(var1, var6);
            class="kw">return;
         }
      } else {
         HytaleLogger.getLogger()
            .at(Level.INFO)
            .log("Received connection from %s to %s", NettyUtil.formatRemoteAddress(var1), NettyUtil.formatLocalAddress(var1));
      }

      boolean var7 = true;
      boolean var8 = true;
      if (var1 class="kw">instanceof QuicStreamChannel var9) {
         var7 = !var9.isInputShutdown();
         var8 = !var9.isOutputShutdown();
      }

      PacketStatsRecorderImpl var10 = new PacketStatsRecorderImpl();
      var1.attr(PacketStatsRecorder.CHANNEL_KEY).set(var10);
      if (var7) {
         Duration var12 = HytaleServer.get().getConfig().getConnectionTimeouts().getInitial();
         var1.attr(ProtocolUtil.PACKET_TIMEOUT_KEY).set(var12);
         var1.pipeline().addLast("packetDecoder", new PacketDecoder());
         RateLimitConfig var14 = HytaleServer.get().getConfig().getRateLimitConfig();
         if (var14.isEnabled()) {
            var1.pipeline().addLast("rateLimit", new RateLimitHandler(var14.getBurstCapacity(), var14.getPacketsPerSecond()));
         }
      }

      if (var8) {
         var1.pipeline().addLast("packetEncoder", new PacketEncoder());
         var1.pipeline().addLast("packetArrayEncoder", NettyUtil.PACKET_ARRAY_ENCODER_INSTANCE);
      }

      if (NettyUtil.PACKET_LOGGER.getLevel() != Level.OFF) {
         var1.pipeline().addLast("logger", NettyUtil.LOGGER);
      }

      InitialPacketHandler var13 = new InitialPacketHandler(new NettyUtil.NettyChannelConnection(var1));
      var1.pipeline().addLast("handler", new PlayerChannelHandler(var13));
      var1.pipeline().addLast(new ChannelHandler[]{new HytaleChannelInitializer.ExceptionHandler()});
      if (var1 class="kw">instanceof QuicStreamChannel var15) {
         var15.parent().attr(GAME_PACKET_HANDLER_ATTR).set(var13);
         var15.updatePriority(StreamManager.GAME_STREAM_PRIORITY);
      }

      var13.registered(null);
   }

   class="kw">private void initAuxiliaryStream(Channel var1, PacketHandler var2) {
      var1.pipeline().addLast("aux_read_timeout", new ReadTimeoutHandler(5L, TimeUnit.SECONDS));
      var1.pipeline().addLast("packetDecoder", new PacketDecoder());
      var1.pipeline().addLast("packetEncoder", new PacketEncoder());
      var1.pipeline().addLast("pending_stream", new PendingStreamHandler(var2));
      var1.pipeline().addLast(new ChannelHandler[]{new HytaleChannelInitializer.AuxiliaryStreamExceptionHandler(var2.getIdentifier())});
   }

   class="kw">public void exceptionCaught(@Nonnull ChannelHandlerContext var1, Throwable var2) {
      ((HytaleLogger.Api)HytaleLogger.getLogger().at(Level.WARNING).withCause(var2)).log("Got exception from netty pipeline in HytaleChannelInitializer!");
      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);
      }
   }

   class="kw">public void channelInactive(@Nonnull ChannelHandlerContext var1) class="kw">throws Exception {
      NettyUtil.closeApplicationConnection(var1.channel());
      super.channelInactive(var1);
   }

   class="kw">private class="kw">static class AuxiliaryStreamExceptionHandler class="kw">extends ChannelInboundHandlerAdapter {
      class="kw">private class="kw">static class="kw">final HytaleLogger LOGGER = HytaleLogger.forEnclosingClass();
      class="kw">private class="kw">final String identifier;

      AuxiliaryStreamExceptionHandler(String var1) {
         this.identifier = var1;
      }

      class="kw">public void exceptionCaught(@Nonnull ChannelHandlerContext var1, Throwable var2) {
         if (!(var2 class="kw">instanceof ClosedChannelException)) {
            ((HytaleLogger.Api)LOGGER.at(Level.WARNING).withCause(var2)).log("Exception in auxiliary stream for %s", this.identifier);
            var1.close();
         }
      }
   }

   class="kw">private class="kw">static class ExceptionHandler class="kw">extends ChannelInboundHandlerAdapter {
      @Nonnull
      class="kw">private class="kw">static class="kw">final HytaleLogger LOGGER = HytaleLogger.forEnclosingClass();
      @Nonnull
      class="kw">private class="kw">static class="kw">final Message MESSAGE_DISCONNECT_TIMEOUT_READ = Message.translation("client.general.disconnect.timeout.read");
      @Nonnull
      class="kw">private class="kw">static class="kw">final Message MESSAGE_DISCONNECT_TIMEOUT_WRITE = Message.translation("client.general.disconnect.timeout.write");
      @Nonnull
      class="kw">private class="kw">static class="kw">final Message MESSAGE_DISCONNECT_TIMEOUT_CONNECTION = Message.translation("client.general.disconnect.timeout.connection");
      @Nonnull
      class="kw">private class="kw">final AtomicBoolean handled = new AtomicBoolean();

      class="kw">private ExceptionHandler() {
      }

      class="kw">public void exceptionCaught(@Nonnull ChannelHandlerContext var1, Throwable var2) {
         if (!(var2 class="kw">instanceof ClosedChannelException)) {
            ChannelHandler var4 = var1.pipeline().get("handler");
            String var3;
            if (var4 class="kw">instanceof PlayerChannelHandler) {
               var3 = ((PlayerChannelHandler)var4).getHandler().getIdentifier();
            } else {
               var3 = NettyUtil.formatRemoteAddress(var1.channel());
            }

            if (this.handled.getAndSet(true)) {
               if (var2 class="kw">instanceof IOException && var2.getMessage() != null) {
                  class="kw">switch (var2.getMessage()) {
                     case "Broken pipe":
                     case "Connection reset by peer":
                     case "An existing connection was forcibly closed by the remote host":
                        class="kw">return;
                  }
               }

               ((HytaleLogger.Api)LOGGER.at(Level.WARNING).withCause(var2)).log("Already handled exception in ExceptionHandler but got another!");
            } else if (var2 class="kw">instanceof TimeoutException) {
               this.handleTimeout(var1, var2, var3);
            } else {
               ((HytaleLogger.Api)LOGGER.at(Level.SEVERE).withCause(var2)).log("Got exception from netty pipeline in ExceptionHandler: %s", var2.getMessage());
               this.gracefulDisconnect(var1, var3, Message.translation("client.general.disconnect.internalServerError").getFormattedMessage());
            }
         }
      }

      class="kw">private void handleTimeout(@Nonnull ChannelHandlerContext var1, Throwable var2, String var3) {
         boolean var4 = var2 class="kw">instanceof ReadTimeoutException;
         boolean var5 = var2 class="kw">instanceof WriteTimeoutException;
         String var6 = var4 ? "Read" : (var5 ? "Write" : "Connection");
         Message var7 = var4 ? MESSAGE_DISCONNECT_TIMEOUT_READ : (var5 ? MESSAGE_DISCONNECT_TIMEOUT_WRITE : MESSAGE_DISCONNECT_TIMEOUT_CONNECTION);
         NettyUtil.TimeoutContext var8 = (NettyUtil.TimeoutContext)var1.channel().attr(NettyUtil.TimeoutContext.KEY).get();
         String var9 = var8 != null ? var8.stage() : "unknown";
         String var10 = var8 != null ? FormatUtil.nanosToString(System.nanoTime() - var8.connectionStartNs()) : "unknown";
         LOGGER.at(Level.INFO).log("%s timeout for %s at stage '%s' after %s connected", var6, var3, var9, var10);
         ((HytaleLogger.Api)NettyUtil.CONNECTION_EXCEPTION_LOGGER.at(Level.FINE).withCause(var2))
            .log("%s timeout for %s at stage '%s' after %s connected", var6, var3, var9, var10);
         this.gracefulDisconnect(var1, var3, var7.getFormattedMessage());
      }

      class="kw">private void gracefulDisconnect(@Nonnull ChannelHandlerContext var1, String var2, FormattedMessage var3) {
         Channel var4 = var1.channel();
         if (var4.isWritable()) {
            var4.writeAndFlush(new ServerDisconnect(var3, DisconnectType.Disconnect))
               .addListener(var1x -> NettyUtil.closeApplicationConnection(var4, QuicApplicationErrorCode.Timeout));
            var4.eventLoop().schedule(() -> {
               if (var4.isOpen()) {
                  LOGGER.at(Level.FINE).log("Force closing %s after graceful disconnect attempt", var2);
                  NettyUtil.closeApplicationConnection(var4, QuicApplicationErrorCode.Timeout);
               }
            }, 1L, TimeUnit.SECONDS);
         } else {
            NettyUtil.closeApplicationConnection(var4, QuicApplicationErrorCode.Timeout);
         }
      }
   }
}