LatencySimulationHandler class

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

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

extends: ChannelDuplexHandler

Поля (10)

МодификаторыТипИмя
final String PIPELINE_KEY
protected final ChannelHandlerContext ctx
protected final long executeAtNanos
private final Object msg
private final ChannelPromise promise
LatencySimulationHandler.DelayedHandler var1
ObjectArrayList var2
ObjectListIterator var3
LatencySimulationHandler.DelayedHandler var4
ChannelPipeline var4

Методы (17)

МодификаторыВозвратСигнатура
public DelayedFlushpublic DelayedFlush(ChannelHandlerContext var1, long var2)
protected DelayedHandlerprotected DelayedHandler(ChannelHandlerContext var1, long var2)
private DelayedReadprivate DelayedRead(ChannelHandlerContext var1, long var2)
public DelayedWritepublic DelayedWrite(ChannelHandlerContext var1, long var2, Object var4, ChannelPromise var5)
public void closevoid close(ChannelHandlerContext var1, ChannelPromise var2)
public int compareTopublic int compareTo(@Nonnull Delayed var1)
public void flushvoid flush(ChannelHandlerContext var1)
public long getDelaypublic long getDelay(@Nonnull TimeUnit var1)
public void handlerRemovedvoid handlerRemoved(ChannelHandlerContext var1)
if if(var1 > 0L)
public void readvoid read(ChannelHandlerContext var1)
public void runpublic void run()
public void runpublic void run()
public void runpublic void run()
static void setLatencyvoid setLatency(@Nonnull Channel var0, long var1, @Nonnull TimeUnit var3)
abstract super super(var1, var2)
public void writevoid write(ChannelHandlerContext var1, Object var2, ChannelPromise var3)

Исходный код

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

class="kw">import io.netty.channel.Channel;
class="kw">import io.netty.channel.ChannelDuplexHandler;
class="kw">import io.netty.channel.ChannelHandlerContext;
class="kw">import io.netty.channel.ChannelPipeline;
class="kw">import io.netty.channel.ChannelPromise;
class="kw">import it.unimi.dsi.fastutil.objects.ObjectArrayList;
class="kw">import it.unimi.dsi.fastutil.objects.ObjectListIterator;
class="kw">import java.time.Duration;
class="kw">import java.util.Comparator;
class="kw">import java.util.concurrent.DelayQueue;
class="kw">import java.util.concurrent.Delayed;
class="kw">import java.util.concurrent.TimeUnit;
class="kw">import java.util.concurrent.atomic.AtomicInteger;
class="kw">import javax.annotation.Nonnull;

class="kw">public class LatencySimulationHandler class="kw">extends ChannelDuplexHandler {
   class="kw">public class="kw">static class="kw">final String PIPELINE_KEY = "latencySimulator";
   class="kw">private class="kw">static class="kw">final AtomicInteger counter = new AtomicInteger();
   class="kw">private class="kw">final DelayQueue<LatencySimulationHandler.DelayedHandler> delayedQueue = new DelayQueue<>();
   @Nonnull
   class="kw">private class="kw">final Thread taskThread;
   class="kw">private class="kw">final long delayNanos;

   class="kw">public LatencySimulationHandler(long var1, @Nonnull TimeUnit var3) {
      this.delayNanos = var3.toNanos(var1);
      this.taskThread = new Thread(() -> {
         try {
            while (!Thread.interrupted()) {
               LatencySimulationHandler.DelayedHandler var1 = this.delayedQueue.take();
               var1.ctx.executor().execute(var1);
            }
         } catch (InterruptedException var2) {
            Thread.currentThread().interrupt();
         }
      }, "latency-simulator-" + counter.getAndIncrement());
      this.taskThread.setDaemon(true);
      this.taskThread.start();
   }

   class="kw">public void read(ChannelHandlerContext var1) class="kw">throws Exception {
      this.delayedQueue.offer(new LatencySimulationHandler.DelayedRead(var1, System.nanoTime() + this.delayNanos));
   }

