DatagramSender interface

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

Файл: com/hypixel/hytale/server/core/io/ice/IceResponder.java

Методы (1)

МодификаторыВозвратСигнатура
abstract void sendvoid send(@Nonnull ByteBuffer var1, @Nonnull SocketAddress var2)

Исходный код

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

class="kw">import com.hypixel.hytale.logger.HytaleLogger;
class="kw">import java.net.InetSocketAddress;
class="kw">import java.net.SocketAddress;
class="kw">import java.nio.ByteBuffer;
class="kw">import java.time.Duration;
class="kw">import java.util.ArrayList;
class="kw">import java.util.Comparator;
class="kw">import java.util.Iterator;
class="kw">import java.util.List;
class="kw">import java.util.Map;
class="kw">import java.util.Set;
class="kw">import java.util.Map.Entry;
class="kw">import java.util.concurrent.ConcurrentHashMap;
class="kw">import java.util.concurrent.CopyOnWriteArraySet;
class="kw">import java.util.logging.Level;
class="kw">import javax.annotation.Nonnull;
class="kw">import javax.annotation.Nullable;

class="kw">public class="kw">final class IceResponder {
   class="kw">private class="kw">static class="kw">final HytaleLogger LOGGER = HytaleLogger.forEnclosingClass();
   class="kw">private class="kw">static class="kw">final int MAX_SESSIONS = 32;
   class="kw">private class="kw">static class="kw">final int MAX_TOTAL_SESSIONS = 256;
   class="kw">private class="kw">static class="kw">final int MAX_REMOTE_CANDIDATES = 16;
   class="kw">private class="kw">static class="kw">final int MAX_VALIDATED_PER_SESSION = 8;
   class="kw">public class="kw">static class="kw">final Duration CHECK_INTERVAL = Duration.ofMillis(50L);
   class="kw">private class="kw">static class="kw">final int MAX_TRANSMISSIONS = 5;
   class="kw">private class="kw">static class="kw">final Duration INITIAL_RTO = Duration.ofMillis(500L);
   class="kw">private class="kw">static class="kw">final Duration PROBE_BUDGET = Duration.ofSeconds(10L);
   class="kw">private class="kw">static class="kw">final Duration SESSION_IDLE_TIMEOUT = Duration.ofSeconds(60L);
   class="kw">private class="kw">static class="kw">final Duration UNVALIDATED_SESSION_TIMEOUT = Duration.ofSeconds(15L);
   class="kw">private class="kw">static class="kw">final Duration RESPONSE_WINDOW = Duration.ofSeconds(5L);
   class="kw">private class="kw">static class="kw">final int MAX_PENDING_CHECKS = 512;
   class="kw">private class="kw">final Map<String, IceResponder.Session> sessionsByUfrag = new ConcurrentHashMap<>();
   class="kw">private class="kw">final Map<ByteBuffer, IceResponder.PendingCheck> checksByTransactionId = new ConcurrentHashMap<>();
   class="kw">private class="kw">final IceResponder.DatagramSender sender;
   class="kw">private long nextCheckNanos;
   @Nullable
   class="kw">private String lastServedUfrag;

   class="kw">public IceResponder(@Nonnull IceResponder.DatagramSender var1) {
      this.sender = var1;
   }

   class="kw">public boolean addSession(@Nonnull IceCredentials var1, @Nonnull IceCredentials var2) {
      class="kw">return this.addSession(var1, var2, List.of());
   }

