Probe class
Пакет: com.hypixel.hytale.server.core.io.ice
Файл: com/hypixel/hytale/server/core/io/ice/IceResponder.java
Поля (4)
| Модификаторы | Тип | Имя |
|---|---|---|
|
int |
attempts |
|
volatile long |
nextSendNanos |
|
volatile byte[] |
transactionId |
|
volatile boolean |
triggered |
Исходный код
Показать/скрыть
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();
}
}
}