   class="kw">public void write(ChannelHandlerContext var1, Object var2, ChannelPromise var3) class="kw">throws Exception {
      this.delayedQueue.offer(new LatencySimulationHandler.DelayedWrite(var1, System.nanoTime() + this.delayNanos, var2, var3));
   }

   class="kw">public void flush(ChannelHandlerContext var1) {
      this.delayedQueue.offer(new LatencySimulationHandler.DelayedFlush(var1, System.nanoTime() + this.delayNanos));
   }

   class="kw">public void handlerRemoved(ChannelHandlerContext var1) class="kw">throws Exception {
      super.handlerRemoved(var1);
      this.taskThread.interrupt();
      ObjectArrayList var2 = new ObjectArrayList(this.delayedQueue);
      var2.sort(Comparator.comparingLong(var0 -> var0.executeAtNanos));
      ObjectListIterator var3 = var2.iterator();

      while (var3.hasNext()) {
         LatencySimulationHandler.DelayedHandler var4 = (LatencySimulationHandler.DelayedHandler)var3.next();
         var4.run();
      }
   }

   class="kw">public void close(ChannelHandlerContext var1, ChannelPromise var2) class="kw">throws Exception {
      super.close(var1, var2);
      this.taskThread.interrupt();
   }

   class="kw">public class="kw">static void setLatency(@Nonnull Channel var0, long var1, @Nonnull TimeUnit var3) {
      ChannelPipeline var4 = var0.pipeline();
      if (var4.get("latencySimulator") == null) {
         if (var1 > 0L) {
            var4.addAfter("packetArrayEncoder", "latencySimulator", new LatencySimulationHandler(var1, var3));
         }
      } else if (var1 <= 0L) {
         var4.remove("latencySimulator");
      } else {
         var4.replace("latencySimulator", "latencySimulator", new LatencySimulationHandler(var1, var3));
      }
   }

   class="kw">private class="kw">static class DelayedFlush class="kw">extends LatencySimulationHandler.DelayedHandler {
      class="kw">public DelayedFlush(ChannelHandlerContext var1, long var2) {
         super(var1, var2);
      }

      @Override
      class="kw">public void run() {
         this.ctx.flush();
      }
   }

   class="kw">private class="kw">abstract class="kw">static class DelayedHandler class="kw">implements Delayed, Runnable {
      class="kw">protected class="kw">final ChannelHandlerContext ctx;
      class="kw">protected class="kw">final long executeAtNanos;

      class="kw">protected DelayedHandler(ChannelHandlerContext var1, long var2) {
         this.ctx = var1;
         this.executeAtNanos = var2;
      }

      @Override
      class="kw">public long getDelay(@Nonnull TimeUnit var1) {
         class="kw">return var1.convert(Duration.ofNanos(this.executeAtNanos - System.nanoTime()));
      }

      class="kw">public int compareTo(@Nonnull Delayed var1) {
         class="kw">return Long.compare(this.executeAtNanos, ((LatencySimulationHandler.DelayedHandler)var1).executeAtNanos);
      }
   }

   class="kw">private class="kw">static class DelayedRead class="kw">extends LatencySimulationHandler.DelayedHandler {
      class="kw">private DelayedRead(ChannelHandlerContext var1, long var2) {
         super(var1, var2);
      }

      @Override
      class="kw">public void run() {
         this.ctx.read();
      }
   }

   class="kw">private class="kw">static class DelayedWrite class="kw">extends LatencySimulationHandler.DelayedHandler {
      class="kw">private class="kw">final Object msg;
      class="kw">private class="kw">final ChannelPromise promise;

      class="kw">public DelayedWrite(ChannelHandlerContext var1, long var2, Object var4, ChannelPromise var5) {
         super(var1, var2);
         this.msg = var4;
         this.promise = var5;
      }

      @Override
      class="kw">public void run() {
         this.ctx.write(this.msg, this.promise);
      }
   }
}