   class="kw">public boolean addSession(@Nonnull IceCredentials var1, @Nonnull IceCredentials var2, @Nonnull List<InetSocketAddress> var3) {
      if (!IceCredentials.isValid(var2.ufrag(), var2.password())) {
         LOGGER.at(Level.WARNING).log("Rejecting ICE session: malformed remote credentials");
         class="kw">return false;
      } else {
         IceResponder.Session var4 = this.sessionsByUfrag.get(var1.ufrag());
         if (var4 != null) {
            addCandidates(var4, var3);
            class="kw">return true;
         } else {
            long var5 = this.sessionsByUfrag.values().stream().filter(var0 -> var0.validated.isEmpty()).count();
            if (var5 >= 32L) {
               LOGGER.at(Level.WARNING).log("Rejecting ICE session: %d already negotiating", var5);
               class="kw">return false;
            } else if (this.sessionsByUfrag.size() >= 256) {
               LOGGER.at(Level.WARNING).log("Rejecting ICE session: %d registered", this.sessionsByUfrag.size());
               class="kw">return false;
            } else {
               IceResponder.Session var7 = this.sessionsByUfrag.computeIfAbsent(var1.ufrag(), var2x -> new IceResponder.Session(var1, var2, System.nanoTime()));
               addCandidates(var7, var3);
               class="kw">return true;
            }
         }
      }
   }

   class="kw">public void addRemoteCandidates(@Nonnull String var1, @Nonnull List<InetSocketAddress> var2) {
      IceResponder.Session var3 = this.sessionsByUfrag.get(var1);
      if (var3 != null) {
         addCandidates(var3, var2);
      }
   }

   class="kw">private class="kw">static void addCandidates(IceResponder.Session var0, List<InetSocketAddress> var1) {
      long var2 = System.nanoTime();

      for (InetSocketAddress var5 : var1) {
         if (var0.probes.size() >= 16) {
            break;
         }

         var0.probes.computeIfAbsent(var5, var2x -> {
            IceResponder.Probe var3 = new IceResponder.Probe();
            var3.nextSendNanos = var2;
            class="kw">return var3;
         });
      }
   }

   class="kw">public void removeSession(@Nonnull String var1) {
      this.sessionsByUfrag.remove(var1);
      this.checksByTransactionId.values().removeIf(var1x -> var1x.localUfrag().equals(var1));
   }

   class="kw">public void clear() {
      this.sessionsByUfrag.clear();
      this.checksByTransactionId.clear();
   }

   class="kw">private void recordPendingCheck(String var1, byte[] var2) {
      if (this.checksByTransactionId.size() < 512) {
         this.checksByTransactionId.put(ByteBuffer.wrap(var2), new IceResponder.PendingCheck(var1, System.nanoTime() + RESPONSE_WINDOW.toNanos()));
      }
   }

   class="kw">public int sessionCount() {
      class="kw">return this.sessionsByUfrag.size();
   }

   class="kw">public boolean isValidated(@Nonnull InetSocketAddress var1) {
      for (IceResponder.Session var3 : this.sessionsByUfrag.values()) {
         if (var3.validated.contains(var1)) {
            class="kw">return true;
         }
      }

      class="kw">return false;
   }

   class="kw">public void handleDatagram(@Nonnull SocketAddress var1, @Nonnull ByteBuffer var2) {
      if (var1 class="kw">instanceof InetSocketAddress var3) {
         StunMessage var4 = StunMessage.parse(var2);
         if (var4 != null) {
            if (var4.isRequest()) {
               this.handleRequest(var4, var2, var3);
            } else if (var4.isResponse()) {
               this.handleResponse(var4, var2, var3);
            }
         }
      }
   }

