IceResponder class

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

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

Поля (40)

МодификаторыТипИмя
final HytaleLogger LOGGER
int attempts
break
final long deadlineNanos
volatile long expiresAtNanos
final IceCredentials local
volatile long nextSendNanos
volatile boolean nominated
final Map<InetSocketAddress, IceResponder.Probe> probes
final IceCredentials remote
return
volatile byte[] transactionId
volatile boolean triggered
final Set<InetSocketAddress> validated
long var2
IceResponder.Probe var2
IceResponder.Session var3
IceResponder.Probe var3
ArrayList var3
IceResponder.Session var4
StunMessage var4
String var4
IceResponder.PendingCheck var4
Iterator var4
long var5
IceResponder.Session var5
IceResponder.Session var5
Entry var5
StunMessage var5
var5
var5
StunMessage var6
IceCredentials var6
IceResponder.Session var6
IceResponder.Probe var6
IceResponder.Session var7
IceCredentials var7
int var7
var7
Entry var9

Методы (45)

МодификаторыВозвратСигнатура
private record PendingCheckrecord PendingCheck(@Nonnull String localUfrag, long expiresAtNanos)
private Probeprivate Probe()
Session Session(IceCredentials var1, IceCredentials var2, long var3)
static void addCandidatesvoid addCandidates(IceResponder.Session var0, List<InetSocketAddress> var1)
abstract addCandidates addCandidates(var4, var3)
abstract addCandidates addCandidates(var7, var3)
abstract addCandidates addCandidates(var3, var2)
public void addRemoteCandidatesvoid addRemoteCandidates(@Nonnull String var1, @Nonnull List<InetSocketAddress> var2)
public boolean addSessionboolean addSession(@Nonnull IceCredentials var1, @Nonnull IceCredentials var2)
public boolean addSessionboolean addSession(@Nonnull IceCredentials var1, @Nonnull IceCredentials var2, @Nonnull List<InetSocketAddress> var3)
public void clearvoid clear()
static Level deadlineLogLevelLevel deadlineLogLevel(boolean var0)
for for(InetSocketAddress var5 : var1)
for for(Entry var5 : var3)
abstract for for(int var8 = 0; var8 < var3.size()
public void handleDatagramvoid handleDatagram(@Nonnull SocketAddress var1, @Nonnull ByteBuffer var2)
private void handleRequestvoid handleRequest(StunMessage var1, ByteBuffer var2, InetSocketAddress var3)
private void handleResponsevoid handleResponse(StunMessage var1, ByteBuffer var2, InetSocketAddress var3)
if if(var4 != null)
if if(var5 >= 32L)
if if(var3 != null)
if if(var1 instanceof InetSocketAddress var3)
if if(var4 != null)
if if(var4 != null)
if if(var5 != null)
if if(var4 != null)
if if(var5 != null)
if if(var1 >= var6.expiresAtNanos)
if if(var1 >= this.nextCheckNanos)
if if(this.lastServedUfrag != null)
if if(var6.attempts < 5)
if if(var6.triggered)
if if(var4 != null && var4.transactionId != null)
if if(var4 != null)
public boolean isValidatedboolean isValidated(@Nonnull InetSocketAddress var1)
void keepAlivevoid keepAlive()
private void recordPendingCheckvoid recordPendingCheck(String var1, byte[] var2)
public void removeSessionvoid removeSession(@Nonnull String var1)
private void sendvoid send(StunMessage var1, InetSocketAddress var2)
abstract void sendvoid send(@Nonnull ByteBuffer var1, @Nonnull SocketAddress var2)
private void sendCheckvoid sendCheck(IceResponder.Session var1, IceCredentials var2, InetSocketAddress var3, @Nullable IceResponder.Probe var4)
private boolean sendOneDueProbeboolean sendOneDueProbe(IceResponder.Session var1, long var2)
public int sessionCountint sessionCount()
public void tickvoid tick()
void tickvoid tick(long var1)

Исходный код

Показать/скрыть
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();
      }
   }
}