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);
}
}