   class="kw">private void handleRequest(StunMessage var1, ByteBuffer var2, InetSocketAddress var3) {
      String var4 = IceCredentials.parseLocalUfrag(var1.getUsername());
      if (var4 != null) {
         IceResponder.Session var5 = this.sessionsByUfrag.get(var4);
         if (var5 != null) {
            if (!StunMessage.verifyMessageIntegrity(var2, var5.local.integrityKey())) {
               LOGGER.at(Level.FINE).log("Dropping ICE check from %s: message integrity failed", var3);
            } else if (!IceAddressFilter.isAcceptable(var3)) {
               LOGGER.at(Level.FINE).log("Dropping ICE check from disallowed source %s", var3);
            } else {
               var5.keepAlive();
               if (var5.validated.size() < 8) {
                  var5.validated.add(var3);
               }

               if (var1.hasAttribute(37)) {
                  var5.nominated = true;
                  LOGGER.at(Level.INFO).log("ICE peer %s nominated its path", var3);
               }

               StunMessage var6 = StunMessage.response(var1.getTransactionId());
               var6.addXorMappedAddress(var3);
               var6.appendMessageIntegrity(var5.local.integrityKey());
               var6.appendFingerprint();
               this.send(var6, var3);
               IceCredentials var7 = var5.remote;
               if (var7 != null && (var5.probes.containsKey(var3) || var5.probes.size() < 16)) {
                  var5.probes.compute(var3, (var0, var1x) -> {
                     IceResponder.Probe var2 = var1x != null ? var1x : new IceResponder.Probe();
                     var2.triggered = true;
                     var2.nextSendNanos = Math.min(var2.nextSendNanos, System.nanoTime());
                     class="kw">return var2;
                  });
               }
            }
         }
      }
   }

   class="kw">private void handleResponse(StunMessage var1, ByteBuffer var2, InetSocketAddress var3) {
      IceResponder.PendingCheck var4 = this.checksByTransactionId.remove(ByteBuffer.wrap(var1.getTransactionId()));
      if (var4 != null) {
         IceResponder.Session var5 = this.sessionsByUfrag.get(var4.localUfrag());
         if (var5 != null) {
            IceCredentials var6 = var5.remote;
            if (!StunMessage.verifyMessageIntegrity(var2, var6.integrityKey())) {
               LOGGER.at(Level.WARNING).log("Dropping ICE response from %s: message integrity failed", var3);
            } else {
               var5.keepAlive();
               if (var5.validated.size() < 8) {
                  var5.validated.add(var3);
               }

               LOGGER.at(Level.INFO).log("ICE check to %s succeeded", var3);
            }
         }
      }
   }

   class="kw">public void tick() {
      this.tick(System.nanoTime());
   }

   void tick(long var1) {
      this.checksByTransactionId.values().removeIf(var2 -> var1 >= var2.expiresAtNanos());
      ArrayList var3 = new ArrayList<>(this.sessionsByUfrag.entrySet());

      for (Entry var5 : var3) {
         IceResponder.Session var6 = var5.getValue();
         if (var1 >= var6.expiresAtNanos) {
            LOGGER.at(Level.FINE).log("Dropping idle ICE session %s", var5.getKey());
            this.sessionsByUfrag.remove(var5.getKey(), var6);
         } else if (var1 >= var6.deadlineNanos && !var6.probes.isEmpty()) {
            if (deadlineLogLevel(var6.validated.isEmpty()) == Level.INFO) {
               LOGGER.at(Level.INFO).log("ICE probing gave up for session %s", var5.getKey());
            } else {
               LOGGER.at(Level.FINE).log("ICE probing stopped with unvalidated candidates for session %s", var5.getKey());
            }

            var6.probes.clear();
         }
      }

      if (var1 >= this.nextCheckNanos) {
         var3.removeIf(var3x -> !this.sessionsByUfrag.containsKey(var3x.getKey()) || var1 >= var3x.getValue().deadlineNanos);
         var3.sort(Comparator.comparing(Entry::getKey));
         int var7 = 0;
         if (this.lastServedUfrag != null) {
            while (var7 < var3.size() && var3.get(var7).getKey().compareTo(this.lastServedUfrag) <= 0) {
               var7++;
            }

            if (var7 == var3.size()) {
               var7 = 0;
            }
         }

         for (int var8 = 0; var8 < var3.size(); var8++) {
            Entry var9 = var3.get((var7 + var8) % var3.size());
            if (this.sendOneDueProbe(var9.getValue(), var1)) {
               this.lastServedUfrag = var9.getKey();
               this.nextCheckNanos = var1 + CHECK_INTERVAL.toNanos();
               class="kw">return;
            }
         }
      }
   }

