PingInfo class
Пакет: com.hypixel.hytale.server.core.io
Файл: com/hypixel/hytale/server/core/io/PacketHandler.java
Поля (24)
| Модификаторы | Тип | Имя |
|---|---|---|
|
public static final int |
FIVE_MINUTE_INDEX |
|
public static final MetricsRegistry<PacketHandler.PingInfo> |
METRICS_REGISTRY |
|
public static final int |
ONE_MINUTE_INDEX |
|
public static final int |
ONE_SECOND_INDEX |
|
public static final double |
PERCENTILE |
|
public static final int |
PING_FREQUENCY |
|
public static final int |
PING_FREQUENCY_MILLIS |
|
public static final TimeUnit |
PING_FREQUENCY_UNIT |
|
public static final int |
PING_HISTORY_LENGTH |
|
public static final int |
PING_HISTORY_MILLIS |
|
public static final TimeUnit |
TIME_UNIT |
|
protected final Metric |
packetQueueMetric |
|
protected final IntPriorityQueue |
pingIdQueue |
|
protected final Lock |
pingLock |
|
protected final HistoricMetric |
pingMetricSet |
|
protected final LongPriorityQueue |
pingTimestampQueue |
|
protected final PongType |
pingType |
|
protected final Lock |
queueLock |
|
int |
var2 |
|
|
var2 |
|
long |
var3 |
|
|
var3 |
|
long |
var5 |
|
long |
var7 |
Методы (10)
| Модификаторы | Возврат | Сигнатура |
|---|---|---|
abstract |
throw new |
IllegalArgumentExceptionthrow new IllegalArgumentException("Got packet for " + var1.type + " but expected " + this.pingType) |
|
public void |
clearpublic void clear() |
|
public Metric |
getPacketQueueMetricpublic Metric getPacketQueueMetric() |
|
public HistoricMetric |
getPingMetricSetpublic HistoricMetric getPingMetricSet() |
|
public PongType |
getPingTypepublic PongType getPingType() |
|
protected void |
handlePacketprotected void handlePacket(@Nonnull Pong var1) |
|
|
if if(var1.type != this.pingType) |
|
|
if if(var1.id != var2) |
|
|
if if(var7 <= 0L) |
|
protected void |
recordSentprotected void recordSent(int var1, long var2) |
Исходный код
Показать/скрыть
class="kw">package com.hypixel.hytale.server.core.io;
class="kw">import com.hypixel.hytale.codec.Codec;
class="kw">import com.hypixel.hytale.codec.codecs.EnumCodec;
class="kw">import com.hypixel.hytale.common.util.DebugFlags;
class="kw">import com.hypixel.hytale.common.util.NetworkUtil;
class="kw">import com.hypixel.hytale.logger.HytaleLogger;
class="kw">import com.hypixel.hytale.metrics.MetricsRegistry;
class="kw">import com.hypixel.hytale.metrics.metric.HistoricMetric;
class="kw">import com.hypixel.hytale.metrics.metric.Metric;
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.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.PacketStatsRecorder;
class="kw">import com.hypixel.hytale.protocol.packets.connection.DisconnectType;
class="kw">import com.hypixel.hytale.protocol.packets.connection.Ping;
class="kw">import com.hypixel.hytale.protocol.packets.connection.Pong;
class="kw">import com.hypixel.hytale.protocol.packets.connection.PongType;
class="kw">import com.hypixel.hytale.protocol.packets.stream.StreamType;
class="kw">import com.hypixel.hytale.server.core.Message;
class="kw">import com.hypixel.hytale.server.core.auth.PlayerAuthentication;
class="kw">import com.hypixel.hytale.server.core.io.adapter.PacketAdapters;
class="kw">import com.hypixel.hytale.server.core.io.handlers.login.AuthenticationPacketHandler;
class="kw">import com.hypixel.hytale.server.core.io.handlers.login.PasswordPacketHandler;
class="kw">import com.hypixel.hytale.server.core.modules.time.WorldTimeResource;
class="kw">import com.hypixel.hytale.server.core.receiver.IPacketReceiver;
class="kw">import com.hypixel.hytale.server.core.util.MessageUtil;
class="kw">import io.netty.channel.local.LocalAddress;
class="kw">import io.netty.channel.unix.DomainSocketAddress;
class="kw">import io.netty.handler.codec.quic.QuicStreamPriority;
class="kw">import it.unimi.dsi.fastutil.ints.IntArrayFIFOQueue;
class="kw">import it.unimi.dsi.fastutil.ints.IntPriorityQueue;
class="kw">import it.unimi.dsi.fastutil.longs.LongArrayFIFOQueue;
class="kw">import it.unimi.dsi.fastutil.longs.LongPriorityQueue;
class="kw">import java.net.InetAddress;
class="kw">import java.net.InetSocketAddress;
class="kw">import java.net.SocketAddress;
class="kw">import java.security.SecureRandom;
class="kw">import java.time.Duration;
class="kw">import java.time.Instant;
class="kw">import java.util.Collections;
class="kw">import java.util.EnumMap;
class="kw">import java.util.Map;
class="kw">import java.util.Objects;
class="kw">import java.util.concurrent.CompletableFuture;
class="kw">import java.util.concurrent.TimeUnit;
class="kw">import java.util.concurrent.atomic.AtomicInteger;
class="kw">import java.util.concurrent.atomic.AtomicLong;
class="kw">import java.util.concurrent.locks.Lock;
class="kw">import java.util.concurrent.locks.ReentrantLock;
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="kw">abstract class PacketHandler class="kw">implements IPacketReceiver, ConnectionHandler {
class="kw">public class="kw">static class="kw">final int MAX_PACKET_ID = 512;
@Nonnull
class="kw">public class="kw">static class="kw">final Map<NetworkChannel, QuicStreamPriority> DEFAULT_STREAM_PRIORITIES = Map.of(
NetworkChannel.Default,
new QuicStreamPriority(0, true),
NetworkChannel.Chunks,
new QuicStreamPriority(0, true),
NetworkChannel.WorldMap,
new QuicStreamPriority(1, true)
);
class="kw">private class="kw">static class="kw">final HytaleLogger LOGIN_TIMING_LOGGER = HytaleLogger.get("LoginTiming");
@Nonnull
class="kw">protected class="kw">final ChannelConnection[] channels = new ChannelConnection[NetworkChannel.COUNT];
@Nonnull
class="kw">protected class="kw">final ProtocolVersion protocolVersion;
@Nullable
class="kw">protected PlayerAuthentication auth;
class="kw">protected boolean queuePackets;
class="kw">protected class="kw">final AtomicInteger queuedPackets = new AtomicInteger();
class="kw">protected class="kw">final SecureRandom pingIdRandom = new SecureRandom();
@Nonnull
class="kw">protected class="kw">final PacketHandler.PingInfo[] pingInfo;
class="kw">private float pingTimer;
class="kw">protected boolean registered;
@Nullable
class="kw">protected Throwable clientReadyForChunksFutureStack;
@Nullable
class="kw">protected CompletableFuture<Void> clientReadyForChunksFuture;
@Nonnull
class="kw">protected class="kw">final PacketHandler.DisconnectReason disconnectReason = new PacketHandler.DisconnectReason();
@Nonnull
class="kw">private class="kw">final Map<StreamType, ChannelConnection> auxiliaryChannels = Collections.synchronizedMap(new EnumMap<>(StreamType.class));
class="kw">private class="kw">final AtomicLong lastStreamOpenTimeNanos = new AtomicLong();
class="kw">private class="kw">static class="kw">final long STREAM_OPEN_MIN_INTERVAL_NANOS = TimeUnit.SECONDS.toNanos(1L);
class="kw">public PacketHandler(@Nonnull ChannelConnection var1, @Nonnull ProtocolVersion var2) {
this.channels[0] = var1;
this.protocolVersion = var2;
this.pingInfo = new PacketHandler.PingInfo[PongType.VALUES.length];
for (PongType var6 : PongType.VALUES) {
this.pingInfo[var6.ordinal()] = new PacketHandler.PingInfo(var6);
}
}
@Nonnull
class="kw">public ChannelConnection getChannel() {
class="kw">return this.channels[0];
}
@Nullable
class="kw">public ChannelConnection getChannel(@Nonnull StreamType var1) {
class="kw">return var1 == StreamType.Game ? this.channels[0] : this.auxiliaryChannels.get(var1);
}
@Nonnull
class="kw">public class="kw">abstract String getIdentifier();
@Nonnull
class="kw">public ProtocolVersion getProtocolVersion() {
class="kw">return this.protocolVersion;
}
@Override
class="kw">public class="kw">final void registered(@Nullable ConnectionHandler var1) {
this.registered = true;
this.registered0(var1);
}
class="kw">protected void registered0(@Nullable ConnectionHandler var1) {
}
@Override
class="kw">public class="kw">final void unregistered(@Nullable ConnectionHandler var1) {
this.registered = false;
this.clearTimeout();
this.unregistered0(var1);
}
class="kw">protected void unregistered0(@Nullable ConnectionHandler var1) {
}
@Override
class="kw">public void handle(@Nonnull ToServerPacket var1) {
this.accept(var1);
}
class="kw">public class="kw">abstract void accept(@Nonnull ToServerPacket var1);
@Override
class="kw">public void logCloseMessage() {
HytaleLogger.getLogger().at(Level.INFO).log("%s was closed.", this.getIdentifier());
}
@Override
class="kw">public void closed(@Nullable NetworkChannel var1) {
this.clearTimeout();
}
class="kw">public void setQueuePackets(boolean var1) {
this.queuePackets = var1;
}
class="kw">public void tryFlush() {
if (this.queuedPackets.getAndSet(0) > 0) {
for (ChannelConnection var4 : this.channels) {
if (var4 != null) {
var4.flush();
}
}
}
}
class="kw">public void write(@Nonnull ToClientPacket... var1) {
if (var1.length != 0) {
ToClientPacket[] var2 = new ToClientPacket[var1.length];
this.handleOutboundAndCachePackets(var1, var2);
NetworkChannel var3 = var1[0].getChannel();
for (int var4 = 1; var4 < var1.length; var4++) {
if (var3 != var1[var4].getChannel()) {
throw new IllegalArgumentException("All packets must be sent on the same channel!");
}
}
ChannelConnection var5 = this.channels[var3.getValue()];
if (this.queuePackets) {
var5.write(var2);
this.queuedPackets.getAndIncrement();
} else {
var5.writeAndFlush(var2);
}
}
}
class="kw">public void write(@Nonnull ToClientPacket[] var1, @Nonnull ToClientPacket var2) {
ToClientPacket[] var3 = new ToClientPacket[var1.length + 1];
this.handleOutboundAndCachePackets(var1, var3);
var3[var3.length - 1] = this.handleOutboundAndCachePacket(var2);
NetworkChannel var4 = var2.getChannel();
for (int var5 = 0; var5 < var1.length; var5++) {
if (var4 != var1[var5].getChannel()) {
throw new IllegalArgumentException("All packets must be sent on the same channel!");
}
}
ChannelConnection var6 = this.channels[var4.getValue()];
if (this.queuePackets) {
var6.write(var3);
this.queuedPackets.getAndIncrement();
} else {
var6.writeAndFlush(var3);
}
}
@Override
class="kw">public void write(@Nonnull ToClientPacket var1) {
this.writePacket(var1, true);
}
@Override
class="kw">public void writeNoCache(@Nonnull ToClientPacket var1) {
this.writePacket(var1, false);
}
class="kw">public boolean writePacket(@Nonnull ToClientPacket var1, boolean var2) {
if (PacketAdapters.__handleOutbound(this, var1)) {
class="kw">return false;
}
ToClientPacket var3;
if (var2) {
var3 = this.handleOutboundAndCachePacket(var1);
} else {
var3 = var1;
}
ChannelConnection var4 = this.channels[var1.getChannel().getValue()];
if (this.queuePackets) {
var4.write(var3);
this.queuedPackets.getAndIncrement();
} else {
var4.writeAndFlush(var3);
}
class="kw">return true;
}
class="kw">private void handleOutboundAndCachePackets(@Nonnull ToClientPacket[] var1, @Nonnull ToClientPacket[] var2) {
for (int var3 = 0; var3 < var1.length; var3++) {
ToClientPacket var4 = var1[var3];
if (!PacketAdapters.__handleOutbound(this, var4)) {
var2[var3] = this.handleOutboundAndCachePacket(var4);
}
}
}
@Nonnull
class="kw">private ToClientPacket handleOutboundAndCachePacket(@Nonnull ToClientPacket var1) {
class="kw">return var1 class="kw">instanceof CachedPacket ? var1 : CachedPacket.cache(var1);
}
class="kw">public void disconnect(@Nonnull Message var1) {
this.disconnect(var1.getFormattedMessage());
}
class="kw">public void disconnect(@Nonnull FormattedMessage var1) {
this.disconnectReason.setServerDisconnectReason(var1);
String var2 = this.getSniHostname();
HytaleLogger.getLogger()
.at(Level.INFO)
.log("Disconnecting %s (SNI: %s) with the message: %s", this.getChannel().formatRemoteAddress(), var2, MessageUtil.formatMessageToPlainString(var1));
this.disconnect0(var1);
}
class="kw">protected void disconnect0(@Nonnull FormattedMessage var1) {
this.getChannel().disconnect(var1);
}
@Nullable
class="kw">public PacketStatsRecorder getPacketStatsRecorder() {
class="kw">return this.getChannel().getPacketStatsRecorder();
}
@Nonnull
class="kw">public PacketHandler.PingInfo getPingInfo(@Nonnull PongType var1) {
class="kw">return this.pingInfo[var1.ordinal()];
}
class="kw">public long getOperationTimeoutThreshold() {
double var1 = this.getPingInfo(PongType.Tick).getPingMetricSet().getAverage(0);
class="kw">return PacketHandler.PingInfo.TIME_UNIT.toMillis(Math.round(var1 * 2.0)) + 3000L;
}
class="kw">public void tickPing(float var1) {
this.pingTimer -= var1;
if (this.pingTimer <= 0.0F) {
this.pingTimer = 1.0F;
this.sendPing();
}
}
class="kw">public void sendPing() {
int var1 = this.pingIdRandom.nextInt();
Instant var2 = Instant.now();
long var3 = System.nanoTime();
for (PacketHandler.PingInfo var8 : this.pingInfo) {
var8.recordSent(var1, var3);
}
this.writeNoCache(
new Ping(
var1,
WorldTimeResource.instantToInstantData(var2),
(int)this.getPingInfo(PongType.Raw).getPingMetricSet().getLastValue(),
(int)this.getPingInfo(PongType.Direct).getPingMetricSet().getLastValue(),
(int)this.getPingInfo(PongType.Tick).getPingMetricSet().getLastValue()
)
);
}
class="kw">public void handlePong(@Nonnull Pong var1) {
this.pingInfo[var1.type.ordinal()].handlePacket(var1);
}
class="kw">protected void initStage(@Nonnull String var1, @Nonnull Duration var2, @Nonnull BooleanSupplier var3) {
this.getChannel().initTimeoutContext(var1, this.getIdentifier());
this.setStageTimeout(var1, var2, var3);
}
class="kw">protected void enterStage(@Nonnull String var1, @Nonnull Duration var2, @Nonnull BooleanSupplier var3) {
this.getChannel().updateTimeoutContext(var1, this.getIdentifier());
this.updatePacketTimeout(var2);
this.setStageTimeout(var1, var2, var3);
}
class="kw">protected void enterStage(@Nonnull String var1, @Nonnull Duration var2) {
this.getChannel().updateTimeoutContext(var1, this.getIdentifier());
this.updatePacketTimeout(var2);
}
class="kw">protected void continueStage(@Nonnull String var1, @Nonnull Duration var2, @Nonnull BooleanSupplier var3) {
this.getChannel().updateTimeoutContext(var1);
this.updatePacketTimeout(var2);
this.setStageTimeout(var1, var2, var3);
}
class="kw">private void setStageTimeout(@Nonnull String var1, @Nonnull Duration var2, @Nonnull BooleanSupplier var3) {
if (this class="kw">instanceof AuthenticationPacketHandler || !(this class="kw">instanceof PasswordPacketHandler) || this.auth != null) {
this.getChannel().setStageTimeout(var1, var2, var3, () -> this.disconnect(Message.translation("client.general.disconnect.stageTimeout")));
}
}
class="kw">private void updatePacketTimeout(@Nonnull Duration var1) {
this.getChannel().setPacketTimeout(var1);
}
class="kw">protected void clearTimeout() {
this.getChannel().clearStageTimeout();
if (this.clientReadyForChunksFuture != null) {
this.clientReadyForChunksFuture.cancel(true);
this.clientReadyForChunksFuture = null;
this.clientReadyForChunksFutureStack = null;
}
}
@Nullable
class="kw">public PlayerAuthentication getAuth() {
class="kw">return this.auth;
}
class="kw">public boolean stillActive() {
class="kw">return this.getChannel().isActive();
}
class="kw">public int getQueuedPacketsCount() {
class="kw">return this.queuedPackets.get();
}
class="kw">public boolean isLocalConnection() {
SocketAddress var1 = this.getChannel().remoteAddress();
if (var1 class="kw">instanceof InetSocketAddress) {
InetAddress var2 = ((InetSocketAddress)var1).getAddress();
class="kw">return NetworkUtil.addressMatchesAny(var2, NetworkUtil.AddressType.ANY_LOCAL, NetworkUtil.AddressType.LOOPBACK);
} else {
class="kw">return var1 class="kw">instanceof DomainSocketAddress || var1 class="kw">instanceof LocalAddress;
}
}
class="kw">public boolean isLANConnection() {
SocketAddress var1 = this.getChannel().remoteAddress();
if (var1 class="kw">instanceof InetSocketAddress) {
InetAddress var2 = ((InetSocketAddress)var1).getAddress();
class="kw">return NetworkUtil.addressMatchesAny(var2);
} else {
class="kw">return var1 class="kw">instanceof DomainSocketAddress || var1 class="kw">instanceof LocalAddress;
}
}
@Nullable
class="kw">public String getSniHostname() {
class="kw">return this.getChannel().getSniHostname();
}
class="kw">public boolean checkStreamOpenRateLimit() {
long var1 = System.nanoTime();
long var3 = this.lastStreamOpenTimeNanos.getAndUpdate(var2 -> var1 - var2 >= STREAM_OPEN_MIN_INTERVAL_NANOS ? var1 : var2);
class="kw">return var1 - var3 < STREAM_OPEN_MIN_INTERVAL_NANOS;
}
@Nonnull
class="kw">public PacketHandler.DisconnectReason getDisconnectReason() {
class="kw">return this.disconnectReason;
}
class="kw">public void setClientReadyForChunksFuture(@Nonnull CompletableFuture<Void> var1) {
if (this.clientReadyForChunksFuture != null) {
throw new IllegalStateException("Tried to hook client ready but something is already waiting for it!", this.clientReadyForChunksFutureStack);
}
HytaleLogger.getLogger().at(Level.WARNING).log("%s Added future for ClientReady packet?", this.getIdentifier());
this.clientReadyForChunksFutureStack = DebugFlags.CAPTURE_THROWABLES ? new Throwable() : null;
this.clientReadyForChunksFuture = var1;
}
@Nullable
class="kw">public CompletableFuture<Void> getClientReadyForChunksFuture() {
class="kw">return this.clientReadyForChunksFuture;
}
@Nonnull
class="kw">public ChannelConnection getChannel(@Nonnull NetworkChannel var1) {
class="kw">return this.channels[var1.getValue()];
}
class="kw">public void setChannel(@Nonnull NetworkChannel var1, @Nonnull ChannelConnection var2) {
this.channels[var1.getValue()] = var2;
}
class="kw">public void setChannel(@Nonnull StreamType var1, @Nullable ChannelConnection var2) {
if (var1 == StreamType.Game) {
throw new IllegalArgumentException("Cannot set Game stream via auxiliary channel API");
}
if (var2 != null) {
this.auxiliaryChannels.put(var1, var2);
} else {
this.auxiliaryChannels.remove(var1);
}
}
class="kw">public boolean compareAndSetChannel(@Nonnull StreamType var1, @Nullable ChannelConnection var2, @Nullable ChannelConnection var3) {
if (var1 == StreamType.Game) {
throw new IllegalArgumentException("Cannot CAS Game stream via auxiliary channel API");
}
class="kw">synchronized (this.auxiliaryChannels) {
ChannelConnection var5 = this.auxiliaryChannels.get(var1);
if (Objects.equals(var5, var2)) {
if (var3 != null) {
this.auxiliaryChannels.put(var1, var3);
} else {
this.auxiliaryChannels.remove(var1);
}
class="kw">return true;
} else {
class="kw">return false;
}
}
}
class="kw">public int getAuxiliaryChannelCount() {
class="kw">return this.auxiliaryChannels.size();
}
class="kw">public class="kw">static void logConnectionTimings(@Nonnull ChannelConnection var0, @Nonnull String var1, @Nonnull Level var2) {
var0.logConnectionTimings(var1, var2);
}
class="kw">static {
LOGIN_TIMING_LOGGER.setLevel(Level.ALL);
}
class="kw">public class="kw">static class DisconnectReason {
@Nullable
class="kw">private FormattedMessage serverDisconnectReason;
@Nullable
class="kw">private DisconnectType clientDisconnectType;
class="kw">protected DisconnectReason() {
}
@Nullable
class="kw">public String getServerDisconnectReason() {
class="kw">return this.serverDisconnectReason != null ? MessageUtil.formatMessageToPlainString(this.serverDisconnectReason) : null;
}
@Nullable
class="kw">public FormattedMessage getServerDisconnectReasonFormatted() {
class="kw">return this.serverDisconnectReason;
}
class="kw">public void setServerDisconnectReason(@Nullable FormattedMessage var1) {
this.serverDisconnectReason = var1;
this.clientDisconnectType = null;
}
@Deprecated
class="kw">public void setServerDisconnectReason(@Nullable String var1) {
this.setServerDisconnectReason(var1 != null ? Message.raw(var1).getFormattedMessage() : null);
}
@Nullable
class="kw">public DisconnectType getClientDisconnectType() {
class="kw">return this.clientDisconnectType;
}
class="kw">public void setClientDisconnectType(DisconnectType var1) {
this.clientDisconnectType = var1;
this.serverDisconnectReason = null;
}
@Nonnull
@Override
class="kw">public String toString() {
class="kw">return "DisconnectReason{serverDisconnectReason='" + this.serverDisconnectReason + "', clientDisconnectType=" + this.clientDisconnectType + "}";
}
}
class="kw">public class="kw">static class PingInfo {
class="kw">public class="kw">static class="kw">final MetricsRegistry<PacketHandler.PingInfo> METRICS_REGISTRY = new MetricsRegistry<PacketHandler.PingInfo>()
.register("PingType", var0 -> var0.pingType, new EnumCodec<>(PongType.class))
.register("PingMetrics", PacketHandler.PingInfo::getPingMetricSet, HistoricMetric.METRICS_CODEC)
.register("PacketQueueMin", var0 -> var0.packetQueueMetric.getMin(), Codec.LONG)
.register("PacketQueueAvg", var0 -> var0.packetQueueMetric.getAverage(), Codec.DOUBLE)
.register("PacketQueueMax", var0 -> var0.packetQueueMetric.getMax(), Codec.LONG);
class="kw">public class="kw">static class="kw">final TimeUnit TIME_UNIT = TimeUnit.MICROSECONDS;
class="kw">public class="kw">static class="kw">final int ONE_SECOND_INDEX = 0;
class="kw">public class="kw">static class="kw">final int ONE_MINUTE_INDEX = 1;
class="kw">public class="kw">static class="kw">final int FIVE_MINUTE_INDEX = 2;
class="kw">public class="kw">static class="kw">final double PERCENTILE = 0.99F;
class="kw">public class="kw">static class="kw">final int PING_FREQUENCY = 1;
class="kw">public class="kw">static class="kw">final TimeUnit PING_FREQUENCY_UNIT = TimeUnit.SECONDS;
class="kw">public class="kw">static class="kw">final int PING_FREQUENCY_MILLIS = 1000;
class="kw">public class="kw">static class="kw">final int PING_HISTORY_MILLIS = 15000;
class="kw">public class="kw">static class="kw">final int PING_HISTORY_LENGTH = 15;
class="kw">protected class="kw">final PongType pingType;
class="kw">protected class="kw">final Lock queueLock = new ReentrantLock();
class="kw">protected class="kw">final IntPriorityQueue pingIdQueue = new IntArrayFIFOQueue(15);
class="kw">protected class="kw">final LongPriorityQueue pingTimestampQueue = new LongArrayFIFOQueue(15);
class="kw">protected class="kw">final Lock pingLock = new ReentrantLock();
@Nonnull
class="kw">protected class="kw">final HistoricMetric pingMetricSet;
class="kw">protected class="kw">final Metric packetQueueMetric = new Metric();
class="kw">public PingInfo(PongType var1) {
this.pingType = var1;
this.pingMetricSet = HistoricMetric.builder(1000L, TimeUnit.MILLISECONDS)
.addPeriod(1L, TimeUnit.SECONDS)
.addPeriod(1L, TimeUnit.MINUTES)
.addPeriod(5L, TimeUnit.MINUTES)
.build();
}
class="kw">protected void recordSent(int var1, long var2) {
this.queueLock.lock();
try {
this.pingIdQueue.enqueue(var1);
this.pingTimestampQueue.enqueue(var2);
} class="kw">finally {
this.queueLock.unlock();
}
}
class="kw">protected void handlePacket(@Nonnull Pong var1) {
if (var1.type != this.pingType) {
throw new IllegalArgumentException("Got packet for " + var1.type + " but expected " + this.pingType);
}
this.queueLock.lock();
int var2;
long var3;
try {
var2 = this.pingIdQueue.dequeueInt();
var3 = this.pingTimestampQueue.dequeueLong();
} class="kw">finally {
this.queueLock.unlock();
}
if (var1.id != var2) {
throw new IllegalArgumentException(String.valueOf(var1.id));
}
long var5 = System.nanoTime();
long var7 = var5 - var3;
if (var7 <= 0L) {
throw new IllegalArgumentException(String.format("Ping must be received after its sent! %s", var7));
}
this.pingLock.lock();
try {
this.pingMetricSet.add(var5, TIME_UNIT.convert(var7, TimeUnit.NANOSECONDS));
this.packetQueueMetric.add(var1.packetQueueSize);
} class="kw">finally {
this.pingLock.unlock();
}
}
class="kw">public PongType getPingType() {
class="kw">return this.pingType;
}
@Nonnull
class="kw">public Metric getPacketQueueMetric() {
class="kw">return this.packetQueueMetric;
}
@Nonnull
class="kw">public HistoricMetric getPingMetricSet() {
class="kw">return this.pingMetricSet;
}
class="kw">public void clear() {
this.pingLock.lock();
try {
this.packetQueueMetric.clear();
this.pingMetricSet.clear();
} class="kw">finally {
this.pingLock.unlock();
}
}
}
}