PendingStreamHandler class

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

Файл: com/hypixel/hytale/server/core/io/stream/PendingStreamHandler.java

extends: ChannelInboundHandlerAdapter

Поля (5)

МодификаторыТипИмя
final HytaleLogger LOGGER
NettyUtil.NettyChannelConnection var12
ConnectionHandler var13
StreamType var5
ChannelConnection var6

Методы (7)

МодификаторыВозвратСигнатура
public void channelInactivevoid channelInactive(@Nonnull ChannelHandlerContext var1)
public void channelReadvoid channelRead(@Nonnull ChannelHandlerContext var1, @Nonnull Object var2)
public void exceptionCaughtvoid exceptionCaught(@Nonnull ChannelHandlerContext var1, @Nonnull Throwable var2)
if if(var2 instanceof Packet var3)
if if(var3 instanceof StreamOpen var4)
if if(var13 == null)
if if(var2 instanceof ProtocolException)

Исходный код

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

class="kw">import com.hypixel.hytale.logger.HytaleLogger;
class="kw">import com.hypixel.hytale.protocol.Packet;
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.ProtocolException;
class="kw">import com.hypixel.hytale.protocol.packets.stream.StreamOpen;
class="kw">import com.hypixel.hytale.protocol.packets.stream.StreamOpenResponse;
class="kw">import com.hypixel.hytale.protocol.packets.stream.StreamType;
class="kw">import com.hypixel.hytale.server.core.io.PacketHandler;
class="kw">import com.hypixel.hytale.server.core.io.netty.NettyUtil;
class="kw">import io.netty.channel.Channel;
class="kw">import io.netty.channel.ChannelHandlerContext;
class="kw">import io.netty.channel.ChannelInboundHandlerAdapter;
class="kw">import io.netty.handler.codec.quic.QuicStreamChannel;
class="kw">import io.netty.util.ReferenceCountUtil;
class="kw">import java.util.logging.Level;
class="kw">import javax.annotation.Nonnull;

class="kw">public class PendingStreamHandler class="kw">extends ChannelInboundHandlerAdapter {
   class="kw">private class="kw">static class="kw">final HytaleLogger LOGGER = HytaleLogger.forEnclosingClass();
   class="kw">private class="kw">static class="kw">final int MAX_AUXILIARY_STREAMS = 4;
   class="kw">private class="kw">final PacketHandler packetHandler;
   class="kw">private class="kw">final StreamManager streamManager;

   class="kw">public PendingStreamHandler(@Nonnull PacketHandler var1) {
      this(var1, StreamManager.getInstance());
   }

   class="kw">public PendingStreamHandler(@Nonnull PacketHandler var1, @Nonnull StreamManager var2) {
      this.packetHandler = var1;
      this.streamManager = var2;
   }

   class="kw">public void channelRead(@Nonnull ChannelHandlerContext var1, @Nonnull Object var2) {
      if (var2 class="kw">instanceof Packet var3) {
         if (var3 class="kw">instanceof StreamOpen var4) {
            StreamType var5 = var4.type;
            if (this.packetHandler.checkStreamOpenRateLimit()) {
               LOGGER.at(Level.WARNING).log("Stream open rate limited for %s requesting %s", this.packetHandler.getIdentifier(), var5.name());
               var1.writeAndFlush(new StreamOpenResponse(var5, false, "Rate limited - try again later")).addListener(var1x -> var1.close());
            } else if (var5 == StreamType.Game) {
               LOGGER.at(Level.WARNING).log("Cannot open Game stream - stream 0 is already the game stream, from %s", this.packetHandler.getIdentifier());
               var1.writeAndFlush(new StreamOpenResponse(var5, false, "Game stream cannot be opened explicitly")).addListener(var1x -> var1.close());
            } else if (!this.streamManager.isSupported(var5)) {
               LOGGER.at(Level.INFO).log("Unsupported stream type %s from %s", var5.name(), this.packetHandler.getIdentifier());
               var1.writeAndFlush(new StreamOpenResponse(var5, false, "Stream type not supported")).addListener(var1x -> var1.close());
            } else {
               ChannelConnection var6 = this.packetHandler.getChannel(var5);
               if (var6 class="kw">instanceof NettyUtil.NettyChannelConnection(Channel var8)) {
                  LOGGER.at(Level.INFO)
                     .log("Replacing stale %s stream for %s (old channel active=%s)", var5.name(), this.packetHandler.getIdentifier(), var8.isActive());
                  this.packetHandler.compareAndSetChannel(var5, var6, null);
                  var8.close();
               }

               if (this.packetHandler.getAuxiliaryChannelCount() >= 4) {
                  LOGGER.at(Level.WARNING).log("Maximum auxiliary streams exceeded for %s requesting %s", this.packetHandler.getIdentifier(), var5.name());
                  var1.writeAndFlush(new StreamOpenResponse(var5, false, "Maximum auxiliary streams exceeded")).addListener(var1x -> var1.close());
               } else {
                  NettyUtil.NettyChannelConnection var12 = new NettyUtil.NettyChannelConnection(var1.channel());
                  ConnectionHandler var13 = this.streamManager.createHandler(var5, this.packetHandler, var12);
                  if (var13 == null) {
                     LOGGER.at(Level.SEVERE).log("Failed to create handler for stream type %s from %s", var5.name(), this.packetHandler.getIdentifier());
                     var1.writeAndFlush(new StreamOpenResponse(var5, false, "Internal error")).addListener(var1x -> var1.close());
                  } else {
                     LOGGER.at(Level.INFO).log("Opening %s stream for %s", var5.name(), this.packetHandler.getIdentifier());
                     var1.pipeline().replace(this, var5.name() + "Handler", new StreamConnectionHandlerAdapter(var13));
                     var1.pipeline().remove("aux_read_timeout");
                     if (var1.channel() class="kw">instanceof QuicStreamChannel var14) {
                        var14.updatePriority(this.streamManager.getStreamPriority(var5));
                     }

                     var1.writeAndFlush(new StreamOpenResponse(var5, true, null));
                  }
               }
            }
         } else {
            LOGGER.at(Level.WARNING).log("Auxiliary stream first packet was not StreamOpen, closing: %s", var3.getClass().getSimpleName());
            var1.close();
         }
      } else {
         LOGGER.at(Level.WARNING)
            .log("Expected Packet but got %s on pending stream from %s", var2.getClass().getSimpleName(), NettyUtil.formatRemoteAddress(var1.channel()));
         ReferenceCountUtil.release(var2);
         var1.close();
      }
   }

   class="kw">public void exceptionCaught(@Nonnull ChannelHandlerContext var1, @Nonnull Throwable var2) {
      if (var2 class="kw">instanceof ProtocolException) {
         LOGGER.at(Level.WARNING).log("Protocol error on pending stream for %s: %s", this.packetHandler.getIdentifier(), var2.getMessage());
      } else {
         ((HytaleLogger.Api)LOGGER.at(Level.WARNING).withCause(var2)).log("Exception in pending stream handler for %s", this.packetHandler.getIdentifier());
      }

      var1.close();
   }

   class="kw">public void channelInactive(@Nonnull ChannelHandlerContext var1) class="kw">throws Exception {
      LOGGER.at(Level.FINE).log("Pending stream closed for %s", this.packetHandler.getIdentifier());
      super.channelInactive(var1);
   }
}