   class="kw">private boolean sendOneDueProbe(IceResponder.Session var1, long var2) {
      Iterator var4 = var1.probes.entrySet().iterator();

      while (var4.hasNext()) {
         Entry var5 = var4.next();
         IceResponder.Probe var6 = var5.getValue();
         if ((!var1.validated.contains(var5.getKey()) || var6.triggered) && var2 >= var6.nextSendNanos) {
            if (var6.attempts < 5) {
               var6.attempts++;
               var6.nextSendNanos = var2 + INITIAL_RTO.toNanos() * (1L << Math.min(var6.attempts - 1, 4));
               this.sendCheck(var1, var1.remote, var5.getKey(), var6);
               if (var6.triggered) {
                  var6.triggered = false;
               }

               class="kw">return true;
            }

            var6.triggered = false;
            if (!var1.validated.contains(var5.getKey())) {
               var4.remove();
            }
         }
      }

      class="kw">return false;
   }

   class="kw">static Level deadlineLogLevel(boolean var0) {
      class="kw">return var0 ? Level.INFO : Level.FINE;
   }

   class="kw">private void sendCheck(IceResponder.Session var1, IceCredentials var2, InetSocketAddress var3, @Nullable IceResponder.Probe var4) {
      StunMessage var5;
      if (var4 != null && var4.transactionId != null) {
         var5 = new StunMessage(1, var4.transactionId);
      } else {
         var5 = StunMessage.request();
         if (var4 != null) {
            var4.transactionId = var5.getTransactionId();
         }
      }

      var5.addUsername(var1.local.buildUsernameFor(var2.ufrag()));
      var5.addPriority(IceCandidate.PEER_REFLEXIVE_PRIORITY);
      var5.addIceRole(false, 0L);
      var5.appendMessageIntegrity(var2.integrityKey());
      var5.appendFingerprint();
      this.recordPendingCheck(var1.local.ufrag(), var5.getTransactionId());
      this.send(var5, var3);
   }

   class="kw">private void send(StunMessage var1, InetSocketAddress var2) {
      try {
         this.sender.send(ByteBuffer.wrap(var1.toBytes()), var2);
      } catch (Exception var4) {
         LOGGER.at(Level.FINE).log("Failed to send an ICE message to %s: %s", var2, var4.getMessage());
      }
   }

   @FunctionalInterface
   class="kw">public class="kw">interface DatagramSender {
      void send(@Nonnull ByteBuffer var1, @Nonnull SocketAddress var2);
   }

   class="kw">private record PendingCheck(@Nonnull String localUfrag, long expiresAtNanos) {
      class="kw">private PendingCheck {
      }
   }

   class="kw">private class="kw">static class="kw">final class Probe {
      int attempts;
      class="kw">volatile long nextSendNanos;
      class="kw">volatile boolean triggered;
      @Nullable
      class="kw">volatile byte[] transactionId;

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

   class="kw">private class="kw">static class="kw">final class Session {
      class="kw">final IceCredentials local;
      class="kw">final IceCredentials remote;
      class="kw">final Set<InetSocketAddress> validated = new CopyOnWriteArraySet<>();
      class="kw">final Map<InetSocketAddress, IceResponder.Probe> probes = new ConcurrentHashMap<>();
      class="kw">final long deadlineNanos;
      class="kw">volatile long expiresAtNanos;
      class="kw">volatile boolean nominated;

      Session(IceCredentials var1, IceCredentials var2, long var3) {
         this.local = var1;
         this.remote = var2;
         this.deadlineNanos = var3 + IceResponder.PROBE_BUDGET.toNanos();
         this.expiresAtNanos = var3 + IceResponder.UNVALIDATED_SESSION_TIMEOUT.toNanos();
      }

      void keepAlive() {
         this.expiresAtNanos = System.nanoTime() + IceResponder.SESSION_IDLE_TIMEOUT.toNanos();
      }
   }
}